feat: remotedevice 二进制分块协议(录像回传+音频下发) + 工具描述增强

回应 GUI 6ea3fde 需服务端配套项([阻塞]两项 + [中]一项):

1. readFrame 支持 0x2 二进制帧(之前遇 0x2 直接断连):
   - 返回 opcode, 二进制帧上限放宽至 8MB
2. handleWS 二进制聚合协议(设备→网关, 录像回传):
   - cmd_data_start 开启按 req_id 聚合 → 0x2 帧追加 → cmd_data_end 聚合完成
   - 结果存 cmdresult(data_base64 字段), device_ctl_cmdresult 可取回
   - 超限防护(声明 total×2 或硬上限 64MB)
3. PushData 下发协议(网关→设备, 音频/TTS):
   - cmd_speech_start → 0x2 分块(8KB) → cmd_speech_end
   - GUI 端已实现接收侧(speakeruse 播放链路打通)
4. device_ctl_cmdrun 工具描述补齐:
   - screensue 带参示例(<内容>/<秒> <内容>)
   - camerasue 录像说明(cmdresult 含 data_base64)
   - speakeruse 文字朗读

测试: 最小 WS 客户端端到端验证二进制分块上传/PushData 下发/离线报错
This commit is contained in:
JianFeeeee
2026-08-21 10:53:38 +08:00
parent 9692eeecd3
commit e42f043bea
3 changed files with 457 additions and 20 deletions

View File

@ -260,6 +260,45 @@ func (r *Registry) PushCmd(deviceID, reqID, command, cmdType string) error {
})
}
// PushData 向设备分块下发二进制数据(网关→设备,如 TTS 音频)。
// 协议(与 GUI 设备桥协商):
// 文本帧 cmd_speech_start {op, req_id, kind, mime, total} → N 个二进制帧(0x2, ≤8KB) → 文本帧 cmd_speech_end {op, req_id}
// kind 为语义标记(如 speech),mime 为数据 MIME 类型。设备聚合后按自身能力处理(播放等)。
func (r *Registry) PushData(deviceID, reqID, kind, mime string, data []byte) error {
r.mu.RLock()
c, ok := r.conns[deviceID]
r.mu.RUnlock()
if !ok {
return fmt.Errorf("device %s not online", deviceID)
}
if err := writeText(c.w, mustJSON(map[string]interface{}{
"op": "cmd_speech_start",
"req_id": reqID,
"kind": kind,
"mime": mime,
"total": len(data),
})); err != nil {
return fmt.Errorf("push data start: %w", err)
}
const chunkSize = 8192
for off := 0; off < len(data); off += chunkSize {
end := off + chunkSize
if end > len(data) {
end = len(data)
}
if err := writeBinary(c.w, data[off:end]); err != nil {
return fmt.Errorf("push data chunk: %w", err)
}
}
if err := writeText(c.w, mustJSON(map[string]interface{}{
"op": "cmd_speech_end",
"req_id": reqID,
})); err != nil {
return fmt.Errorf("push data end: %w", err)
}
return nil
}
// AwaitResult 等待某请求的结果(带超时)。
func (r *Registry) AwaitResult(reqID string, timeout time.Duration) (map[string]interface{}, error) {
ch := make(chan map[string]interface{}, 1)
@ -342,43 +381,50 @@ func httpUpgrade(w http.ResponseWriter, r *http.Request) (net.Conn, *bufio.ReadW
return conn, rw, nil
}
func readFrame(r *bufio.Reader) ([]byte, bool, error) {
// readFrame 读取一个 WS 帧。返回 (payload, isClose, err)。
// opcode: 0x1 文本 / 0x2 二进制(设备→网关大体积数据分块,如录像回传)。
func readFrame(r *bufio.Reader) ([]byte, bool, byte, error) {
b0, err := r.ReadByte()
if err != nil {
return nil, true, err
return nil, true, 0, err
}
opcode := b0 & 0x0f
b1, err := r.ReadByte()
if err != nil {
return nil, true, err
return nil, true, opcode, err
}
masked := b1&0x80 != 0
length := uint64(b1 & 0x7f)
if length == 126 {
var ext [2]byte
if _, err := io.ReadFull(r, ext[:]); err != nil {
return nil, true, err
return nil, true, opcode, err
}
length = uint64(binary.BigEndian.Uint16(ext[:]))
} else if length == 127 {
var ext [8]byte
if _, err := io.ReadFull(r, ext[:]); err != nil {
return nil, true, err
return nil, true, opcode, err
}
length = binary.BigEndian.Uint64(ext[:])
}
if length > 1<<20 {
return nil, true, fmt.Errorf("frame too large")
// 二进制帧允许更大(录像分块聚合,单帧仍限 8MB 防滥用)
maxFrame := uint64(1 << 20)
if opcode == 0x2 {
maxFrame = 8 << 20
}
if length > maxFrame {
return nil, true, opcode, fmt.Errorf("frame too large")
}
var maskKey [4]byte
if masked {
if _, err := io.ReadFull(r, maskKey[:]); err != nil {
return nil, true, err
return nil, true, opcode, err
}
}
payload := make([]byte, length)
if _, err := io.ReadFull(r, payload); err != nil {
return nil, true, err
return nil, true, opcode, err
}
if masked {
for i := range payload {
@ -386,23 +432,32 @@ func readFrame(r *bufio.Reader) ([]byte, bool, error) {
}
}
switch opcode {
case 0x1:
return payload, false, nil
case 0x1, 0x2:
return payload, false, opcode, nil
case 0x8:
return nil, true, nil
return nil, true, opcode, nil
case 0xa:
return nil, false, nil
return nil, false, opcode, nil
case 0x9:
return nil, false, errPing
return nil, false, opcode, errPing
default:
return nil, false, fmt.Errorf("unsupported opcode %x", opcode)
return nil, false, opcode, fmt.Errorf("unsupported opcode %x", opcode)
}
}
var errPing = fmt.Errorf("ping")
func writeText(w *bufio.Writer, payload []byte) error {
if err := writeFrameHeader(w, 0x1, len(payload)); err != nil {
return writeFrame(w, 0x1, payload)
}
// writeBinary 发送 WS 二进制帧(0x2):网关→设备大体积数据(如 TTS 音频)分块下发。
func writeBinary(w *bufio.Writer, payload []byte) error {
return writeFrame(w, 0x2, payload)
}
func writeFrame(w *bufio.Writer, opcode byte, payload []byte) error {
if err := writeFrameHeader(w, opcode, len(payload)); err != nil {
return err
}
if _, err := w.Write(payload); err != nil {
@ -491,8 +546,12 @@ func (r *Registry) handleWS(conn net.Conn, rw *bufio.ReadWriter) {
}
}()
// 二进制分块聚合状态(设备→网关,如录像回传):
// cmd_data_start 开启 → 0x2 帧追加 → cmd_data_end 聚合存入 cmdresult
var dataAccum *dataAccumulator
for {
payload, isClose, err := readFrame(rw.Reader)
payload, isClose, opcode, err := readFrame(rw.Reader)
if err != nil {
if err == errPing {
if werr := writePong(rw.Writer); werr != nil {
@ -505,6 +564,23 @@ func (r *Registry) handleWS(conn net.Conn, rw *bufio.ReadWriter) {
if isClose {
return
}
if opcode == 0x2 {
// 二进制帧:处于聚合状态时追加数据块,否则忽略
if dataAccum != nil {
dataAccum.chunks = append(dataAccum.chunks, payload)
dataAccum.got += len(payload)
// 防滥用:超出声明 total 的 2 倍或硬上限 64MB 时放弃聚合
limit := int64(dataAccum.total)*2 + 1024
if limit < 64<<20 {
limit = 64 << 20
}
if int64(dataAccum.got) > limit {
log.Printf("[remotedevice] data accumulation exceeded limit for req %s, dropped", dataAccum.reqID)
dataAccum = nil
}
}
continue
}
var msg map[string]interface{}
if err := json.Unmarshal(payload, &msg); err != nil {
continue
@ -563,10 +639,70 @@ func (r *Registry) handleWS(conn net.Conn, rw *bufio.ReadWriter) {
if reqID != "" {
r.deliverResult(reqID, msg)
}
case "cmd_data_start":
reqID, _ := msg["req_id"].(string)
if reqID == "" {
continue
}
total, _ := msg["total"].(float64)
kind, _ := msg["kind"].(string)
mime, _ := msg["mime"].(string)
dataAccum = &dataAccumulator{
reqID: reqID,
kind: kind,
mime: mime,
total: int(total),
}
case "cmd_data_end":
reqID, _ := msg["req_id"].(string)
status, _ := msg["status"].(string)
if dataAccum == nil || dataAccum.reqID != reqID {
continue
}
acc := dataAccum
dataAccum = nil
if status != "ok" {
r.SaveResult(reqID, map[string]interface{}{
"op": "cmd_result", "req_id": reqID, "status": "error",
"error": "device reported transfer failure",
})
r.deliverResult(reqID, map[string]interface{}{
"op": "cmd_result", "req_id": reqID, "status": "error",
"error": "device reported transfer failure",
})
continue
}
data := make([]byte, 0, acc.got)
for _, c := range acc.chunks {
data = append(data, c...)
}
res := map[string]interface{}{
"op": "cmd_result",
"req_id": reqID,
"status": "ok",
"kind": acc.kind,
"mime": acc.mime,
"size": len(data),
"expected": acc.total,
// base64 编码完整二进制(录像 mp4 等),供 agent/上层取回后解码使用
"data_base64": base64.StdEncoding.EncodeToString(data),
}
r.SaveResult(reqID, res)
r.deliverResult(reqID, res)
}
}
}
// dataAccumulator 聚合设备→网关的二进制分块传输(如录像回传)。
type dataAccumulator struct {
reqID string
kind string
mime string
total int
chunks [][]byte
got int
}
func metaFromMsg(msg map[string]interface{}) DeviceMeta {
var meta DeviceMeta
if d, ok := msg["device"].(map[string]interface{}); ok {