一个基于纯 Python 和 asyncio 构建的轻量级多智能体协作框架。
- Agent 注册中心 — Agent 注册能力并可被其他 Agent 发现
- 结构化任务 — 带优先级、预算、截止时间和状态追踪的任务定义
- 智能调度器 — 基于能力的路由、负载均衡、优先队列、超时重试
- 消息总线 — 异步点对点和广播消息,支持持久化
- 共享状态 — 键值状态板,支持事件监听
- 任务流水线 — 顺序链式任务,自动传递输出到输入
- 评分系统 — 追踪 Agent 性能历史,用于智能调度
- 预算控制 — Token 和成本限制,自动执行
- 事件系统 — 本地处理器 + Webhook HTTP 回调
- DAG 调度器 — 依赖图调度,支持并行执行、循环检测、ASCII 可视化
- 重试与容错 — 指数退避、备用 Agent、死信队列
- Agent 沙箱 — 隔离执行,资源限制,违规检测
- 监控仪表板 — 实时 ASCII 仪表板和 JSON 报告导出
- 协议版本化 — 版本协商、向后兼容、优雅降级
- 持久化 — 状态快照、回滚、自动恢复
- CLI 工具 —
python -m agent_protocol demo|status|tasks|run
import asyncio
from agent_protocol import BaseAgent, Task, AgentRegistry, Dispatcher, MessageBus, SharedState
class MyAgent(BaseAgent):
async def handle_task(self, task: Task) -> dict:
return {"result": f"已处理: {task.input['data']}"}
async def main():
registry = AgentRegistry()
bus = MessageBus()
state = SharedState()
dispatcher = Dispatcher(registry, bus, state)
agent = MyAgent("my-agent", capabilities=["process"])
registry.register(agent)
task = Task(type="process", input={"data": "hello"})
result = await dispatcher.execute(task)
print(result.result)
asyncio.run(main())from agent_protocol import DAGScheduler
dag = DAGScheduler()
dag.add_node("collect", "collect", input_data={"source": "db"})
dag.add_node("clean", "clean", dependencies=["collect"])
dag.add_node("analyze", "analyze", dependencies=["clean"])
print(dag.visualize()) # ASCII 可视化
results = await dag.execute(dispatcher) # 自动并行执行from agent_protocol import RetryPolicy, FaultTolerantExecutor
policy = RetryPolicy(max_retries=3, base_delay=1.0, exponential_base=2.0)
executor = FaultTolerantExecutor(default_policy=policy)
executor.set_fallback("primary-agent", "backup-agent")
result = await executor.execute_with_retry(task, agent, registry=registry)python examples/demo.py # 基础演示
python examples/demo_advanced.py # 高级演示(DAG、重试、监控)pip install pytest pytest-asyncio
python -m pytest tests/ -vMIT