mirror of
https://gitcode.com/JianFeeeee/HomeAgent.git
synced 2026-10-03 15:53:56 +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(3ceb69f 接口变更漏改测试)。 3. GUI camerasue 平台分支 ffmpeg 参数原硬编码 Linux v4l2(/dev/video0),Windows 上必然失败。 现按平台探测:win32=dshow(枚举设备名取第一个视频设备)、 darwin=avfoundation、linux=v4l2;录像编码 Windows 交给 mp4 muxer 默认。
This commit is contained in:
@ -1380,7 +1380,8 @@ function executeHomeagentCmd(capability, reqId) {
|
|||||||
.split(/[ >\n]/)[0];
|
.split(/[ >\n]/)[0];
|
||||||
switch (name) {
|
switch (name) {
|
||||||
case "camerasue": {
|
case "camerasue": {
|
||||||
// 摄像头:camerasue=抓拍单张;camerasue <秒>=录制 N 秒视频,返回 base64
|
// 摄像头:camerasue=抓拍单张;camerasue <秒>=录制 N 秒视频
|
||||||
|
// 平台分支:Windows=dshow(设备名自动探测),macOS=avfoundation,Linux=v4l2
|
||||||
const argStr = String(capability || "")
|
const argStr = String(capability || "")
|
||||||
.replace(/^camerasue/, "")
|
.replace(/^camerasue/, "")
|
||||||
.trim();
|
.trim();
|
||||||
@ -1390,24 +1391,54 @@ function executeHomeagentCmd(capability, reqId) {
|
|||||||
const os = require("os");
|
const os = require("os");
|
||||||
const path = require("path");
|
const path = require("path");
|
||||||
const fs = require("fs");
|
const fs = require("fs");
|
||||||
|
// 探测平台可用的 ffmpeg 输入参数(缓存结果避免重复探测)
|
||||||
|
let camInput = null;
|
||||||
|
function resolveCameraInput(cb) {
|
||||||
|
if (camInput) return cb(camInput);
|
||||||
|
const plat = process.platform;
|
||||||
|
if (plat === "win32") {
|
||||||
|
// dshow:先枚举设备名取第一个视频设备
|
||||||
|
cp.execFile(
|
||||||
|
"ffmpeg",
|
||||||
|
["-hide_banner", "-list_devices", "true", "-f", "dshow", "-i", "video= dummy"],
|
||||||
|
{ timeout: 8000 },
|
||||||
|
(err, _so, se) => {
|
||||||
|
const out = String(se || "");
|
||||||
|
const m = out.match(/"([^"]+)"\s*\((?:video|默认)|"([^"]+)"[\s\S]{0,200}?\(video/)
|
||||||
|
|| out.match(/"([^"]+)"[^\n]*\(video/i);
|
||||||
|
const name = m ? (m[1] || m[2]) : null;
|
||||||
|
if (name) {
|
||||||
|
camInput = { pre: ["-f", "dshow", "-i", "video=" + name] };
|
||||||
|
} else {
|
||||||
|
camInput = { pre: ["-f", "dshow", "-i", "video=USB Camera"] }; // 常见默认名兑底
|
||||||
|
}
|
||||||
|
cb(camInput);
|
||||||
|
},
|
||||||
|
);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
if (plat === "darwin") {
|
||||||
|
camInput = { pre: ["-f", "avfoundation", "-i", "0:0"] }; // 默认摄像头
|
||||||
|
return cb(camInput);
|
||||||
|
}
|
||||||
|
camInput = { pre: ["-f", "v4l2", "-i", "/dev/video0"] }; // Linux
|
||||||
|
return cb(camInput);
|
||||||
|
}
|
||||||
|
resolveCameraInput((cam) => {
|
||||||
if (isVideo) {
|
if (isVideo) {
|
||||||
// 录像:ffmpeg 录 N 秒 mp4 到临时文件
|
// 录像:ffmpeg 录 N 秒 mp4 到临时文件
|
||||||
const outFile = path.join(os.tmpdir(), "ha_cam_" + Date.now() + ".mp4");
|
const outFile = path.join(os.tmpdir(), "ha_cam_" + Date.now() + ".mp4");
|
||||||
const args = [
|
const args = [
|
||||||
"-f",
|
...cam.pre,
|
||||||
"v4l2",
|
|
||||||
"-i",
|
|
||||||
"/dev/video0",
|
|
||||||
"-t",
|
"-t",
|
||||||
String(durMatch),
|
String(durMatch),
|
||||||
"-pix_fmt",
|
"-pix_fmt",
|
||||||
"yuv420p",
|
"yuv420p",
|
||||||
"-c:v",
|
|
||||||
"libx264",
|
|
||||||
"-f",
|
|
||||||
"mp4",
|
|
||||||
outFile,
|
|
||||||
];
|
];
|
||||||
|
if (process.platform !== "win32") {
|
||||||
|
args.push("-c:v", "libx264"); // Windows dshow→mp4 由扩展名驱动原生编码器
|
||||||
|
}
|
||||||
|
args.push("-f", "mp4", outFile);
|
||||||
cp.execFile(
|
cp.execFile(
|
||||||
"ffmpeg",
|
"ffmpeg",
|
||||||
args,
|
args,
|
||||||
@ -1446,19 +1477,16 @@ function executeHomeagentCmd(capability, reqId) {
|
|||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
// 抓拍单张 jpeg
|
// 抓拍单张 jpeg
|
||||||
const args = [
|
const args = [
|
||||||
"-f",
|
...cam.pre,
|
||||||
"v4l2",
|
"-frames:v",
|
||||||
"-i",
|
"1",
|
||||||
"/dev/video0",
|
"-f",
|
||||||
"-frames:v",
|
"image2pipe",
|
||||||
"1",
|
"-vcodec",
|
||||||
"-f",
|
"mjpeg",
|
||||||
"image2pipe",
|
"pipe:1",
|
||||||
"-vcodec",
|
];
|
||||||
"mjpeg",
|
|
||||||
"pipe:1",
|
|
||||||
];
|
|
||||||
cp.execFile(
|
cp.execFile(
|
||||||
"ffmpeg",
|
"ffmpeg",
|
||||||
args,
|
args,
|
||||||
@ -1483,6 +1511,7 @@ function executeHomeagentCmd(capability, reqId) {
|
|||||||
);
|
);
|
||||||
},
|
},
|
||||||
);
|
);
|
||||||
|
}); // resolveCameraInput 回调闭合
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
case "screensue": {
|
case "screensue": {
|
||||||
|
|||||||
@ -9,6 +9,8 @@ import (
|
|||||||
"net"
|
"net"
|
||||||
"net/http"
|
"net/http"
|
||||||
"net/http/httptest"
|
"net/http/httptest"
|
||||||
|
"os"
|
||||||
|
"path/filepath"
|
||||||
"strings"
|
"strings"
|
||||||
"sync"
|
"sync"
|
||||||
"testing"
|
"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 协议)=====
|
// ===== 端到端:PushData 下发音频(网关→设备 cmd_speech 协议)=====
|
||||||
|
|
||||||
func TestWSPushDataAudio(t *testing.T) {
|
func TestWSPushDataAudio(t *testing.T) {
|
||||||
|
|||||||
@ -8,6 +8,7 @@ import (
|
|||||||
"fmt"
|
"fmt"
|
||||||
"log"
|
"log"
|
||||||
"net/http"
|
"net/http"
|
||||||
|
"path/filepath"
|
||||||
"strings"
|
"strings"
|
||||||
"sync"
|
"sync"
|
||||||
"time"
|
"time"
|
||||||
@ -104,6 +105,15 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error {
|
|||||||
log.Printf("[remotedevice] register devicectl channel: %v", err)
|
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 注入 ----------------
|
// ---- 设备主动上报事件 → agent 注入 ----------------
|
||||||
// 摄像头发现异常/传感器报警等场景:设备经 WS op=event 上报,
|
// 摄像头发现异常/传感器报警等场景:设备经 WS op=event 上报,
|
||||||
// 插件将其格式化为文本经 SDK InjectText 异步注入 agent(source=device/{id},
|
// 插件将其格式化为文本经 SDK InjectText 异步注入 agent(source=device/{id},
|
||||||
|
|||||||
@ -11,6 +11,8 @@ import (
|
|||||||
"log"
|
"log"
|
||||||
"net"
|
"net"
|
||||||
"net/http"
|
"net/http"
|
||||||
|
"os"
|
||||||
|
"path/filepath"
|
||||||
"strings"
|
"strings"
|
||||||
"sync"
|
"sync"
|
||||||
"time"
|
"time"
|
||||||
@ -49,6 +51,7 @@ type Registry struct {
|
|||||||
acceptFn func(token string) bool
|
acceptFn func(token string) bool
|
||||||
cmdPending map[string]chan map[string]interface{} // reqID -> 结果 channel
|
cmdPending map[string]chan map[string]interface{} // reqID -> 结果 channel
|
||||||
results map[string]resultEntry // reqID -> 已留档结果
|
results map[string]resultEntry // reqID -> 已留档结果
|
||||||
|
mediaDir string // 设备回传媒体落盘目录;空则退化为 base64 内联
|
||||||
}
|
}
|
||||||
|
|
||||||
// resultEntry 保存一次 cmdrun 的结果(供 device_ctl_cmdresult 查询)。
|
// resultEntry 保存一次 cmdrun 的结果(供 device_ctl_cmdresult 查询)。
|
||||||
@ -120,6 +123,43 @@ func deviceSupportsTool(caps []string, tool string) bool {
|
|||||||
return !hasKnown // 未声明任何已知能力 → 全能力兼容
|
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 {
|
func NewRegistry() *Registry {
|
||||||
return &Registry{
|
return &Registry{
|
||||||
devices: make(map[string]*DeviceMeta),
|
devices: make(map[string]*DeviceMeta),
|
||||||
@ -743,8 +783,28 @@ func (r *Registry) handleWS(conn net.Conn, rw *bufio.ReadWriter) {
|
|||||||
"mime": acc.mime,
|
"mime": acc.mime,
|
||||||
"size": len(data),
|
"size": len(data),
|
||||||
"expected": acc.total,
|
"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.SaveResult(reqID, res)
|
||||||
r.deliverResult(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)
|
ticker := time.NewTicker(15 * time.Second)
|
||||||
defer ticker.Stop()
|
defer ticker.Stop()
|
||||||
|
|
||||||
|
// writeCh 不 close:Subscribe 回调闭包持有它,handler 退出后回调仍可能被
|
||||||
|
// 总线异步触发,close 后再发送会 panic(send on closed channel,生产日志中
|
||||||
|
// 单日数千次)。writer goroutine 通过 done 退出;发送侧 select on done 防泄漏。
|
||||||
writeCh := make(chan string, 64)
|
writeCh := make(chan string, 64)
|
||||||
defer close(writeCh)
|
writerDone := make(chan struct{})
|
||||||
|
|
||||||
go func() {
|
go func() {
|
||||||
defer func() {
|
defer func() {
|
||||||
if r := recover(); r != nil {
|
if r := recover(); r != nil {
|
||||||
log.Printf("[SSE] writer panic: %v", r)
|
log.Printf("[SSE] writer panic: %v", r)
|
||||||
}
|
}
|
||||||
|
close(writerDone)
|
||||||
}()
|
}()
|
||||||
for line := range writeCh {
|
for {
|
||||||
fmt.Fprintf(w, "%s\n", line)
|
select {
|
||||||
flusher.Flush()
|
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 {
|
for _, unsub := range unsubs {
|
||||||
unsub()
|
unsub()
|
||||||
}
|
}
|
||||||
|
<-writerDone // 等 writer 退出,保证 handler 返回后无残余写入
|
||||||
}()
|
}()
|
||||||
for {
|
for {
|
||||||
select {
|
select {
|
||||||
|
|||||||
@ -52,6 +52,7 @@ func (m *mockPluginMgr) ReloadPlugins() (string, error) {
|
|||||||
return "reloaded", nil
|
return "reloaded", nil
|
||||||
}
|
}
|
||||||
func (m *mockPluginMgr) ReloadOne(name string) error { return 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 {
|
func (m *mockPluginMgr) PluginMetas() map[string]sdk.PluginMeta {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user