首页
/ UFO 项目 Agent Interaction Protocol(AIP)深度解析:五层架构、WebSocket 消息协议与弹性编排机制

UFO 项目 Agent Interaction Protocol(AIP)深度解析:五层架构、WebSocket 消息协议与弹性编排机制

2026-09-15 18:02:54作者:郜逊炳

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 交付可部署的编排原语。

AIP 五层架构图

各层职责与源码对应关系如下:

名称 核心职责 关键源码
L1 Message Schema Layer 强类型 Pydantic 消息契约 aip/messages.py
L2 Transport Abstraction Layer 协议无关的 Transport 接口与 WebSocket 实现 aip/transport/base.pyaip/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 校验的 ClientMessageServerMessage 契约,涵盖消息方向、用途与任务状态迁移。所有消息在模式层完成校验,从根本上阻止畸形消息进入协议管道,从而实现早期错误检测、简化调试。

职责 实现 支撑目标
消息契约 Pydantic 模型 + 校验 人类可读 + 机器可验证
结构化元数据 系统信息、能力声明 统一能力发现(G2)
ID 关联 显式请求/响应链接 确定性排序(G4)

2.2 L2:传输抽象层(Transport Abstraction Layer)

提供协议无关的 Transport 接口与生产级 WebSocket 实现。抽象层允许在不改动协议逻辑的前提下替换传输实现,支撑未来协议演进。

特性 收益 支撑目标
可配置 ping/超时 连接健康监控 G3
大载荷支持 承载复杂任务定义 G1
传输逻辑解耦 未来扩展(HTTP/3、gRPC 等) G5
低延迟持久会话 消除逐请求开销 G1

在源码中,Transport 抽象基类定义了 connectsendreceiveclosewait_closed 五个异步抽象方法,并通过 TransportState 枚举(DISCONNECTEDCONNECTINGCONNECTEDDISCONNECTINGERROR)追踪连接状态(见 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 客户端类型:DEVICECONSTELLATION
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 命令与结果结构

Commandaip/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"

Resultaip/messages.py)封装命令执行结果:statusSUCCESS/FAILURE/SKIPPED/NONE)、errorresult 载荷、namespacecall_id(与 Command.call_id 对应,实现请求-结果关联)。

TaskStatus 枚举(aip/messages.py)定义了任务生命周期状态:CONTINUE(任务进行中)、COMPLETED(成功完成)、FAILED(执行出错)、OK(确认/健康检查通过)、ERROR(协议级错误)。

3.3 协议合规校验

MessageValidatoraip/messages.py)提供静态方法做协议级约束校验,例如:

  • validate_registration:要求消息类型为 REGISTER 且必须携带 client_id
  • validate_task_request:要求 TASK 消息必须携带 requestclient_id
  • validate_command_results:要求 COMMAND_RESULTS 消息必须携带 prev_response_idaction_results 非空;
  • validate_server_messageCOMMAND 消息必须携带 actionsresponse_id

这套校验与 Pydantic 字段级校验叠加,形成"字段合法 + 协议语义合法"的双层防线。

3.4 二进制传输扩展

除文本 JSON 消息外,aip/messages.py 还定义了二进制传输协议族:BinaryMetadata(文本帧元数据:文件名、MIME 类型、大小、校验和)、FileTransferStart(分块传输开始:块大小、总块数)、FileTransferComplete(传输完成 + MD5 校验和)、ChunkMetadata(单块序号与校验)。这些模型配合 aip/protocol/base.py 中的 send_binary_messagesend_filereceive_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 状态

传输层通过 WebSocketAdapteraip/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 协议基类核心流程

AIPProtocolaip/protocol/base.py)是全部协议逻辑的枢纽,核心方法:

  • send_messageaip/protocol/base.py):先正向执行出站中间件链,再将 Pydantic 模型序列化为 UTF-8 JSON,最后经 transport.send 发出;
  • receive_messageaip/protocol/base.py):从传输层接收字节、按目标类型(ClientMessage/ServerMessage)反序列化,再反向执行入站中间件;
  • register_handler / dispatch_messageaip/protocol/base.py):按消息类型注册异步处理器,dispatch_messagemsg.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_deviceregister_as_constellationsend_registration_confirmationsend_registration_error
aip/protocol/task_execution.py 任务分配、命令下发、结果上报 send_task_requestsend_task_assignmentsend_commandssend_command_resultssend_task_end
aip/protocol/heartbeat.py 心跳保活 send_heartbeat
aip/protocol/command.py 命令执行协议 命令级收发封装
aip/protocol/device_info.py 设备遥测刷新 设备信息请求/响应

以注册为例,aip/protocol/registration.pyregister_as_device 会构造 REGISTER 消息(携带 device_idplatformregistration_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_commandsaip/protocol/task_execution.py),将 List[Command] 打包进一条 ServerMessage(type=COMMAND);客户端用 send_command_results 回传 List[Result] 并携带 prev_response_id 与服务器请求关联。任务结束时,服务端通过 send_task_end 发送 TASK_END(携带 statusresult/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_fordefault_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_BACKOFFLINEAR_BACKOFFIMMEDIATENONE,默认指数退避,参数默认值:

参数 默认值 说明
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_disconnectedaip/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_devicerequest_device_infodisconnect_deviceis_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 包装既有 UFOWebSocketClientstart() 创建 connect_and_listen 任务并等待 connected_event;同时装配 HeartbeatManagerReconnectionStrategy(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.pysend_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.pyClientMessageTypeServerMessageTypeClientType)。

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.mdaip/resilience/ 源码,以及 tests/aip/test_resilience.pytests/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

仓库内置 LoggingMiddlewareaip/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.pyTransport 的异步接口即可无缝接入,无需改动协议逻辑(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.mdprotocols.mdtransport.mdendpoints.mdresilience.md 五份文档继续深入。

登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
kernelkernel
deepin linux kernel
C
34
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.21 K
2.82 K
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
952
1.87 K
docsdocs
暂无描述
Markdown
906
5.84 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
537
614
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
864
1.36 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
4.29 K
1.04 K
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.4 K
1.48 K
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
550
402
flutter_flutterflutter_flutter
本仓库是 Flutter SDK 与 Flutter Engine 的 OpenHarmony 适配版本,由 CPF-Flutter 团队维护。开发者可使用熟悉的 Flutter 技术栈开发 OpenHarmony 应用,3.35.7 及以后的适配版本可基于本仓库源码构建支持 OpenHarmony 的 Flutter Engine。
Dart
1.19 K
348