mirror of
https://gitcode.com/JianFeeeee/HomeAgent.git
synced 2026-09-22 18:08:04 +00:00
fix(remotedevice): 媒体回传落盘 + webui SSE panic + GUI 相机跨平台
1. remotedevice 媒体落盘(核心改动) 设备录像/照片二进制聚合后写入 <data>/device_media/<reqID>.<ext>, cmd_result 返回 file 路径,不再 base64 内联——10s 录像数 MB 的 base64 会撑爆 LLM 上下文与工具结果管道。未配置目录时保持旧内联行为。 新增 TestWSBinaryMediaToFile 覆盖。 2. webui SSE 'send on closed channel' panic(生产单日 4924 次) handleChatEvents 的 defer close(writeCh) 与 Subscribe 回调闭包竞态: handler 退出后总线仍可能异步触发回调向已关闭 channel 发送。 改为 writer goroutine select on done 退出,不 close channel; defer 中等待 writerDone 保证无残余写入。 顺带补 mockPluginMgr.StopAndUnload(ae42e48 接口变更漏改测试)。 3. GUI camerasue 平台分支 ffmpeg 参数原硬编码 Linux v4l2(/dev/video0),Windows 上必然失败。 现按平台探测:win32=dshow(枚举设备名取第一个视频设备)、 darwin=avfoundation、linux=v4l2;录像编码 Windows 交给 mp4 muxer 默认。
This commit is contained in:
@ -9,6 +9,8 @@ import (
|
||||
"net"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
@ -207,6 +209,73 @@ func TestWSBinaryChunkUpload(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// ===== 媒体落盘:SetMediaDir 后 cmd_result 返回 file 路径而非 base64 内联 =====
|
||||
func TestWSBinaryMediaToFile(t *testing.T) {
|
||||
reg := NewRegistry()
|
||||
token := "test-token-123"
|
||||
reg.SetAcceptToken(func(provided string) bool { return provided == token })
|
||||
mediaDir := t.TempDir()
|
||||
reg.SetMediaDir(mediaDir)
|
||||
|
||||
srv := httptest.NewServer(http.HandlerFunc(reg.ServeWS))
|
||||
defer srv.Close()
|
||||
|
||||
cli := dialTestWS(t, srv.URL, token)
|
||||
defer cli.close()
|
||||
|
||||
cli.sendText([]byte(`{"op":"hello","device":{"device_id":"gui-media","name":"媒体机","kind":"computer","caps":["cmd"]}}`))
|
||||
if _, _, err := cli.readMsg(); err != nil { // hello_ack
|
||||
t.Fatalf("read hello_ack: %v", err)
|
||||
}
|
||||
|
||||
videoData := make([]byte, 30000)
|
||||
for i := range videoData {
|
||||
videoData[i] = byte(i % 253)
|
||||
}
|
||||
go func() {
|
||||
time.Sleep(50 * time.Millisecond)
|
||||
cli.sendText(mustJSON(map[string]interface{}{
|
||||
"op": "cmd_data_start", "req_id": "req-file-1",
|
||||
"kind": "camera_video", "mime": "video/mp4",
|
||||
"total": len(videoData),
|
||||
}))
|
||||
const chunk = 8192
|
||||
for off := 0; off < len(videoData); off += chunk {
|
||||
end := off + chunk
|
||||
if end > len(videoData) {
|
||||
end = len(videoData)
|
||||
}
|
||||
cli.sendBinary(videoData[off:end])
|
||||
}
|
||||
cli.sendText(mustJSON(map[string]interface{}{
|
||||
"op": "cmd_data_end", "req_id": "req-file-1", "status": "ok",
|
||||
}))
|
||||
}()
|
||||
|
||||
res, err := reg.AwaitResult("req-file-1", 5*time.Second)
|
||||
if err != nil {
|
||||
t.Fatalf("await result: %v", err)
|
||||
}
|
||||
// 落盘模式:file 字段存在且内容一致;不应再有 data_base64
|
||||
fp, ok := res["file"].(string)
|
||||
if !ok || fp == "" {
|
||||
t.Fatalf("expected file path in result, got %v", res)
|
||||
}
|
||||
if _, hasB64 := res["data_base64"]; hasB64 {
|
||||
t.Fatal("data_base64 should be absent in file mode")
|
||||
}
|
||||
if want := filepath.Join(mediaDir, "req-file-1.mp4"); fp != want {
|
||||
t.Fatalf("file path = %s, want %s", fp, want)
|
||||
}
|
||||
got, err := os.ReadFile(fp)
|
||||
if err != nil {
|
||||
t.Fatalf("read media file: %v", err)
|
||||
}
|
||||
if string(got) != string(videoData) {
|
||||
t.Fatal("media file content mismatch")
|
||||
}
|
||||
}
|
||||
|
||||
// ===== 端到端:PushData 下发音频(网关→设备 cmd_speech 协议)=====
|
||||
|
||||
func TestWSPushDataAudio(t *testing.T) {
|
||||
|
||||
@ -8,6 +8,7 @@ import (
|
||||
"fmt"
|
||||
"log"
|
||||
"net/http"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
@ -104,6 +105,15 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error {
|
||||
log.Printf("[remotedevice] register devicectl channel: %v", err)
|
||||
}
|
||||
|
||||
// ---- 媒体落盘目录:<data>/device_media ----------------
|
||||
// 设备回传的录像/照片等二进制聚合后写入此目录,cmd_result 返回 file 路径,
|
||||
// 避免 base64 内联撑爆 LLM 上下文。目录由 logManager/运维定期清理。
|
||||
if dataDir, err := s.Settings().GetCore("daemon.data_dir"); err == nil {
|
||||
if dd, ok := dataDir.(string); ok && dd != "" {
|
||||
p.registry.SetMediaDir(filepath.Join(dd, "device_media"))
|
||||
}
|
||||
}
|
||||
|
||||
// ---- 设备主动上报事件 → agent 注入 ----------------
|
||||
// 摄像头发现异常/传感器报警等场景:设备经 WS op=event 上报,
|
||||
// 插件将其格式化为文本经 SDK InjectText 异步注入 agent(source=device/{id},
|
||||
|
||||
@ -11,6 +11,8 @@ import (
|
||||
"log"
|
||||
"net"
|
||||
"net/http"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
@ -49,6 +51,7 @@ type Registry struct {
|
||||
acceptFn func(token string) bool
|
||||
cmdPending map[string]chan map[string]interface{} // reqID -> 结果 channel
|
||||
results map[string]resultEntry // reqID -> 已留档结果
|
||||
mediaDir string // 设备回传媒体落盘目录;空则退化为 base64 内联
|
||||
}
|
||||
|
||||
// resultEntry 保存一次 cmdrun 的结果(供 device_ctl_cmdresult 查询)。
|
||||
@ -120,6 +123,43 @@ func deviceSupportsTool(caps []string, tool string) bool {
|
||||
return !hasKnown // 未声明任何已知能力 → 全能力兼容
|
||||
}
|
||||
|
||||
// SetMediaDir 设置设备回传媒体的落盘目录。
|
||||
// 非空时 cmd_data_end 聚合完成后写入该目录,cmd_result 返回 file 路径
|
||||
// (大体积 base64 内联会撑爆 LLM 上下文与工具结果管道);空则保持旧的内联行为。
|
||||
func (r *Registry) SetMediaDir(dir string) {
|
||||
r.mu.Lock()
|
||||
r.mediaDir = dir
|
||||
r.mu.Unlock()
|
||||
}
|
||||
|
||||
// mediaExt 按 mime/kind 推断扩展名。
|
||||
func mediaExt(mime, kind string) string {
|
||||
m := strings.ToLower(mime)
|
||||
switch {
|
||||
case strings.Contains(m, "mp4"):
|
||||
return ".mp4"
|
||||
case strings.Contains(m, "webm"):
|
||||
return ".webm"
|
||||
case strings.Contains(m, "jpeg"), strings.Contains(m, "jpg"):
|
||||
return ".jpg"
|
||||
case strings.Contains(m, "png"):
|
||||
return ".png"
|
||||
case strings.Contains(m, "wav"):
|
||||
return ".wav"
|
||||
case strings.Contains(m, "mpeg"), strings.Contains(m, "mp3"):
|
||||
return ".mp3"
|
||||
}
|
||||
k := strings.ToLower(kind)
|
||||
if strings.Contains(k, "video") {
|
||||
return ".mp4"
|
||||
}
|
||||
if strings.Contains(k, "image") || strings.Contains(k, "camera_photo") {
|
||||
return ".jpg"
|
||||
}
|
||||
return ".bin"
|
||||
}
|
||||
|
||||
// NewRegistry 返回初始化后的设备注册表。
|
||||
func NewRegistry() *Registry {
|
||||
return &Registry{
|
||||
devices: make(map[string]*DeviceMeta),
|
||||
@ -743,8 +783,28 @@ func (r *Registry) handleWS(conn net.Conn, rw *bufio.ReadWriter) {
|
||||
"mime": acc.mime,
|
||||
"size": len(data),
|
||||
"expected": acc.total,
|
||||
// base64 编码完整二进制(录像 mp4 等),供 agent/上层取回后解码使用
|
||||
"data_base64": base64.StdEncoding.EncodeToString(data),
|
||||
}
|
||||
// 媒体落盘模式:写入 <mediaDir>/<reqID>.<ext>,cmd_result 返回 file 路径。
|
||||
// 大体积 base64 内联会撑爆 LLM 上下文(一段 10s 录像即数 MB),
|
||||
// agent 应拿路径后用 files/describe_image/ocr 等工具消费。
|
||||
r.mu.RLock()
|
||||
mediaDir := r.mediaDir
|
||||
r.mu.RUnlock()
|
||||
if mediaDir != "" {
|
||||
if err := os.MkdirAll(mediaDir, 0755); err == nil {
|
||||
fp := filepath.Join(mediaDir, reqID+mediaExt(acc.mime, acc.kind))
|
||||
if werr := os.WriteFile(fp, data, 0644); werr == nil {
|
||||
res["file"] = fp
|
||||
} else {
|
||||
log.Printf("[remotedevice] media write %s: %v", fp, werr)
|
||||
}
|
||||
} else {
|
||||
log.Printf("[remotedevice] media dir %s: %v", mediaDir, err)
|
||||
}
|
||||
}
|
||||
// 未配置落盘目录时保持旧行为:base64 内联返回(小体积数据仍可用)
|
||||
if _, hasFile := res["file"]; !hasFile {
|
||||
res["data_base64"] = base64.StdEncoding.EncodeToString(data)
|
||||
}
|
||||
r.SaveResult(reqID, res)
|
||||
r.deliverResult(reqID, res)
|
||||
|
||||
@ -1412,18 +1412,26 @@ func (h *Handler) handleChatEvents(w http.ResponseWriter, r *http.Request) {
|
||||
ticker := time.NewTicker(15 * time.Second)
|
||||
defer ticker.Stop()
|
||||
|
||||
// writeCh 不 close:Subscribe 回调闭包持有它,handler 退出后回调仍可能被
|
||||
// 总线异步触发,close 后再发送会 panic(send on closed channel,生产日志中
|
||||
// 单日数千次)。writer goroutine 通过 done 退出;发送侧 select on done 防泄漏。
|
||||
writeCh := make(chan string, 64)
|
||||
defer close(writeCh)
|
||||
|
||||
writerDone := make(chan struct{})
|
||||
go func() {
|
||||
defer func() {
|
||||
if r := recover(); r != nil {
|
||||
log.Printf("[SSE] writer panic: %v", r)
|
||||
}
|
||||
close(writerDone)
|
||||
}()
|
||||
for line := range writeCh {
|
||||
fmt.Fprintf(w, "%s\n", line)
|
||||
flusher.Flush()
|
||||
for {
|
||||
select {
|
||||
case line := <-writeCh:
|
||||
fmt.Fprintf(w, "%s\n", line)
|
||||
flusher.Flush()
|
||||
case <-done:
|
||||
return
|
||||
}
|
||||
}
|
||||
}()
|
||||
|
||||
@ -1501,6 +1509,7 @@ func (h *Handler) handleChatEvents(w http.ResponseWriter, r *http.Request) {
|
||||
for _, unsub := range unsubs {
|
||||
unsub()
|
||||
}
|
||||
<-writerDone // 等 writer 退出,保证 handler 返回后无残余写入
|
||||
}()
|
||||
for {
|
||||
select {
|
||||
|
||||
@ -52,6 +52,7 @@ func (m *mockPluginMgr) ReloadPlugins() (string, error) {
|
||||
return "reloaded", nil
|
||||
}
|
||||
func (m *mockPluginMgr) ReloadOne(name string) error { return nil }
|
||||
func (m *mockPluginMgr) StopAndUnload(name string) error { return nil }
|
||||
func (m *mockPluginMgr) PluginMetas() map[string]sdk.PluginMeta {
|
||||
return nil
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user