数据截至 (上游 commit 96983c73ed09)
第 6 章 · AIP 协议与设备层
这章讲什么: 同一份 agent 代码,怎么做到「跑在本机就直接点鼠标,跑在服务端就远程驱动另一台机器」。答案是一个接口 + 一层协议。
6.1 要解决的小问题
第 4 章讲的动作链路,末端是 pywinauto。但 Galaxy 的场景里,大脑跑在编排机上,手脚在另外三台机器上。
如果为远程场景另写一套 agent,就要维护两份几乎相同的代码。UFO 的做法是:把「动作往哪送」这件事收进一个接口,上层完全不知道差别。
6.2 两种部署形态
怎么读这张图:上下两块是两种部署;方框内是一个进程,虚线是跨进程。
【本地形态】 python -m ufo
┌────────────────────────────────────────────────────────┐
│ Session → Agent → 管线 → LocalDispatcher │
│ └─▶ Computer → MCP → pywinauto│
└────────────────────────────────────────────────────────┘
【远程形态】 服务端 + 设备端
┌───────────── 服务端 ─────────────┐ ┌───── 设备端 ─────┐
│ Session → Agent → 管线 │ │ UFOClient │
│ → WebSocketDispatcher ┄┄┄┄┼┄AIP┄┄▶│ → CommandRouter │
│ ◀┄┄┄┄┄ Result ┄┄┄┄┄┄┄┄┄┄┄┄┼┄┄┄┄┄┄┄│ → MCP → 真实操作 │
└──────────────────────────────────┘ └──────────────────┘
上层的 Session、Agent、管线在两种形态里是同一份代码。差别只在 _init_context 里挂哪个 dispatcher(对比 ufo/module/sessions/session.py:95-104 挂 LocalCommandDispatcher,与 ufo/module/sessions/service_session.py:43-54 挂 WebSocketCommandDispatcher)。
6.3 AIP 分成五层
aip/__init__.py:16-25 的 Architecture: 段把五层从上到下写在了模块文档里:
| 层 | 目录 | 管什么 |
|---|---|---|
| Messages | aip/messages.py | 强类型消息定义(Pydantic) |
| Protocol | aip/protocol/ | 协议逻辑:注册、任务执行、心跳、设备信息 |
| Transport | aip/transport/ | 传输抽象(目前只有 WebSocket) |
| Endpoints | aip/endpoints/ | 三种端点:设备服务端、设备客户端、星座端 |
| Resilience | aip/resilience/ | 重连、心跳、超时 |
核心类只做序列化和中间件
AIPProtocol(aip/protocol/base.py:22)本身很薄:
# 示意,非源码 —— 演示协议层的职责边界
async def send_message(self, msg):
for mw in self.middleware_chain:
msg = await mw.process_outgoing(msg) # 出站中间件
data = msg.model_dump_json().encode("utf-8") # Pydantic 序列化
await self.transport.send(data) # 交给传输层
传输无关是刻意的。 类注释写着 transport-agnostic and works with any Transport implementation(:32),留了将来换 HTTP/3 或 gRPC 的口子。
一个错误处理的细节
发送失败时,如果错误信息里含 closed 或 not connected,只打 DEBUG 日志而不是 ERROR(aip/protocol/base.py:81-90)。理由写在注释里:正常断连时刷一堆 ERROR 会吓人。这是很实际的可观测性打磨。
6.4 消息类型
两个方向的消息各有一个类:
ServerMessage(aip/messages.py:299) | ClientMessage(aip/messages.py:343) | |
|---|---|---|
| 谁发 | 服务端(大脑) | 设备端(手脚) |
| 关键字段 | actions: List[Command]、agent_name、process_name、root_name、response_id | action_results: List[Result]、client_id、prev_response_id |
| 类型枚举 | ServerMessageType | ClientMessageType |
注意 ServerMessage 里带着 agent_name / process_name / root_name 三件套——设备端要靠这三个值去 ComputerManager 里挑对应的工具集(UFOClient.execute_step,ufo/client/ufo_client.py:53-67)。这跟第 4 章本地形态里的三元组缓存键是同一套东西。
配对靠 response_id:服务端发出时记进 pending 表,设备端回包时带上,set_result 把结果塞进对应的 future(ufo/module/dispatcher.py:206-252)。
6.5 设备生命周期
ConstellationDeviceManager(galaxy/client/device_manager.py:33)管一整套:
| 阶段 | 方法 | 干什么 |
|---|---|---|
| 注册 | register_device(:174) | 把设备写进注册表 |
| 连接 | connect_device(:221) | 建 WebSocket |
| 派活 | assign_task_to_device(:543) | 忙就排队,闲就直接执行 |
| 断连 | _handle_device_disconnection(:380) | 清理 + 安排重连 |
| 重连 | _reconnect_device(:461) | 退避重试 |
忙闲与排队
派活时先看设备状态(:543-598):
# 示意,非源码 —— 演示忙闲分流
if device_registry.is_device_busy(device_id):
future = task_queue_manager.enqueue_task(device_id, task_request) # 排队
return await future # 等到轮到自己
else:
return await self._execute_task_on_device(device_id, task_request) # 直接跑
一台设备同时只跑一个任务。 这很合理——GUI 是共享的物理资源,两个任务同时点鼠标必然互相破坏。DAG 的并行度因此受限于设备数,而不是节点数。
失败不抛异常,只回结果
_execute_task_on_device 把三类失败都转成带 TaskStatus.FAILED 的 ExecutionResult(:600-717):
| 失败类型 | error_category |
|---|---|
ConnectionError(执行途中掉线) | connection_error |
asyncio.TimeoutError | timeout_error |
| 其他异常 | general_error |
方法文档明确写着 always returns, never raises。理由很直白:一个节点挂了不能让整张 DAG 崩,失败要作为数据流回编排器,让 ConstellationAgent 决定要不要改图补救。
finally 块里还固定做三件事:把设备设回 IDLE、发状态变更事件、检查队列里有没有下一个任务(:720-732)。
6.6 服务端形态
服务端在 ufo/server/:
| 文件 | 角色 |
|---|---|
app.py | FastAPI 应用 |
ws/handler.py | WebSocket 处理器(UFOWebSocketHandler,:62) |
services/session_manager.py | 会话管理(SessionManager,:52) |
services/client_connection_manager.py | 连接管理 |
UFOWebSocketHandler.connect 在握手时就要求一条注册消息,解析出 client_type 才放行(ufo/server/ws/handler.py:100-190)。星座客户端还要额外走一次 _validate_constellation_client(:219)。
会话归属检查
SessionManager 里有一个专门的异常类 SessionOwnershipError(ufo/server/services/session_manager.py:20),以及 register_result_sender / is_authorized_result_sender 这对方法(:231、:245)。
这不是过度设计——仓库的 tests/security/test_session_id_reuse.py 开头写明了它修的是一个已公布的安全问题:早期版本只要客户端猜中 session_id 就能复用别人的内存会话对象。
安全测试清单
tests/security/ 下五个文件,每个对应一类攻击面:
| 文件 | 防什么 |
|---|---|
test_session_id_reuse.py | 跨客户端复用 session_id |
test_task_name_traversal.py | 任务名做路径穿越(对应第 1 章的 sanitize_task_name) |
test_ws_role_spoof.py | WebSocket 角色伪装 |
test_ws_shared_handler_isolation.py | 共享 handler 的会话串扰 |
test_mobile_mcp_auth.py | 移动端 MCP 服务器鉴权 |
最后一条对应 config/ufo/mcp.yaml:157 里 MobileAgent 的 auth: "${UFO_MCP_API_KEY}"——HTTP 型 MCP 服务器是需要带鉴权的。
6.7 跨平台设备端
设备端不必是 Windows。HTTP 型 MCP 服务器有三个(ufo/client/mcp/http_servers/):
| 文件 | 面向 | 默认端口 |
|---|---|---|
linux_mcp_server.py | Linux(shell) | 8010 |
mobile_mcp_server.py | Android | 8020 / 8021 |
hardware_mcp_server.py | 硬件 | 8006 |
对应的 agent 侧策略在 ufo/agents/processors/strategies/linux_agent_strategy.py 和 mobile_agent_strategy.py,状态机在 ufo/agents/states/linux_agent_state.py / mobile_agent_state.py,会话在 ufo/module/sessions/linux_session.py / mobile_session.py。
骨架一样,只换策略和工具集——这是第 2 章那套设计的直接兑现。
设备自己能上报硬件信息:DeviceInfoProvider.collect_system_info 收集平台、OS 版本、CPU 核数、内存、主机名、IP 和特性列表(ufo/client/device_info_provider.py:64-100)。这些信息进了 Galaxy 的设备清单,再进 prompt。
6.8 代码地图
| 主题 | 文件路径 | 符号名 |
|---|---|---|
| 协议分层总述 | aip/__init__.py | (模块 docstring 描述五层结构) |
| 消息定义 | aip/messages.py | ServerMessage、ClientMessage、Command、Result、ClientType |
| 协议核心 | aip/protocol/base.py | AIPProtocol、send_message、receive_message |
| 任务执行子协议 | aip/protocol/task_execution.py | TaskExecutionProtocol、send_command、send_command_results |
| 注册 / 心跳 / 设备信息 | aip/protocol/ | registration.py、heartbeat.py、device_info.py |
| 传输层 | aip/transport/websocket.py | (WebSocket 传输实现) |
| 三种端点 | aip/endpoints/ | DeviceServerEndpoint、DeviceClientEndpoint、ConstellationEndpoint |
| 韧性组件 | aip/resilience/ | heartbeat_manager.py、reconnection.py、timeout.py |
| 远程 dispatcher | ufo/module/dispatcher.py | WebSocketCommandDispatcher、make_server_response、set_result |
| 设备端执行器 | ufo/client/ufo_client.py | UFOClient.execute_step、execute_actions |
| 设备信息采集 | ufo/client/device_info_provider.py | DeviceInfoProvider.collect_system_info |
| 设备生命周期 | galaxy/client/device_manager.py | ConstellationDeviceManager、assign_task_to_device、_execute_task_on_device |
| 服务端 WebSocket | ufo/server/ws/handler.py | UFOWebSocketHandler、connect、handle_task_request |
| 会话归属校验 | ufo/server/services/session_manager.py | SessionManager、SessionOwnershipError、is_authorized_result_sender |
| 跨平台 MCP 服务器 | ufo/client/mcp/http_servers/ | linux_mcp_server.py、mobile_mcp_server.py、hardware_mcp_server.py |
| 安全回归测试 | tests/security/ | test_session_id_reuse.py、test_task_name_traversal.py、test_ws_role_spoof.py |