UFO 项目 Agent Interaction Protocol(AIP)深度解析:五层架构、WebSocket 消息协议与弹性编排机制
AIP(Agent Interaction Protocol)是 UFO 项目中连接 ConstellationClient、设备 Agent 服务与设备客户端的事件驱动通信协议,为持续演化的 DAG 工作流与长时推理循环提供正确的编排底座。本文将围绕 AIP 的六大设计目标(G1–G6)、五层架构、消息契约、弹性连接与扩展机制展开系统讲解,并结合仓库源码与测试用例给出可验证的实现细节,帮助读者掌握在 UFO 生态中搭建、扩展与运维 AIP 通信层的完整方案。
1. 为什么需要 AIP:从短生命周期 HTTP 协调到持久化事件驱动
UFO 的编排模型要求通信基座在 DAG 持续演化(continuous DAG evolution)、动态 Agent 参与(dynamic agent participation) 与细粒度事件传播(fine-grained event propagation) 三种场景下保持正确性。而传统基于 HTTP 的协调方案(如 A2A、ACP)假设交互是短生命周期、无状态的,在实际运行中会暴露出握手开销、能力视图过期、任务中途部分失败时恢复脆弱等缺陷,难以支撑 UFO 中持续演化的工作流与长时推理循环。
AIP 由此被设计为 UFO 的"神经系统"(nervous system),以轻量但能容忍演化的协议形态满足六个目标:
| 编号 | 设计目标 | 核心要义 |
|---|---|---|
| G1 | 持久化双向会话 | 消除逐请求建立连接的开销 |
| G2 | 异构能力发现统一 | 通过多源画像(multi-source profiling)归一化能力发现 |
| G3 | 细粒度可靠性 | 心跳与超时管理器负责断连与失败检测 |
| G4 | 会话内确定性命令排序 | 保证命令执行顺序可预测 |
| G5 | 可组合扩展性 | 支持新增消息类型与弹性策略 |
| G6 | 透明重连与任务延续 | 在瞬时故障下实现无缝恢复 |
从源码结构看,这六个目标贯穿了整个 aip/ 包的设计:aip/messages.py 承载 G1/G2/G4 的消息契约,aip/transport/ 承载 G1/G3/G5 的传输抽象,aip/protocol/ 承载 G4/G5 的协议编排,aip/resilience/ 承载 G3/G6 的可靠性机制,aip/endpoints/ 则将全部下层能力整合为可部署的角色门面。
1.1 与 Legacy HTTP 协调的对比
| Legacy HTTP 协调 | AIP WebSocket 设计 |
|---|---|
| ❌ 短生命周期请求 | ✅ 持久化会话(G1) |
| ❌ 无状态交互 | ✅ 会话感知的任务管理 |
| ❌ 高延迟开销 | ✅ 低延迟事件流 |
| ❌ 重连支持差 | ✅ 断连无缝恢复(G6) |
| ❌ 手动状态同步 | ✅ 自动 DAG 状态传播 |
| ❌ 部分失败时脆弱 | ✅ 细粒度可靠性(G3) |
2. 五层架构总览
AIP 采用持久化、双向的 WebSocket 传输,并将编排基座拆分为五个逻辑层,每层负责可靠性与适应性的一类独立关注点。完整架构中:L1 定义语义契约,L2 提供传输灵活性,L3 实现协议逻辑,L4 保障运行期弹性,L5 交付可部署的编排原语。
各层职责与源码对应关系如下:
| 层 | 名称 | 核心职责 | 关键源码 |
|---|---|---|---|
| L1 | Message Schema Layer | 强类型 Pydantic 消息契约 | aip/messages.py |
| L2 | Transport Abstraction Layer | 协议无关的 Transport 接口与 WebSocket 实现 | aip/transport/base.py、aip/transport/websocket.py |
| L3 | Protocol Orchestration Layer | 注册、任务、心跳、命令等模块化处理器 | aip/protocol/ |
| L4 | Resilience and Health Management Layer | 心跳、超时、重连与会话恢复 | aip/resilience/ |
| L5 | Endpoint Orchestration Layer | 面向角色的门面端点 | aip/endpoints/ |
2.1 L1:消息模式层(Message Schema Layer)
本层定义强类型、经 Pydantic 校验的 ClientMessage、ServerMessage 契约,涵盖消息方向、用途与任务状态迁移。所有消息在模式层完成校验,从根本上阻止畸形消息进入协议管道,从而实现早期错误检测、简化调试。
| 职责 | 实现 | 支撑目标 |
|---|---|---|
| 消息契约 | Pydantic 模型 + 校验 | 人类可读 + 机器可验证 |
| 结构化元数据 | 系统信息、能力声明 | 统一能力发现(G2) |
| ID 关联 | 显式请求/响应链接 | 确定性排序(G4) |
2.2 L2:传输抽象层(Transport Abstraction Layer)
提供协议无关的 Transport 接口与生产级 WebSocket 实现。抽象层允许在不改动协议逻辑的前提下替换传输实现,支撑未来协议演进。
| 特性 | 收益 | 支撑目标 |
|---|---|---|
| 可配置 ping/超时 | 连接健康监控 | G3 |
| 大载荷支持 | 承载复杂任务定义 | G1 |
| 传输逻辑解耦 | 未来扩展(HTTP/3、gRPC 等) | G5 |
| 低延迟持久会话 | 消除逐请求开销 | G1 |
在源码中,Transport 抽象基类定义了 connect、send、receive、close、wait_closed 五个异步抽象方法,并通过 TransportState 枚举(DISCONNECTED、CONNECTING、CONNECTED、DISCONNECTING、ERROR)追踪连接状态(见 aip/transport/base.py)。
2.3 L3:协议编排层(Protocol Orchestration Layer)
实现注册、任务执行、心跳与命令分发的模块化处理器。每个处理器可独立测试、独立替换,在保证有序状态迁移(G4)的同时支持可组合扩展(G5)。
| 组件 | 用途 | 设计 |
|---|---|---|
AIPProtocol 基类 |
通用处理器基础设施 | 可扩展基类 |
| 处理器模块 | 注册、任务、心跳、命令 | 可插拔处理器 |
| 中间件钩子 | 日志、指标、认证 | 可组合扩展(G5) |
| 状态迁移 | 有序消息处理 | 确定性排序(G4) |
2.4 L4:弹性与健康管理层(Resilience and Health Management Layer)
容错保障:本层保证细粒度可靠性(G3)与瞬时断连下的任务无缝延续(G6),防止级联故障。
| 组件 | 机制 | 支撑目标 |
|---|---|---|
HeartbeatManager |
周期性保活信号 | G3 |
TimeoutManager |
可配置超时策略 | G3 |
ReconnectionStrategy |
指数退避 + 抖动(jitter) | G6 |
| 会话恢复 | 自动状态恢复 | G6 |
2.5 L5:端点编排层(Endpoint Orchestration Layer)
提供面向角色的门面,将下层能力整合为可部署组件,统一连接生命周期、任务路由与健康监控,并通过一致的下层能力实现强化 G1–G6。
| 端点 | 角色 | 职责 |
|---|---|---|
ConstellationEndpoint |
编排者 | 全局 Agent 注册表、任务分配、DAG 协调 |
DeviceServerEndpoint |
服务端 | WebSocket 连接管理、任务分发、结果聚合 |
DeviceClientEndpoint |
执行者 | 本地任务执行、MCP 工具调用、遥测上报 |
端点集成收益:✅ 连接生命周期管理(G1、G6)|✅ 角色化协议变体(G5)|✅ 健康监控集成(G3)|✅ 任务路由与会话管理(G4)。
3. L1 源码深读:Pydantic 强类型消息契约
消息层是 AIP 的地基,全部定义在 aip/messages.py 中,采用 Pydantic v2 的 BaseModel + Field 做声明式校验。
3.1 双向消息基类
ClientMessage(客户端 → 服务端,aip/messages.py)关键字段:
| 字段 | 类型 | 说明 |
|---|---|---|
type |
ClientMessageType |
消息类型(TASK、HEARTBEAT、COMMAND_RESULTS、REGISTER、TASK_END、DEVICE_INFO_REQUEST/RESPONSE、ERROR) |
status |
TaskStatus |
当前任务状态 |
client_type |
ClientType |
客户端类型:DEVICE 或 CONSTELLATION |
session_id |
str |
会话上下文 |
task_name |
str |
人类可读任务名 |
client_id / target_id |
str |
客户端 ID / 目标设备 ID(Constellation 客户端使用) |
request |
str |
任务请求文本 |
action_results |
List[Result] |
命令执行结果列表 |
request_id / prev_response_id |
str |
请求/响应关联 ID |
metadata |
Dict[str, Any] |
系统信息、能力声明等扩展元数据 |
ServerMessage(服务端 → 客户端,aip/messages.py)关键字段:
| 字段 | 类型 | 说明 |
|---|---|---|
type |
ServerMessageType |
TASK、COMMAND、TASK_END、HEARTBEAT、ERROR、DEVICE_INFO_REQUEST/RESPONSE |
status |
TaskStatus |
任务状态 |
user_request |
str |
原始用户请求文本 |
agent_name / process_name / root_name |
str |
执行上下文标识 |
actions |
List[Command] |
待执行命令列表(COMMAND 消息) |
messages |
List[str] |
日志消息 |
result |
Any |
TASK_END 或 DEVICE_INFO_RESPONSE 的结果载荷 |
response_id |
str |
响应关联 ID |
3.2 命令与结果结构
Command(aip/messages.py)是编排者分发的原子工作单元:
| 字段 | 类型 | 说明 |
|---|---|---|
tool_name |
str |
工具/动作名,如 "click_input" |
parameters |
Dict[str, Any] |
类型化参数,如 {"target": "Save Button", "button": "left"} |
tool_type |
Literal["data_collection", "action"] |
工具类别 |
call_id |
str |
唯一标识,如 "cmd_001" |
Result(aip/messages.py)封装命令执行结果:status(SUCCESS/FAILURE/SKIPPED/NONE)、error、result 载荷、namespace 与 call_id(与 Command.call_id 对应,实现请求-结果关联)。
TaskStatus 枚举(aip/messages.py)定义了任务生命周期状态:CONTINUE(任务进行中)、COMPLETED(成功完成)、FAILED(执行出错)、OK(确认/健康检查通过)、ERROR(协议级错误)。
3.3 协议合规校验
MessageValidator(aip/messages.py)提供静态方法做协议级约束校验,例如:
validate_registration:要求消息类型为REGISTER且必须携带client_id;validate_task_request:要求TASK消息必须携带request与client_id;validate_command_results:要求COMMAND_RESULTS消息必须携带prev_response_id且action_results非空;validate_server_message:COMMAND消息必须携带actions与response_id。
这套校验与 Pydantic 字段级校验叠加,形成"字段合法 + 协议语义合法"的双层防线。
3.4 二进制传输扩展
除文本 JSON 消息外,aip/messages.py 还定义了二进制传输协议族:BinaryMetadata(文本帧元数据:文件名、MIME 类型、大小、校验和)、FileTransferStart(分块传输开始:块大小、总块数)、FileTransferComplete(传输完成 + MD5 校验和)、ChunkMetadata(单块序号与校验)。这些模型配合 aip/protocol/base.py 中的 send_binary_message、send_file、receive_file 等异步方法,可实现"先发 JSON 元数据帧、再发二进制数据帧"的双帧传输,以及 1MB 默认分块(可通过 chunk_size 调整)的大文件断点续传。对应测试可参考 tests/aip/test_binary_transfer.py。
4. L2 源码深读:Transport 抽象与 WebSocket 实现
4.1 Transport 抽象接口
aip/transport/base.py 定义了 Transport 抽象基类,要求实现:
connect(url, **kwargs):建立连接(失败抛ConnectionError);send(data: bytes)/receive() -> bytes:字节级收发;close()/wait_closed():幂等关闭与优雅退出。
实现必须满足异步(async/await)、状态查询线程安全、对瞬时错误具备韧性。is_connected 属性基于 TransportState.CONNECTED 判定。
4.2 WebSocketTransport 实现
aip/transport/websocket.py 是默认生产级实现,构造参数如下:
| 参数 | 默认值 | 说明 |
|---|---|---|
ping_interval |
30.0 秒 |
心跳 ping 间隔 |
ping_timeout |
180.0 秒 |
ping 响应超时 |
close_timeout |
10.0 秒 |
优雅关闭超时 |
max_size |
100 * 1024 * 1024(100MB) |
最大消息尺寸 |
websocket |
None |
传入既有连接(FastAPI 服务端场景)即直接进入 CONNECTED 状态 |
传输层通过 WebSocketAdapter(aip/transport/adapters.py)屏蔽 websockets 客户端库与 FastAPI 服务端对象的差异,支持三类接收方式:receive()(文本帧,JSON 消息)、receive_binary()(二进制帧,文件/图片)、receive_auto()(自动探测帧类型)。断连时会将状态置为 DISCONNECTED 并抛 ConnectionError,供上层弹性机制捕获触发重连。
4.3 传输层配置实例
设备客户端端点实际配置(aip/endpoints/client_endpoint.py):
transport = WebSocketTransport(
ping_interval=20, ping_timeout=180, max_size=100 * 1024 * 1024
)
protocol = AIPProtocol(transport)
Constellation 端点配置(aip/endpoints/constellation_endpoint.py):
transport = WebSocketTransport(
ping_interval=30, ping_timeout=30, max_size=100 * 1024 * 1024
)
protocol = AIPProtocol(transport)
5. L3 源码深读:AIPProtocol 与模块化处理器
5.1 协议基类核心流程
AIPProtocol(aip/protocol/base.py)是全部协议逻辑的枢纽,核心方法:
send_message(aip/protocol/base.py):先正向执行出站中间件链,再将 Pydantic 模型序列化为 UTF-8 JSON,最后经transport.send发出;receive_message(aip/protocol/base.py):从传输层接收字节、按目标类型(ClientMessage/ServerMessage)反序列化,再反向执行入站中间件;register_handler/dispatch_message(aip/protocol/base.py):按消息类型注册异步处理器,dispatch_message依msg.type路由到对应处理器列表;send_error/send_ack:服务端通用的错误上报与HEARTBEAT/TaskStatus.OK确认消息;send_binary_message/receive_binary_message/send_file/receive_file:二进制与分块文件传输(默认 1MB 分块、可选 MD5 校验和校验)。
发送路径对 ConnectionError/IOError 做了精细分级:正常断连场景("closed"、"not connected")仅记录 DEBUG 日志,避免正常关机时刷 ERROR 日志干扰排障。
5.2 模块化处理器族
aip/protocol/ 目录下按关注点拆分了五个可插拔处理器:
| 模块 | 职责 | 关键方法 |
|---|---|---|
| aip/protocol/registration.py | Agent 注册与能力广播 | register_as_device、register_as_constellation、send_registration_confirmation、send_registration_error |
| aip/protocol/task_execution.py | 任务分配、命令下发、结果上报 | send_task_request、send_task_assignment、send_commands、send_command_results、send_task_end |
| aip/protocol/heartbeat.py | 心跳保活 | send_heartbeat 等 |
| aip/protocol/command.py | 命令执行协议 | 命令级收发封装 |
| aip/protocol/device_info.py | 设备遥测刷新 | 设备信息请求/响应 |
以注册为例,aip/protocol/registration.py 的 register_as_device 会构造 REGISTER 消息(携带 device_id、platform、registration_time 元数据),发送后同步等待服务端 ServerMessage 确认,status == TaskStatus.OK 即注册成功。register_as_constellation 则额外携带 targeted_device_id,声明该 Constellation 客户端将定向管理的目标设备。
5.3 任务执行协议细节
任务阶段与消息类型对应关系(实现在 aip/protocol/task_execution.py):
| 阶段 | 消息类型 | 内容 |
|---|---|---|
| 分配 | TASK |
TaskStar 定义、目标设备、命令 |
| 执行 | (内部) | MCP 工具调用、本地计算 |
| 上报 | TASK_END |
状态、日志、评估输出、结果 |
服务端下发命令使用 send_commands(aip/protocol/task_execution.py),将 List[Command] 打包进一条 ServerMessage(type=COMMAND);客户端用 send_command_results 回传 List[Result] 并携带 prev_response_id 与服务器请求关联。任务结束时,服务端通过 send_task_end 发送 TASK_END(携带 status、result/error)。
6. L4 源码深读:弹性与健康管理
6.1 HeartbeatManager
aip/resilience/heartbeat_manager.py 按客户端维度管理心跳:
default_interval默认 30 秒,可对单个客户端覆盖(start_heartbeat(client_id, interval));- 内部用
asyncio.Task维护每个客户端的_heartbeat_loop,循环先asyncio.sleep(interval),若协议仍连接则调用send_heartbeat发送保活信号,否则告警并交由连接管理器处理断连; - 提供
stop_heartbeat/stop_all(幂等取消任务)与is_running/get_interval查询接口。
6.2 TimeoutManager
aip/resilience/timeout.py 封装 asyncio.wait_for,default_timeout 默认 120 秒,支持逐操作覆盖:
with_timeout(coro, timeout, operation_name):超时抛asyncio.TimeoutError并记录详细日志;with_timeout_or_none(...):超时返回None,适合可降级的探测场景。
AIPEndpoint 基类(aip/endpoints/base.py)通过 send_with_timeout / receive_with_timeout 将收发操作纳入超时治理。
6.3 ReconnectionStrategy(指数退避 + 抖动)
aip/resilience/reconnection.py 提供四种策略枚举:EXPONENTIAL_BACKOFF、LINEAR_BACKOFF、IMMEDIATE、NONE,默认指数退避,参数默认值:
| 参数 | 默认值 | 说明 |
|---|---|---|
max_retries |
5 |
最大重连尝试次数 |
initial_backoff |
1.0 秒 |
初始退避时间 |
max_backoff |
60.0 秒 |
退避时间上限 |
backoff_multiplier |
2.0 |
指数退避倍数 |
退避计算逻辑(aip/resilience/reconnection.py)为 backoff = min(initial_backoff * (backoff_multiplier ** retry_count), max_backoff),即重连间隔按 1s → 2s → 4s → … 指数增长并封顶 60 秒。handle_disconnection 的恢复流程为:① 取消该设备的全部待处理任务(cancel_device_tasks);② 通知上层(on_device_disconnected);③ 按策略重连(attempt_reconnection);④ 成功回调 on_reconnect。重连成功后重试计数归零(reset)。
6.4 设备断连处理
设备断连时端点侧执行:标记为 DISCONNECTED、移出调度池、触发自动重连(G6);若在任务执行中断连,则任务标记为 FAILED 并传播至 ConstellationAgent 触发 DAG 编辑。ConstellationEndpoint 的 on_device_disconnected(aip/endpoints/constellation_endpoint.py)会同步调用 cancel_device_tasks 清理该设备关联任务,防止孤儿任务。
7. L5 源码深读:三类端点
7.1 ConstellationEndpoint(编排者)
aip/endpoints/constellation_endpoint.py 封装 Galaxy 的 WebSocketConnectionManager,对外提供 connect_to_device(按 AgentProfile 建连)、send_task_to_device、request_device_info、disconnect_device、is_device_connected 等编排原语,重连策略为 max_retries=5, initial_backoff=1.0, max_backoff=60.0。
7.2 DeviceServerEndpoint(服务端)
aip/endpoints/server_endpoint.py 包装 UFO 既有 UFOWebSocketHandler 实现完全向后兼容:handle_websocket 直接委托给既有处理器;cancel_device_tasks 通过 ws_manager.get_device_sessions(device_id) 找到该设备全部会话并调用 session_manager.cancel_task(session_id, reason=...) 逐个取消。服务端不主动重连,等待客户端重连。
7.3 DeviceClientEndpoint(执行者)
aip/endpoints/client_endpoint.py 包装既有 UFOWebSocketClient,start() 创建 connect_and_listen 任务并等待 connected_event;同时装配 HeartbeatManager 与 ReconnectionStrategy(max_retries=3, initial_backoff=2.0, max_backoff=60.0)。reconnect_device 通过重新 start() 实现透明重连,is_connected 委托给底层客户端。
8. 核心能力一:Agent 注册与多源画像(G2)
每个 Agent 由 AgentProfile 表示,聚合三个来源实现异构能力统一(G2):
| 来源 | 提供者 | 信息 |
|---|---|---|
| 用户配置 | ConstellationClient | 端点 URL、用户偏好、设备身份 |
| 服务清单 | Device Agent Service | 支持的工具、能力、运行元数据 |
| 客户端遥测 | Device Agent Client | 操作系统、硬件规格、GPU 状态、运行时指标 |
多级画像收益:✅ 基于实时能力精准分配任务(G2)|✅ 对环境变化(如 GPU 可用性)透明适应|✅ 设备状态变化无需手动更新|✅ 规模化场景下的明智调度决策。
动态画像更新(tip):客户端遥测持续刷新,编排者始终看到最新设备状态——这对 GPU 感知调度或跨设备负载均衡至关重要(G2)。注册协议中
metadata字段(平台、注册时间、targeted_device_id等)即是画像数据的载体,详见 aip/protocol/registration.py 与 注册流程详解。
9. 核心能力二:任务分发与结果回传(G1、G4)
AIP 使用长生命周期 WebSocket 会话横跨多次任务执行,消除逐请求连接开销并保留上下文(G1)。
任务执行时序(完整生命周期:从分配、中间执行步骤到状态更新):
sequenceDiagram
participant CC as ConstellationClient
participant DAS as Device Service
participant DAC as Device Client
CC->>DAS: TASK message (TaskStar)
DAS->>DAC: Stream task payload
DAC->>DAC: Execute using MCP tools
DAC->>DAS: Stream execution logs
DAS->>CC: TASK_END (status, logs, results)
CC->>CC: Update TaskConstellation
CC->>CC: Notify ConstellationAgent
每条箭头代表一次消息交换,垂直生命线展示事件的时间顺序。注意执行过程中日志持续回传,支持实时监控。
异步执行(warning):任务异步执行,编排者可同时向不同设备分配多个任务,结果以非确定顺序到达。相关消息格式见 messages.md,任务图语义见 TaskConstellation 文档 与 TaskStar 文档。
10. 核心能力三:命令执行(G4)
在每个任务内部,AIP 确定性执行单条命令并保持顺序,实现精确控制与错误处理(G4)。
命令结构:
| 字段 | 用途 | 示例 |
|---|---|---|
tool_name |
工具/动作名 | "click_input" |
parameters |
类型化参数 | {"target": "Save Button", "button": "left"} |
tool_type |
类别 | "action" 或 "data_collection" |
call_id |
唯一标识 | "cmd_001" |
执行保证:✅ 会话内顺序执行(确定性顺序,G4)|✅ 支持命令批量(降低网络开销)|✅ 结构化结果携带状态码与错误详情|✅ 超时传播支撑精确恢复策略(G3)。
命令批量示例:
{
"actions": [
{"tool_name": "click", "parameters": {"target": "File"}, "call_id": "1"},
{"tool_name": "click", "parameters": {"target": "Save As"}, "call_id": "2"},
{"tool_name": "type", "parameters": {"text": "document.pdf"}, "call_id": "3"}
]
}
三条命令在一条消息中发送、顺序执行。底层实现见 命令执行协议 与 aip/protocol/task_execution.py 的 send_commands,命令结果通过 Result.call_id 与请求一一对应。
11. 消息协议全景
所有 AIP 消息使用 Pydantic 模型自动校验、序列化并保证类型安全。
11.1 双向消息类型
| 方向 | 消息类型 | 用途 |
|---|---|---|
| 客户端 → 服务端 | REGISTER |
初始能力广播 |
COMMAND_RESULTS |
返回命令执行结果 | |
TASK_END |
通知任务完成 | |
HEARTBEAT |
保活信号 | |
DEVICE_INFO_RESPONSE |
设备遥测更新 | |
| 服务端 → 客户端 | TASK |
任务分配 |
COMMAND |
命令执行请求 | |
DEVICE_INFO_REQUEST |
请求遥测刷新 | |
HEARTBEAT |
保活确认 | |
| 双向 | ERROR |
错误条件上报 |
说明:
DEVICE_INFO_REQUEST在客户端与服务端消息枚举中均存在——客户端可主动请求设备信息,服务端也可下发刷新指令,体现双向遥测能力。
11.2 消息关联字段
每条消息包含:
timestamp:ISO 8601 格式时间戳;request_id/response_id:唯一标识;prev_response_id:将响应链接到请求;session_id:会话上下文。
完整字段定义与消息参考见 messages.md,枚举定义见 aip/messages.py(ClientMessageType、ServerMessageType、ClientType)。
12. 弹性连接协议(G3、G6)
网络不稳定处理(warning):AIP 通过细粒度可靠性机制与透明重连,确保在瞬时网络故障或设备断连时编排持续可用。
12.1 设备断连流程
连接状态迁移:
stateDiagram-v2
[*] --> CONNECTED
CONNECTED --> DISCONNECTED: Connection lost
DISCONNECTED --> CONNECTED: Reconnection succeeds
DISCONNECTED --> [*]: Timeout / Manual removal
note right of DISCONNECTED
• Excluded from scheduling
• Tasks marked FAILED
• Auto-reconnect triggered
end note
DISCONNECTED 状态相当于"隔离区":设备暂时移出调度池,同时自动重连持续尝试;若超时仍失败则永久移除。
| 事件 | 编排者动作 | 设备动作 |
|---|---|---|
| 设备断连 | 标记为 DISCONNECTED;移出调度;触发自动重连(G6) |
N/A |
| 重连成功 | 标记为 CONNECTED;恢复调度 |
会话恢复(G6) |
| 任务中断连 | 任务标记为 FAILED;传播至 ConstellationAgent;触发 DAG 编辑 |
N/A |
12.2 ConstellationClient 断连
双向故障处理(danger):当 ConstellationClient 断连时,所有关联的 Device Agent Services 必须:① 接收终止信号;② 中止该客户端关联的全部进行中任务;③ 防止资源泄漏与僵尸进程;④ 保持端到端一致性。
保证:✅ 无孤儿任务|✅ 客户端-服务端边界状态同步|✅ 连接恢复后快速复原(G6)|✅ TaskConstellation 状态一致(G4)。
实现级验证可参考 resilience.md、aip/resilience/ 源码,以及 tests/aip/test_resilience.py、tests/galaxy/client/test_device_disconnection_reconnection.py 等测试。
13. 扩展机制(G5)
AIP 提供多个扩展点,支持在不修改核心协议的前提下满足领域特定需求。
13.1 协议中间件
在消息管道中注入自定义处理(出站正向、入站反向执行):
from aip.protocol.base import ProtocolMiddleware
class AuditMiddleware(ProtocolMiddleware):
async def process_outgoing(self, msg):
log_to_audit_trail(msg)
return msg
async def process_incoming(self, msg):
log_to_audit_trail(msg)
return msg
仓库内置 LoggingMiddleware(aip/protocol/base.py)即基于该抽象实现,通过 add_middleware 挂载(aip/protocol/base.py)。中间件也可从 aip/extensions/middleware.py 扩展。
13.2 自定义消息处理器
为新消息类型注册处理器:
protocol.register_handler("custom_type", handle_custom_message)
register_handler 内部按消息类型维护处理器列表,dispatch_message 自动路由并隔离单个处理器的异常(aip/protocol/base.py)。
13.3 可插拔传输层
替换默认 WebSocket 传输:
from aip.transport import CustomTransport
protocol.transport = CustomTransport(config)
只要实现 aip/transport/base.py 中 Transport 的异步接口即可无缝接入,无需改动协议逻辑(G5)。更多扩展指引见 protocols.md。
14. 与 UFO 生态的集成
| 组件 | 集成点 | 收益 |
|---|---|---|
| MCP Servers | 命令执行模型对齐 MCP 消息格式 | 系统动作与 LLM 工具调用统一接口 |
| TaskConstellation | 经 AIP 消息实时状态同步 | 规划 DAG 始终反映分布式执行状态 |
| 配置系统 | Agent 端点、能力经 UFO 配置管理 | 集中管理、类型安全校验 |
| 日志与监控 | 全协议层日志覆盖 | 调试、性能监控、审计追踪 |
AIP 抽象了网络/设备异构性,使编排者将全部 Agent 视作单一事件驱动控制平面中的一等公民。相关文档:TaskConstellation(DAG 编排器)、ConstellationAgent(编排 Agent)、MCP 集成指南、配置系统。
下一步阅读:📖 消息参考|🔧 协议指南|🌐 传输层|🔌 端点|🛡️ 弹性。
15. 总结与关键要点
AIP 将分布式工作流执行转化为连贯、安全、自适应的系统,使推理与执行在多样化 Agent 与环境中无缝收敛。
| 方面 | 影响 | 目标 |
|---|---|---|
| 持久性 | 长连接降低开销、保持上下文 | G1 |
| 低延迟 | WebSocket 支持实时事件传播 | G1 |
| 能力发现 | 多源画像统一异构 Agent | G2 |
| 可靠性 | 心跳、超时、自动重连保证优雅降级 | G3、G6 |
| 确定性 | 顺序命令执行、显式 ID 关联 | G4 |
| 可扩展性 | 中间件钩子、可插拔传输、自定义处理器 | G5 |
| 开发者体验 | 强类型消息、清晰错误降低集成成本 | G5 |
通过将编排拆解为五个逻辑层——每层解决特定的可靠性与适应性关注点——AIP 使 UFO 能够在 DAG 演化(G4、G5)、Agent 更替(G3、G6)与异构执行环境(G1、G2)下保持正确性与可用性。读者可从 aip/ 源码、tests/aip/ 测试集及配套的 messages.md、protocols.md、transport.md、endpoints.md、resilience.md 五份文档继续深入。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust4.26 K641- DDeepSeek-V4.1-FlashDeepSeek-V4.1-Flash 是一个多模态混合专家(MoE)模型,拥有 5520 亿骨干参数,并支持最多一百万 token 的上下文长度。该模型原生支持图像和文本输入,并以自回归方式生成文本Python860
SlideSCIPPT插件,支持素材库、AI助手、一键添加图片标题,复制粘贴位置、一键图片对齐、一键插入Markdown(加粗、超链接等行内样式、代码块、LaTeX等块级样式)、便捷导出图片!C#601
Agent-Reach给你的 AI Agent 一键装上互联网能力。13 个平台(网页/GitHub/YouTube/小红书/B站/Twitter/Reddit 等)多后端路由,当下最稳的接入方式替你选好、装好、体检好。GitHub 主仓库同步镜像。Python1284
new-apiAI模型聚合管理中转分发系统,一个应用管理您的所有AI模型,支持将多种大模型转为统一格式调用,支持OpenAI、Claude、Gemini等格式,可供个人或者企业内部管理与分发渠道使用。🍥 A Unified AI Model Management & Distribution System. Aggregate all your LLMs into one app and access them via an OpenAI-compatible API, with native support for Claude (Messages) and Gemini formats.Go23245
JeecgBoot🔥企业级低代码平台集成了AI应用平台,帮助企业快速实现低代码开发和构建AI应用!前后端分离架构 SpringBoot,SpringCloud、Mybatis,Ant Design4、 Vue3.0、TS+vite!强大的代码生成器让前后端代码一键生成,无需写任何代码! 引领AI低代码开发模式: AI生成->OnlineCoding-> 代码生成-> 手工MERGE,显著的提高效率,又不失灵活~Java37451
