Files
HomeAgent/scripts/kernel-stress/devclient.py
JianFeeeee f0be8cbaa5 feat(remotedevice): 设备能力 outputch 化 —— 每设备一个 device/<id> 通道 + 设备指令类工具补授权闸
用户指出 remotedevice 的能力应当 outputch 化。查证后发现比"应当"更严重:
设备方向**根本没有出站实现**。

## 查到的三个缺口

1. `devicectlDevice` 一直声明 `OutputCapabilities() = CapStructured`(对外宣称可作输出目标),
   但 `Execute` 的 switch 里**没有 "output" 分支** ⇒ `output_send__devicectl` 必然拿到
   `unknown device tool output`,回模型"通过 [devicectl] 通道发送失败"。
2. 全插件没有 `RegisterOutputChannel`,也没有任何 EmitOutput/output_send 路径:
   设备方向只有「工具(请求-响应)」与「设备→agent 注入」,**agent → 设备是断的**
   (唯一的下行通道是个 HTTP 端点 `/api/v1/device/push`,不在 agent 的工具/通道模型里)。
3. 寻址是聚合的:所有设备共用一个名字 `devicectl`,没有 `device/<id>`;且设备 caps 只在
   插件内部软检查(`SupportsTool`),**绕过**了内核的 `AllowedOutputs` 授权闸 ——
   驻留子只要拿到 `device_ctl_cmdrun` 就能指挥**任意**设备。

## 按方向切分(不是一刀切)

**出站/消息类 → 每设备一个输出通道 `device/<id>`**(与入站同名):
- 上线注册、掉线注销(caps 由设备声明的 caps 映射:文本恒有;有屏→图/文件;
  speaker→音频;可跑命令(cmd/cmdrun/cmdresult)或未声明已知能力→全能力,与
  `deviceSupportsTool` 的旧设备兼容规则一致)。断连不注销会留下死通道骗模型。
- 于是自动获得:内核按 caps 在**发送前**拦(送图给纯文本音箱直接拒);
  `output_list_channels` 能列出设备;`AllowedOutputs` 可按设备收窄给驻留子。
- 上下线钩子用 `Registry.SetPresenceHandler`(**同步回调**)而不是既有的 `ChangeChan`
  (那是 select+default,缓冲满会丢事件;丢一次就留下死通道或漏注册)。
- `devicectl` 保留为聚合通道,并把它"声明了却不实现"的 output 补实:按
  `meta.device_id`(或 meta 就是设备 id / args.device_id)路由;缺省时返回**可执行**的
  报错(列出在线设备),而不是含糊失败。

**RPC 类保留为工具**(`devicedetect`/`screensee`/`computeruse`/`clipboard*`/
`device_ctl_status|cmdrun|cmdresult`):它们的返回值(图像/命令输出/状态)必须进模型
上下文,做成通道会丢掉这个语义。

**并给设备指令类工具补上同一道授权闸**(core/toolcall.go):`device_id` 指向的设备
必须是本 agent 被授权的 `device/<id>`。根 agent 默认完整授权 ⇒ 无行为变化;
驻留子收窄后,"拿到工具就能指挥任意设备"的缺口被堵上(新增 3 条 core 测试钉住)。

**设备端参考实现**(`internal/devicebridge/client/bridge.go` + waiter):新增 `op=push`
分发与 `OnPush` 回调(文本/结构化;二进制走既有 `cmd_speech_*` → `DataHandler`),
waiter 把它打到终端。

## 设计口径(用户当场纠偏,已写进代码注释与 harness README)

**主动转发只有 webui 与 cli 两个交互界面**(webui 订阅 EventAgentOutput 渲染气泡、
cli 用同步回程写回终端)。其它通道一律要求 agent **显式** `output_send__<通道>`。
我第一版给 remotedevice 加了 `EventAgentOutput` 订阅来自动回投设备——那是凭空造了
第三个转发者,违背"输出是 agent 的主动调用",已撤回(该测试一并删除)。

## 验收

- 单测:caps 映射词表;设备上线→注册通道(含同名 inputch)/掉线→注销;push 真落到
  WS 设备;聚合通道 `devicectl` 的寻址(无 device_id 报可执行错误、按 meta 投递、
  指定不存在设备报错);core 授权闸 3 例。
- 全量 `go test ./...` = 37 包 ok / 0 FAIL;`-race`(agent/plugins/plugin/devicebridge/sdk)干净。
- **真二进制端到端**(私有 netns + mock LLM + 真 WS 设备客户端 scripts/kernel-stress/devclient.py):
  设备上线 → 通道表出现 `device/pydev-1`(caps=[text file image audio structured],由
  `caps:["cmd"]` 映射)→ agent 经 `output_send__device/pydev-1` 主动发送 → 设备收到
  `{"op":"push","payload":"内核推给你的消息","type":"text"}` → 设备掉线 → 通道从表中消失。
2026-09-13 12:23:28 +08:00

161 lines
5.7 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

#!/usr/bin/env python3
"""最小远程设备客户端WebSocket无第三方依赖
用途:在**真内核**上验证"agent → 设备"这条出站链路 ——
`output_send__device/<id>` 应该以 `op=push` 的帧落到设备。
流程:握手 → hello自报 id/kind/caps→ bind → 打印收到的帧;
看到 push 就往 --out 文件里写一行(便于 shell 断言)。
注意:实例跑在私有 netns 里,设备网关的 127.0.0.1:9890 在 netns 内,
所以本脚本要用 nsenter 进同一个 netns 跑,例如:
nsenter -t <homed-pid> -n python3 devclient.py --port 9890 --token X --id pydev-1 --out /tmp/push.txt
"""
import argparse
import base64
import hashlib
import json
import os
import socket
import struct
import threading
import time
GUID = "258EAFA5-E914-47DA-95CA-C5AB0DC85B11"
class WS:
def __init__(self, host, port, path, timeout=30):
self.sock = socket.create_connection((host, port), timeout=timeout)
self.sock.settimeout(timeout)
self._handshake(host, port, path)
def _handshake(self, host, port, path):
key = base64.b64encode(os.urandom(16)).decode()
req = (
f"GET {path} HTTP/1.1\r\n"
f"Host: {host}:{port}\r\n"
"Upgrade: websocket\r\n"
"Connection: Upgrade\r\n"
f"Sec-WebSocket-Key: {key}\r\n"
"Sec-WebSocket-Version: 13\r\n\r\n"
)
self.sock.sendall(req.encode())
buf = b""
while b"\r\n\r\n" not in buf:
chunk = self.sock.recv(4096)
if not chunk:
raise RuntimeError("握手未完成: 连接关闭")
buf += chunk
head = buf.decode("utf-8", "replace")
if "101" not in head.split("\r\n")[0]:
raise RuntimeError("握手被拒: " + head.split("\r\n")[0])
expect = base64.b64encode(hashlib.sha1((key + GUID).encode()).digest()).decode()
if expect.lower() not in head.lower():
raise RuntimeError("Sec-WebSocket-Accept 校验失败")
self.buf = buf.split(b"\r\n\r\n", 1)[1]
# ---- 发送 ----
def send(self, opcode, payload=b""):
header = bytes([0x80 | opcode])
mask = os.urandom(4)
n = len(payload)
if n < 126:
header += bytes([0x80 | n])
elif n < 65536:
header += bytes([0x80 | 126]) + struct.pack(">H", n)
else:
header += bytes([0x80 | 127]) + struct.pack(">Q", n)
masked = bytes(b ^ mask[i % 4] for i, b in enumerate(payload))
self.sock.sendall(header + mask + masked)
def send_json(self, obj):
self.send(0x1, json.dumps(obj).encode())
# ---- 接收 ----
def _read(self, n):
while len(self.buf) < n:
chunk = self.sock.recv(65536)
if not chunk:
raise RuntimeError("连接关闭")
self.buf += chunk
out, self.buf = self.buf[:n], self.buf[n:]
return out
def recv_frame(self):
b0, b1 = self._read(2)
opcode = b0 & 0x0F
ln = b1 & 0x7F
if ln == 126:
ln = struct.unpack(">H", self._read(2))[0]
elif ln == 127:
ln = struct.unpack(">Q", self._read(8))[0]
masked = b1 & 0x80
mask = self._read(4) if masked else None
payload = self._read(ln)
if mask:
payload = bytes(b ^ mask[i % 4] for i, b in enumerate(payload))
return opcode, payload
def main():
ap = argparse.ArgumentParser()
ap.add_argument("--host", default="127.0.0.1")
ap.add_argument("--port", type=int, default=9890)
ap.add_argument("--token", required=True, help="设备接入 tokenconfig_remotedevice.ws_token")
ap.add_argument("--id", default="pydev-1")
ap.add_argument("--kind", default="computer")
ap.add_argument("--caps", default="cmd")
ap.add_argument("--seconds", type=float, default=25)
ap.add_argument("--out", default="", help="收到 push 时写一行到此文件")
args = ap.parse_args()
ws = WS(args.host, args.port, "/api/v1/device/ws?token=" + args.token)
ws.send_json({"op": "hello", "device": {
"device_id": args.id, "name": "python 设备", "kind": args.kind,
"caps": [c for c in args.caps.split(",") if c],
}})
op, payload = ws.recv_frame()
ack = json.loads(payload)
print("[devclient] hello_ack:", ack, flush=True)
ws.send_json({"op": "bind", "device_id": args.id, "token": args.token})
op, payload = ws.recv_frame()
bind = json.loads(payload)
print("[devclient] bind_ack:", bind, flush=True)
if not bind.get("ok"):
raise SystemExit("bind 被拒: %s" % bind)
deadline = time.time() + args.seconds
ws.sock.settimeout(1.0)
while time.time() < deadline:
try:
op, payload = ws.recv_frame()
except socket.timeout:
continue
except Exception as e: # 服务端关闭
print("[devclient] 连接结束:", e, flush=True)
break
if op == 0x9: # ping → pong
ws.send(0xA, payload)
continue
if op == 0x2:
print("[devclient] 收到二进制帧 %d 字节" % len(payload), flush=True)
continue
if op != 0x1:
continue
try:
m = json.loads(payload)
except Exception:
print("[devclient] 非 JSON 帧:", payload[:120], flush=True)
continue
print("[devclient] 收到帧:", json.dumps(m, ensure_ascii=False)[:300], flush=True)
if m.get("op") == "push" and args.out:
with open(args.out, "a") as f:
f.write(json.dumps(m, ensure_ascii=False) + "\n")
print("[devclient] ✅ 收到 pushagent → 设备链路通)", flush=True)
if __name__ == "__main__":
main()