mirror of
https://gitcode.com/JianFeeeee/HomeAgent.git
synced 2026-09-21 09:28:14 +00:00
feat(waiter): daemon模式 + localuse本机外设插件
waiter 新增 --daemon 后台驻留模式 (daemon.go): - 维持 homed 连接,TUI实例经 Unix socket 接入 - 单客户端串行模型:每个TUI独占homed响应,新客户端回放缓冲(256行) - 设备桥场景可无 homed 运行(纯设备桥驻留) - 设备桥看护循环:WS断开自动重连(3-5s间隔) - conn.go 新增 daemonConn 类型,dial() 优先检测 daemon localuse 插件 (internal/plugins/localuse/): - local_screensee / local_camerasue / local_speakeruse - local_screensue / local_clipboardsee / local_clipboardsue / local_computeruse - 跨平台实现(Linux/macOS/Windows),能力与 waiter device.go 对齐 - headless服务器上缺依赖工具自动返回安装提示 SDK sync: third_party SetToolBlocks 多模态类型同步 README 补充 daemon 模式使用文档
This commit is contained in:
@ -19,6 +19,12 @@ type Conn interface {
|
||||
}
|
||||
|
||||
func dial(cfg *Config) (Conn, error) {
|
||||
// 优先连接运行中的 daemon(后台驻留模式)
|
||||
if daemonIsRunning() {
|
||||
if c, err := dialDaemon(); err == nil {
|
||||
return c, nil
|
||||
}
|
||||
}
|
||||
if cfg.Remote != "" {
|
||||
return dialRemote(cfg.Remote, cfg.APIKey)
|
||||
}
|
||||
@ -52,6 +58,39 @@ func dialLocal(socket, apiKey string) (Conn, error) {
|
||||
return lc, nil
|
||||
}
|
||||
|
||||
// daemonConn 是连接到运行中 waiter daemon 的轻量封装。
|
||||
// 协议与 localConn 完全一致(行式 \n 分隔),但不做 /auth(daemon 已集中鉴权)。
|
||||
type daemonConn struct {
|
||||
conn net.Conn
|
||||
r *bufio.Reader
|
||||
}
|
||||
|
||||
func dialDaemon() (Conn, error) {
|
||||
sock := daemonSocketPath()
|
||||
c, err := net.DialTimeout("unix", sock, 3*time.Second)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("daemon %s: %w", sock, err)
|
||||
}
|
||||
return &daemonConn{conn: c, r: bufio.NewReader(c)}, nil
|
||||
}
|
||||
|
||||
func (c *daemonConn) Send(line string) error {
|
||||
_, err := fmt.Fprintf(c.conn, "%s\n", line)
|
||||
return err
|
||||
}
|
||||
|
||||
func (c *daemonConn) ReadLine() (string, error) {
|
||||
s, err := c.r.ReadString('\n')
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
return strings.TrimSuffix(s, "\n"), nil
|
||||
}
|
||||
|
||||
func (c *daemonConn) Close() error {
|
||||
return c.conn.Close()
|
||||
}
|
||||
|
||||
type localConn struct {
|
||||
conn net.Conn
|
||||
r *bufio.Reader
|
||||
|
||||
378
cmd/waiter/daemon.go
Normal file
378
cmd/waiter/daemon.go
Normal file
@ -0,0 +1,378 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"fmt"
|
||||
"log"
|
||||
"net"
|
||||
"os"
|
||||
"os/signal"
|
||||
"path/filepath"
|
||||
"runtime"
|
||||
"strings"
|
||||
"sync"
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/devicebridge/client"
|
||||
)
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Daemon 模式:后台驻留,维持 homed 连接 + 设备桥 + 消息缓冲
|
||||
//
|
||||
// 工作原理:
|
||||
// - daemon 保持一个到 homed 的持久连接
|
||||
// - TUI 实例通过 Unix socket 连接到 daemon
|
||||
// - daemon 为每个 TUI 客户端分配独立的 homed 响应(通过 homeMu 序列化)
|
||||
// - 新客户端连入时回放缓冲的历史消息(方便重连后看到上下文)
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
const (
|
||||
daemonSocketName = "waiter.sock"
|
||||
msgBufCap = 256 // 环形缓冲最近 N 行 homed 输出
|
||||
)
|
||||
|
||||
type msgEntry struct {
|
||||
line string
|
||||
seq uint64
|
||||
}
|
||||
|
||||
type daemonHandler struct {
|
||||
// homed 连接
|
||||
homeMu sync.Mutex
|
||||
homeConn net.Conn
|
||||
homeR *bufio.Reader
|
||||
homeCfg *Config
|
||||
|
||||
// 消息缓冲(新客户端连入时回放)
|
||||
bufMu sync.Mutex
|
||||
buf []msgEntry
|
||||
bufSeq uint64
|
||||
bufCap int
|
||||
|
||||
// 生命周期
|
||||
stopCh chan struct{}
|
||||
}
|
||||
|
||||
func newDaemonHandler() *daemonHandler {
|
||||
return &daemonHandler{
|
||||
bufCap: msgBufCap,
|
||||
stopCh: make(chan struct{}),
|
||||
}
|
||||
}
|
||||
|
||||
// ===== 消息缓冲 =====
|
||||
|
||||
func (h *daemonHandler) appendBuf(line string) {
|
||||
h.bufMu.Lock()
|
||||
defer h.bufMu.Unlock()
|
||||
h.bufSeq++
|
||||
h.buf = append(h.buf, msgEntry{line: line, seq: h.bufSeq})
|
||||
if len(h.buf) > h.bufCap {
|
||||
h.buf = h.buf[len(h.buf)-h.bufCap:]
|
||||
}
|
||||
}
|
||||
|
||||
func (h *daemonHandler) replayBuffer() []string {
|
||||
h.bufMu.Lock()
|
||||
defer h.bufMu.Unlock()
|
||||
lines := make([]string, 0, len(h.buf))
|
||||
for _, e := range h.buf {
|
||||
lines = append(lines, e.line)
|
||||
}
|
||||
return lines
|
||||
}
|
||||
|
||||
// ===== homed 连接 =====
|
||||
|
||||
func (h *daemonHandler) connectHome(cfg *Config) error {
|
||||
h.homeCfg = cfg
|
||||
if cfg.Remote != "" {
|
||||
return fmt.Errorf("daemon: remote mode not supported")
|
||||
}
|
||||
if cfg.Socket == "" {
|
||||
cfg.Socket = discoverSocket("")
|
||||
}
|
||||
c, err := net.DialTimeout("unix", cfg.Socket, 5*time.Second)
|
||||
if err != nil {
|
||||
return fmt.Errorf("daemon: connect home: %w", err)
|
||||
}
|
||||
h.homeConn = c
|
||||
h.homeR = bufio.NewReader(c)
|
||||
log.Printf("[daemon] connected to home %s", cfg.Socket)
|
||||
return nil
|
||||
}
|
||||
|
||||
func (h *daemonHandler) closeHome() {
|
||||
if h.homeConn != nil {
|
||||
h.homeConn.Close()
|
||||
h.homeConn = nil
|
||||
}
|
||||
}
|
||||
|
||||
func (h *daemonHandler) reconnectHome() {
|
||||
cfg := h.homeCfg
|
||||
if cfg == nil {
|
||||
cfg = discoverConfig("")
|
||||
}
|
||||
if cfg.Socket == "" && cfg.Remote == "" {
|
||||
cfg.Socket = discoverSocket("")
|
||||
}
|
||||
for i := 0; i < 30; i++ {
|
||||
select {
|
||||
case <-h.stopCh:
|
||||
return
|
||||
default:
|
||||
}
|
||||
h.closeHome()
|
||||
time.Sleep(2 * time.Second)
|
||||
if err := h.connectHome(cfg); err != nil {
|
||||
log.Printf("[daemon] reconnect home (%d/30): %v", i+1, err)
|
||||
continue
|
||||
}
|
||||
log.Printf("[daemon] reconnected to home")
|
||||
return
|
||||
}
|
||||
log.Printf("[daemon] gave up reconnecting to home")
|
||||
}
|
||||
|
||||
// handleClient 处理单个 TUI 客户端:
|
||||
// 1. 回放缓冲历史
|
||||
// 2. 读客户端输入 → 转发到 homed
|
||||
// 3. 读 homed 响应 → 回写给该客户端(独占响应,不广播)
|
||||
func (h *daemonHandler) handleClient(c net.Conn) {
|
||||
defer c.Close()
|
||||
|
||||
cid := fmt.Sprintf("%s", c.RemoteAddr())
|
||||
log.Printf("[daemon] client %s connected", cid)
|
||||
defer log.Printf("[daemon] client %s disconnected", cid)
|
||||
|
||||
// 1) 回放缓冲(新客户端看到最近对话上下文)
|
||||
for _, line := range h.replayBuffer() {
|
||||
fmt.Fprintf(c, "%s\n", line)
|
||||
}
|
||||
|
||||
// 2) 循环:读客户端 → 转发 homed → 读 homed 响应 → 回写客户端
|
||||
reader := bufio.NewReader(c)
|
||||
for {
|
||||
c.SetReadDeadline(time.Now().Add(5 * time.Minute))
|
||||
line, err := reader.ReadString('\n')
|
||||
if err != nil {
|
||||
if netErr, ok := err.(net.Error); ok && netErr.Timeout() {
|
||||
continue
|
||||
}
|
||||
return
|
||||
}
|
||||
line = strings.TrimSuffix(line, "\n")
|
||||
if line == "" {
|
||||
continue
|
||||
}
|
||||
|
||||
// 转发到 homed(加锁保证请求-响应配对)
|
||||
h.homeMu.Lock()
|
||||
if h.homeConn == nil {
|
||||
h.homeMu.Unlock()
|
||||
fmt.Fprintf(c, `{"type":"error","error":"not connected to home"}`+"\n")
|
||||
continue
|
||||
}
|
||||
_, sendErr := fmt.Fprintf(h.homeConn, "%s\n", line)
|
||||
if sendErr != nil {
|
||||
h.homeMu.Unlock()
|
||||
fmt.Fprintf(c, `{"type":"error","error":"send failed"}`+"\n")
|
||||
continue
|
||||
}
|
||||
|
||||
// 读 homed 响应(所有帧:reasoning_delta / content_delta / tool_call / response / error)
|
||||
for {
|
||||
h.homeConn.SetReadDeadline(time.Now().Add(60 * time.Second))
|
||||
respLine, readErr := h.homeR.ReadString('\n')
|
||||
if readErr != nil {
|
||||
h.homeMu.Unlock()
|
||||
log.Printf("[daemon] home read error during client %s: %v", cid, readErr)
|
||||
h.reconnectHome()
|
||||
// 回写错误给客户端
|
||||
fmt.Fprintf(c, `{"type":"error","error":"home disconnected"}`+"\n")
|
||||
goto nextMessage
|
||||
}
|
||||
respLine = strings.TrimSuffix(respLine, "\n")
|
||||
if respLine == "" {
|
||||
continue
|
||||
}
|
||||
// 写入缓冲 + 回写给发起请求的客户端
|
||||
h.appendBuf(respLine)
|
||||
fmt.Fprintf(c, "%s\n", respLine)
|
||||
|
||||
// 检查是否是终结帧
|
||||
if strings.Contains(respLine, `"type":"response"`) || strings.Contains(respLine, `"type":"error"`) {
|
||||
break
|
||||
}
|
||||
}
|
||||
h.homeMu.Unlock()
|
||||
|
||||
nextMessage:
|
||||
}
|
||||
}
|
||||
|
||||
// ===== 启动入口 =====
|
||||
|
||||
func runDaemon(cfg *Config) {
|
||||
dh := newDaemonHandler()
|
||||
|
||||
// 连接 homed(设备桥场景下可失败——被控主机无需 homed)
|
||||
if cfg.Socket != "" || cfg.Remote != "" {
|
||||
if err := dh.connectHome(cfg); err != nil {
|
||||
log.Printf("[daemon] home connect failed: %v (continue with device bridge only)", err)
|
||||
dh.homeConn = nil
|
||||
}
|
||||
} else {
|
||||
log.Printf("[daemon] no home socket configured, running device bridge only")
|
||||
}
|
||||
defer dh.closeHome()
|
||||
|
||||
// 启动设备桥(设备网关场景下为核心职责)
|
||||
startDaemonDeviceBridge(cfg)
|
||||
|
||||
// 监听 Unix socket
|
||||
sockPath := daemonSocketPath()
|
||||
os.Remove(sockPath)
|
||||
os.MkdirAll(filepath.Dir(sockPath), 0755)
|
||||
|
||||
ln, err := net.Listen("unix", sockPath)
|
||||
if err != nil {
|
||||
fmt.Fprintf(os.Stderr, "daemon: listen %s: %v\n", sockPath, err)
|
||||
os.Exit(1)
|
||||
}
|
||||
defer func() {
|
||||
ln.Close()
|
||||
os.Remove(sockPath)
|
||||
}()
|
||||
|
||||
log.Printf("[daemon] listening on %s", sockPath)
|
||||
fmt.Printf("waiter daemon started\n socket: %s\n press Ctrl+C to stop\n", sockPath)
|
||||
|
||||
sigCh := make(chan os.Signal, 1)
|
||||
signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM)
|
||||
go func() {
|
||||
<-sigCh
|
||||
log.Printf("[daemon] shutting down")
|
||||
close(dh.stopCh)
|
||||
dh.closeHome()
|
||||
ln.Close()
|
||||
}()
|
||||
|
||||
// 单客户端模式:串行处理(同一时刻只有一个 TUI 连接)
|
||||
// 这与 homed CLI 插件的行为一致——一个连接对应一个活跃会话。
|
||||
for {
|
||||
conn, err := ln.Accept()
|
||||
if err != nil {
|
||||
select {
|
||||
case <-dh.stopCh:
|
||||
log.Printf("[daemon] stopped")
|
||||
return
|
||||
default:
|
||||
log.Printf("[daemon] accept error: %v", err)
|
||||
continue
|
||||
}
|
||||
}
|
||||
dh.handleClient(conn)
|
||||
}
|
||||
}
|
||||
|
||||
// ===== Socket 工具 =====
|
||||
|
||||
func daemonSocketPath() string {
|
||||
home, _ := os.UserHomeDir()
|
||||
if home == "" {
|
||||
home = "/tmp"
|
||||
}
|
||||
return filepath.Join(home, ".homeagent", daemonSocketName)
|
||||
}
|
||||
|
||||
func daemonIsRunning() bool {
|
||||
sock := daemonSocketPath()
|
||||
c, err := net.DialTimeout("unix", sock, 500*time.Millisecond)
|
||||
if err != nil {
|
||||
return false
|
||||
}
|
||||
c.Close()
|
||||
return true
|
||||
}
|
||||
|
||||
func startDaemonDeviceBridge(cfg *Config) {
|
||||
dg := cfg.DeviceGateway
|
||||
dt := cfg.DeviceToken
|
||||
if dg == "" || dt == "" {
|
||||
return
|
||||
}
|
||||
// 设备桥重连循环:WS 断开时自动重连
|
||||
go runDeviceBridgeLoop(dg, dt)
|
||||
}
|
||||
|
||||
// runDeviceBridgeLoop 无限重连循环:建立设备桥 → 等待断开 → 重连。
|
||||
func runDeviceBridgeLoop(gateway, token string) {
|
||||
for {
|
||||
bridge, err := connectDeviceBridge(gateway, token)
|
||||
if err != nil {
|
||||
log.Printf("[daemon] device bridge connect failed: %v, retrying in 5s", err)
|
||||
time.Sleep(5 * time.Second)
|
||||
continue
|
||||
}
|
||||
log.Printf("[daemon] device bridge connected, waiting...")
|
||||
bridge.Wait() // 阻塞直到断开
|
||||
log.Printf("[daemon] device bridge disconnected, reconnecting in 3s")
|
||||
time.Sleep(3 * time.Second)
|
||||
}
|
||||
}
|
||||
|
||||
// connectDeviceBridge 创建并启动一次设备桥,返回 bridge 实例供 Wait()。
|
||||
func connectDeviceBridge(gateway, token string) (*client.Bridge, error) {
|
||||
hostname, _ := os.Hostname()
|
||||
if hostname == "" {
|
||||
hostname = "local"
|
||||
}
|
||||
deviceID := "waiter-" + sanitizeID(hostname)
|
||||
|
||||
caps := []string{
|
||||
"status", "cmdrun", "deviceinfo", "cmdresult",
|
||||
"computeruse", "screensee", "clipboardsee", "clipboardsue",
|
||||
"camerasue", "speakeruse", "screensue",
|
||||
}
|
||||
|
||||
info := map[string]interface{}{
|
||||
"hostname": hostname,
|
||||
"platform": runtime.GOOS,
|
||||
"arch": runtime.GOARCH,
|
||||
"cpus": runtime.NumCPU(),
|
||||
}
|
||||
|
||||
// 确保 gateway URL 格式正确
|
||||
gw := gateway
|
||||
if !strings.HasPrefix(gw, "ws://") && !strings.HasPrefix(gw, "wss://") {
|
||||
gw = "ws://" + gw
|
||||
}
|
||||
if !strings.Contains(gw, "/api/v1/device/ws") {
|
||||
gw = gw + "/api/v1/device/ws"
|
||||
}
|
||||
|
||||
bridge := client.New(gw, token, deviceID, hostname, caps, info)
|
||||
|
||||
// 注册命令处理器
|
||||
cr := client.NewCmdRouter()
|
||||
cr.Handle("homeagent-", handleHomeagentCmd)
|
||||
cr.HandleDefault(handleShellCmd)
|
||||
bridge.OnCmd(func(reqID, command string) {
|
||||
cr.Dispatch(reqID, command)
|
||||
})
|
||||
|
||||
if err := bridge.Start(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// 设置全局变量供 sendBridgeResult 使用
|
||||
deviceBridge = bridge
|
||||
deviceBridgeID = deviceID
|
||||
auth := true // daemon 模式默认授权(配置已指定)
|
||||
bridge.SetAuthorized(auth)
|
||||
|
||||
return bridge, nil
|
||||
}
|
||||
@ -126,6 +126,7 @@ func main() {
|
||||
deviceGateway := flag.String("device", "", "remotedevice 网关地址(如 127.0.0.1:9890),启动设备桥")
|
||||
deviceToken := flag.String("device-token", "", "设备接入 token")
|
||||
deviceAuthorized := flag.Bool("device-authorized", false, "客户端本地授权(允许远程操控本机;也可在 waiter.yaml 配 device_authorized: true)")
|
||||
daemonMode := flag.Bool("daemon", false, "后台驻留模式:维持 homed 连接 + 设备桥,等待 TUI 实例接入")
|
||||
testCap := flag.String("test-cap", "", "测试本地能力(screensue/speakeruse/screensee/clipboardsee/clipboardsue/computeruse/camerasue),如 --test-cap screensue")
|
||||
testCapArgs := flag.String("test-cap-args", "", "测试能力的参数")
|
||||
flag.Parse()
|
||||
@ -144,6 +145,12 @@ func main() {
|
||||
cfg.Socket = discoverSocket("")
|
||||
}
|
||||
|
||||
// Daemon 模式:后台驻留
|
||||
if *daemonMode {
|
||||
runDaemon(cfg)
|
||||
return
|
||||
}
|
||||
|
||||
oneShotMsg := *chat
|
||||
if oneShotMsg == "" {
|
||||
oneShotMsg = *say
|
||||
|
||||
Reference in New Issue
Block a user