跳到主要内容

数据截至 (上游 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 → 真实操作 │
└──────────────────────────────────┘ └──────────────────┘

上层的 SessionAgent、管线在两种形态里是同一份代码。差别只在 _init_context 里挂哪个 dispatcher(对比 ufo/module/sessions/session.py:95-104LocalCommandDispatcher,与 ufo/module/sessions/service_session.py:43-54WebSocketCommandDispatcher)。


6.3 AIP 分成五层

aip/__init__.py:16-25Architecture: 段把五层从上到下写在了模块文档里:

目录管什么
Messagesaip/messages.py强类型消息定义(Pydantic)
Protocolaip/protocol/协议逻辑:注册、任务执行、心跳、设备信息
Transportaip/transport/传输抽象(目前只有 WebSocket)
Endpointsaip/endpoints/三种端点:设备服务端、设备客户端、星座端
Resilienceaip/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 的口子。

一个错误处理的细节

发送失败时,如果错误信息里含 closednot 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_nameprocess_nameroot_nameresponse_idaction_results: List[Result]client_idprev_response_id
类型枚举ServerMessageTypeClientMessageType

注意 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.FAILEDExecutionResult(:600-717):

失败类型error_category
ConnectionError(执行途中掉线)connection_error
asyncio.TimeoutErrortimeout_error
其他异常general_error

方法文档明确写着 always returns, never raises。理由很直白:一个节点挂了不能让整张 DAG 崩,失败要作为数据流回编排器,让 ConstellationAgent 决定要不要改图补救。

finally 块里还固定做三件事:把设备设回 IDLE、发状态变更事件、检查队列里有没有下一个任务(:720-732)。


6.6 服务端形态

服务端在 ufo/server/:

文件角色
app.pyFastAPI 应用
ws/handler.pyWebSocket 处理器(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.pyWebSocket 角色伪装
test_ws_shared_handler_isolation.py共享 handler 的会话串扰
test_mobile_mcp_auth.py移动端 MCP 服务器鉴权

最后一条对应 config/ufo/mcp.yaml:157MobileAgentauth: "${UFO_MCP_API_KEY}"——HTTP 型 MCP 服务器是需要带鉴权的。


6.7 跨平台设备端

设备端不必是 Windows。HTTP 型 MCP 服务器有三个(ufo/client/mcp/http_servers/):

文件面向默认端口
linux_mcp_server.pyLinux(shell)8010
mobile_mcp_server.pyAndroid8020 / 8021
hardware_mcp_server.py硬件8006

对应的 agent 侧策略在 ufo/agents/processors/strategies/linux_agent_strategy.pymobile_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.pyServerMessageClientMessageCommandResultClientType
协议核心aip/protocol/base.pyAIPProtocolsend_messagereceive_message
任务执行子协议aip/protocol/task_execution.pyTaskExecutionProtocolsend_commandsend_command_results
注册 / 心跳 / 设备信息aip/protocol/registration.pyheartbeat.pydevice_info.py
传输层aip/transport/websocket.py(WebSocket 传输实现)
三种端点aip/endpoints/DeviceServerEndpointDeviceClientEndpointConstellationEndpoint
韧性组件aip/resilience/heartbeat_manager.pyreconnection.pytimeout.py
远程 dispatcherufo/module/dispatcher.pyWebSocketCommandDispatchermake_server_responseset_result
设备端执行器ufo/client/ufo_client.pyUFOClient.execute_stepexecute_actions
设备信息采集ufo/client/device_info_provider.pyDeviceInfoProvider.collect_system_info
设备生命周期galaxy/client/device_manager.pyConstellationDeviceManagerassign_task_to_device_execute_task_on_device
服务端 WebSocketufo/server/ws/handler.pyUFOWebSocketHandlerconnecthandle_task_request
会话归属校验ufo/server/services/session_manager.pySessionManagerSessionOwnershipErroris_authorized_result_sender
跨平台 MCP 服务器ufo/client/mcp/http_servers/linux_mcp_server.pymobile_mcp_server.pyhardware_mcp_server.py
安全回归测试tests/security/test_session_id_reuse.pytest_task_name_traversal.pytest_ws_role_spoof.py