Python AI助手-带MCP服务器的多代理系统
使用多代理架构和模型上下文协议(MCP)服务器的人工智能助手的最小实现。
特性
🤖 代理系统
- 基于角色的代理 具有专业能力
- 存储器系统 -代理可以通过TTL记住和回忆信息
- 任务历史记录 -跟踪所有已完成的任务
- 回调 -挂接到任务完成事件
- 能力标记 -根据他们的技能找到代理人
- 工具访问 -代理可以使用MCP服务器工具
- 优先级 -为任务选择确定代理的优先级
- 国家管理 -跟踪代理状态(空闲/繁忙)
- 自动重试 -失败时自动重试
- 指标 -成功率、持续时间、任务计数
- 学习 -代理从反馈中学习并改进
🔧 MCP服务器
- 工具注册 包含元数据和描述
- 执行日志记录 -跟踪所有工具调用
- 统计 -用平均值监控成功/失败率
- 错误处理 -优雅的故障管理
- 参数模式 -定义工具接口
- 缓存 -缓存工具性能结果
- 速率限制 -防止服务器过载
- 钩子 -执行前/执行后钩子
- 工具元数据 -描述和模式
🎯 多代理编排
- 授权 -将任务分配给特定代理
- 汽车代表团 -根据能力自动选择最佳代理
- 并行执行 -在超时的情况下同时运行多个代理
- 广播 -将任务发送给所有代理
- 共享内存 -基于TTL的跨代理数据共享
- 事件日志记录 -全系统活动跟踪
- 状态监测 -实时系统健康状况
- 中间件 -拦截和修改任务
- 排行榜 -按业绩对代理人进行排名
- 负载平衡 -高效地分配任务
🔄 工作流引擎
- 顺序执行 -有序任务处理
- 并行工作流 -同时执行独立步骤
- 依赖项 -任务等待先决条件
- 上下文传递 -在步骤之间共享结果
- 多代理工作流 -协调不同的代理
- 有条件的步骤 -根据条件执行
- 错误处理程序 -每一步自定义错误处理
- 重试逻辑 -失败时自动重试
- 工作流统计信息 -跟踪执行指标
🤝 协作
- Agent协商 -通过协商选择最佳代理
- 多Agent协作 -代理在任务上协同工作
- 投票 -代理人对决定进行投票
- 合作历史 -追踪团队努力
📡 事件系统
- 事件总线 -组件之间的发布/订阅消息传递
- 事件订阅 -订阅特定活动
- 事件历史 -查询过去的事件
- 事件过滤 -按类型和时间筛选
⏰ 调度
- 延迟任务 -安排任务以备将来执行
- 周期性任务 -周期性任务执行
- 任务取消 -取消计划任务
- 自动执行 -运行挂起的任务
🛡️ 韧性
- 断路器 -防止级联故障
- 负载平衡器 -在代理之间分配负载(循环、最不繁忙、性能最佳)
- 速率限制器 -控制请求速率
- 重试逻辑 -具有回退功能的自动重试
💾 坚持
- 国家管理 -保存/加载系统状态
- 检查点 -创建命名快照
- 指标收集 -综合指标跟踪
- 出口、进口 -在系统之间传递知识
安装
pip install -r requirements.txt快速开始
from mcp_server import MCPServer
from agent import Agent
from multi_agent_system import MultiAgentSystem
# Create MCP server with rate limiting
file_server = MCPServer("file_ops", "File operations", rate_limit=100)
file_server.register_tool("read", lambda path: f"Reading {path}", cacheable=True)
# Create agent with priority
coder = Agent("coder", "Code writer", [file_server], ["coding"], priority=1)
# Create system with max workers
system = MultiAgentSystem(max_workers=5)
system.add_agent(coder)
# Execute task
result = system.delegate("coder", "Write function")
print(result)用法示例
基本代理委托
result = system.delegate("researcher", "Find AI trends")并行执行
tasks = [
{"agent": "researcher", "task": "Research topic A"},
{"agent": "coder", "task": "Write module B"}
]
# With timeout
results = system.parallel_execute(tasks, timeout=30)汽车代表团
# Automatically select best agent by capability
result = system.auto_delegate("Find research papers", "research")
print(f"Selected: {result['agent']}")代理商排行榜
leaderboard = system.get_agent_leaderboard()
for rank, entry in enumerate(leaderboard, 1):
print(f"{rank}. {entry['name']}: {entry['metrics']['success_rate']:.2%}")代理内存
# Basic memory
agent.remember("key", "value")
value = agent.recall("key")
# Memory with TTL (expires after 60 seconds)
agent.remember("temp_key", "temp_value", ttl=60)
# Memory management
agent.forget("key")
agent.clear_memory()共享内存
# Basic shared memory
system.share_data("project_name", "AI Assistant")
data = system.get_shared_data("project_name")
# Shared memory with TTL
system.share_data("session_id", "abc123", ttl=3600)广播
results = system.broadcast("Status check")按能力查找
web_agents = system.find_agent_by_capability("web")具有依赖关系的工作流
from workflow import Workflow
workflow = Workflow("pipeline", max_retries=3)
workflow.add_step("researcher", "Gather data")
workflow.add_step("analyst", "Process data", depends_on=[0])
workflow.add_step("coder", "Generate report", depends_on=[1])
# Sequential execution
results = workflow.execute(system)
# Parallel execution (by dependency levels)
results = workflow.execute(system, parallel=True)
# Conditional steps
workflow.add_step("notifier", "Send alert",
condition=lambda r: r.get(0, {}).get("status") == "completed")
# Error handlers
workflow.add_error_handler(0, lambda e, step, ctx: {"status": "recovered"})
# Workflow stats
stats = workflow.get_stats()
print(f"Avg duration: {stats['avg_duration']:.3f}s")回调
def on_complete(result):
print(f"Task done: {result['task']}")
agent.add_callback(on_complete)MCP工具使用
result = agent.use_tool("file_ops", "read", {"path": "/data.txt"})系统状态
status = system.get_system_status()
print(status)MCP服务器统计信息
stats = server.get_stats()
print(f"Total executions: {stats['total_executions']}")
print(f"Success rate: {stats['success_rate']:.2%}")
print(f"Avg duration: {stats['avg_duration']:.3f}s")
print(f"Cache size: {stats['cache_size']}")MCP服务器挂钩
def before_hook(tool_name, params):
print(f"Executing {tool_name}")
def after_hook(log_entry):
print(f"Completed in {log_entry['duration']:.3f}s")
server.add_hook("before", before_hook)
server.add_hook("after", after_hook)代理指标
metrics = agent.get_metrics()
print(f"Total tasks: {metrics['total_tasks']}")
print(f"Success rate: {metrics['success_rate']:.2%}")
print(f"Avg duration: {metrics['avg_duration']:.3f}s")建筑
MultiAgentSystem
├── Agent (researcher)
│ ├── Memory
│ ├── Callbacks
│ └── MCP Servers
│ └── web_operations
├── Agent (coder)
│ └── MCP Servers
│ └── file_operations
└── Agent (analyst)
└── MCP Servers
├── file_operations
└── web_operations组件
代理商(agent.py)
具有角色、功能、内存和MCP服务器访问权限的个人AI助手。
MCP服务器(mcp_server.py)
使用日志和统计信息注册和执行函数的工具提供程序。
多 代理 系统multi_agent_system.py)
通过委派、并行执行和共享状态协调多个代理。
工作流程(workflow.py)
通过依赖关系解析管理顺序任务执行。
运行示例
基础示例
python example.py高级功能
python advanced_example.py高级用法
Agent学习
from learning import AgentLearning
learning = AgentLearning(agent)
learning.learn_from_feedback("Research task", "Excellent", 5.0)
best_tasks = learning.get_best_task_types()
# Export/import knowledge
knowledge = learning.export_knowledge()
learning.import_knowledge(knowledge)协作
from collaboration import AgentCollaboration
collab = AgentCollaboration(system)
# Negotiate best agent
best = collab.negotiate(["agent1", "agent2"], "Complex task")
# Collaborate on task
result = collab.collaborate(["agent1", "agent2"], "Team task")
# Voting
choice = collab.vote(["agent1", "agent2"], "Pick option", ["A", "B", "C"])事件总线
from event_bus import EventBus
bus = EventBus()
# Subscribe to events
bus.subscribe("task_complete", lambda e: print(e))
# Publish events
bus.publish("task_complete", {"task": "done"})
# Query event history
events = bus.get_events(event_type="task_complete", since=time.time()-3600)调度
from scheduler import TaskScheduler
scheduler = TaskScheduler(system)
# Delayed task (run after 10 seconds)
scheduler.schedule("agent", "Task", delay=10)
# Recurring task (run every 60 seconds)
scheduler.schedule_recurring("agent", "Task", interval=60)
# Execute pending tasks
results = scheduler.run_pending()
# Cancel recurring task
scheduler.cancel_recurring(0)负载平衡
from resilience import LoadBalancer
balancer = LoadBalancer(system)
# Set strategy: round_robin, least_busy, best_performance
balancer.set_strategy("best_performance")
agent = balancer.select_agent("capability")断路器
from resilience import CircuitBreaker
breaker = CircuitBreaker(failure_threshold=5, timeout=60)
# Protected call
result = breaker.call(system.delegate, "agent", "task")
print(f"Circuit state: {breaker.state}") # closed, open, half-open速率限制器
from resilience import RateLimiter
limiter = RateLimiter(max_requests=10, window=60)
if limiter.allow():
# Process request
pass
else:
wait = limiter.wait_time()
print(f"Rate limited, wait {wait:.1f}s")状态管理
from persistence import StateManager
manager = StateManager("state.json")
# Save current state
manager.save_state(system)
# Create checkpoint
manager.checkpoint(system, "backup1")
# Load state
state = manager.load_state()
# Restore from checkpoint
state = manager.restore_checkpoint("backup1")指标收集
from persistence import MetricsCollector
metrics = MetricsCollector()
# Record task execution
metrics.record("agent", result)
# Get summary
summary = metrics.get_summary()
print(f"Success rate: {summary['success_rate']:.2%}")
print(f"Avg duration: {summary['avg_duration']:.3f}s")
# Reset metrics
metrics.reset()API 参考
代理
Agent(name, role, mcp_servers=[], capabilities=[], priority=0)
.remember(key, value, ttl=None)
.recall(key) -> value
.forget(key)
.clear_memory()
.use_tool(server_name, tool_name, params)
.process(task, context=None)
.get_metrics() -> dict
.get_available_tools() -> dictMCP服务器
MCPServer(name, description="", rate_limit=None)
.register_tool(name, func, description="", params_schema=None, cacheable=False)
.add_hook(hook_type, func) # "before" or "after"
.execute(tool_name, params, use_cache=True)
.clear_cache()
.get_stats() -> dict
.list_tools() -> listMultiAgent 系统
MultiAgentSystem(max_workers=10)
.add_agent(agent)
.remove_agent(agent_name)
.delegate(agent_name, task, context=None, priority=0)
.auto_delegate(task, capability, context=None)
.parallel_execute(tasks, timeout=None)
.broadcast(task)
.find_agent_by_capability(capability)
.get_best_agent(capability)
.share_data(key, value, ttl=None)
.get_shared_data(key)
.get_system_status() -> dict
.get_agent_leaderboard() -> list
.add_middleware(middleware_func)工作流程
Workflow(name, max_retries=3)
.add_step(agent_name, task, depends_on=[], condition=None)
.add_error_handler(step_id, handler)
.execute(system, parallel=False)
.get_stats() -> dict许可证
麻省理工学院
