AI Agent编排:复杂工作流的自动化实现

构建智能、可扩展的多Agent协作系统


1. 摘要/引言 (Abstract / Introduction)

1.1 问题陈述

在人工智能快速发展的今天,单个大语言模型(LLM)虽然展现出了惊人的能力,但在面对复杂的多步骤任务时,仍存在明显的局限性。如何让多个AI智能体(Agent)像专业团队一样高效协作,处理从信息检索、数据分析到决策制定的复杂工作流,已成为当前AI应用开发领域的核心挑战。

1.2 核心方案

本文将深入探讨AI Agent编排的核心概念、技术架构和实现方法。我们将学习如何设计和实现一个灵活的编排系统,使多个Agent能够各司其职、无缝协作,共同完成复杂任务。文章将结合理论讲解和实际代码实现,展示如何从零开始构建一个多Agent工作流系统。

1.3 主要成果/价值

读完本文,您将:

  • 深刻理解AI Agent编排的核心概念和重要性
  • 掌握主流Agent编排框架的设计思想和使用方法
  • 具备设计和实现自定义Agent编排系统的能力
  • 了解如何优化多Agent协作的性能和可靠性
  • 获得可直接应用于实际项目的最佳实践经验

1.4 文章导览

本文将从基础概念开始,逐步深入到架构设计和代码实现。我们将首先探讨AI Agent编排的背景和动机,然后介绍核心概念和理论基础,接着通过实际项目演示如何实现一个完整的编排系统,最后讨论性能优化、最佳实践和未来发展趋势。


2. 目标读者与前置知识 (Target Audience & Prerequisites)

2.1 目标读者

本文主要面向以下人群:

  • 有一定Python编程基础的AI应用开发者
  • 对大语言模型和AI Agent有初步了解的工程师
  • 希望构建复杂AI应用系统的架构师
  • 对自动化工作流和多Agent协作感兴趣的技术爱好者

2.2 前置知识

阅读本文前,建议您具备以下基础知识:

  • 熟练掌握Python编程语言
  • 了解基本的大语言模型(如GPT、Claude)的使用方法
  • 对API调用和异步编程有一定了解
  • 熟悉基本的软件工程概念和设计模式

3. 文章目录 (Table of Contents)

  1. 摘要/引言
  2. 目标读者与前置知识
  3. 文章目录
  4. 问题背景与动机
  5. 核心概念与理论基础
  6. 环境准备
  7. 分步实现
  8. 关键代码解析与深度剖析
  9. 结果展示与验证
  10. 性能优化与最佳实践
  11. 常见问题与解决方案
  12. 未来展望与扩展方向
  13. 总结
  14. 参考资料
  15. 附录

4. 问题背景与动机 (Problem Background & Motivation)

4.1 从单一模型到多Agent系统的演进

在过去的几年中,大语言模型(LLMs)的发展取得了突破性进展。从GPT-3到Claude,再到各种开源模型,单个AI模型的能力已经达到了令人惊叹的水平。然而,随着应用场景的不断复杂化,单一模型的局限性也日益凸显。

让我们思考一个典型的企业级场景:一个市场分析报告的生成过程。这个过程可能包括:

  1. 从多个数据源收集最新的市场数据
  2. 清洗和预处理这些数据
  3. 进行统计分析和趋势预测
  4. 检索相关的行业新闻和研究报告
  5. 综合所有信息生成结构化的分析报告
  6. 根据反馈进行多轮修改和优化

对于单个LLM来说,要独立完成所有这些任务几乎是不可能的。每个子任务都需要不同的专业知识、工具使用能力和推理方式。这就催生了多Agent系统的需求:让多个专门化的Agent协同工作,每个Agent负责自己擅长的部分,通过有效的编排机制完成整体任务。

4.2 现有解决方案的局限性

在Agent编排概念成熟之前,开发者们尝试了多种方法来解决复杂任务的自动化问题:

4.2.1 硬编码的工作流系统

早期的解决方案通常是基于硬编码的工作流引擎。开发者需要预先定义好每一个步骤、每一个条件分支,以及每一种可能的异常处理逻辑。

优点:

  • 执行流程明确可控
  • 性能相对可预测
  • 适合高度结构化的场景

缺点:

  • 缺乏灵活性,难以应对变化
  • 开发和维护成本极高
  • 无法处理模糊或开放式的任务
  • 难以利用LLM的推理能力
4.2.2 单一Agent的"思考-行动"循环

随着LLM能力的增强,出现了ReAct(Reasoning + Acting)等框架,让单个Agent能够通过思考、行动、观察的循环来完成任务。

优点:

  • 能够处理相对复杂的任务
  • 具有一定的适应性
  • 可以利用LLM的推理能力

缺点:

  • 单个Agent的能力和注意力有限
  • 缺乏专业分工,效率低下
  • 在多步骤任务中容易出现错误累积
  • 难以并行处理任务
4.2.3 简单的多Agent组合

一些系统尝试将多个Agent简单组合在一起,但缺乏有效的编排机制。

优点:

  • 实现了初步的分工
  • 可以发挥不同Agent的特长

缺点:

  • 协作效率低下
  • 缺乏统一的协调机制
  • 难以处理复杂的依赖关系
  • 错误恢复和容错能力差

4.3 为什么需要专门的Agent编排?

正是由于现有方案的种种不足,专门的AI Agent编排技术变得至关重要。一个好的编排系统应该能够:

  1. 实现高效的专业分工:让不同的Agent专注于自己擅长的任务
  2. 管理复杂的依赖关系:处理任务之间的先后顺序和数据依赖
  3. 提供灵活的控制流:支持条件分支、循环、并行执行等复杂逻辑
  4. 确保可靠的通信机制:让Agent之间能够高效、准确地传递信息
  5. 实现智能的任务分配:根据Agent的能力和负载动态分配任务
  6. 提供完善的错误处理:能够检测、恢复和容错
  7. 支持可观察性和调试:便于监控系统状态和排查问题
  8. 具备良好的可扩展性:能够方便地添加新的Agent和功能

4.4 实际应用场景

AI Agent编排的应用场景非常广泛,让我们来看几个典型的例子:

4.4.1 智能客服系统
  • 用户意图识别Agent:分析用户问题,确定服务类型
  • 知识库检索Agent:从知识库中查找相关信息
  • 问题解决Agent:根据检索结果生成解决方案
  • 情感分析Agent:监控用户情绪,必要时转接人工
  • 质量评估Agent:评估服务质量,收集反馈
4.4.2 软件开发助手
  • 需求分析Agent:理解用户需求,生成需求文档
  • 架构设计Agent:设计系统架构和技术方案
  • 代码生成Agent:根据设计生成代码
  • 代码审查Agent:检查代码质量和安全性
  • 测试用例生成Agent:生成测试用例
  • 文档编写Agent:编写技术文档
4.4.3 科研文献分析系统
  • 文献检索Agent:从多个数据库检索相关文献
  • 文献筛选Agent:根据标准筛选相关文献
  • 内容提取Agent:提取文献中的关键信息
  • 综合分析Agent:分析研究趋势和空白
  • 报告生成Agent:生成综述报告

这些场景都有一个共同点:任务复杂、步骤多、需要多种专业能力,且各步骤之间存在复杂的依赖关系。这正是Agent编排技术大显身手的地方。


5. 核心概念与理论基础 (Core Concepts & Theoretical Foundation)

5.1 什么是AI Agent?

在深入探讨Agent编排之前,我们首先需要明确什么是AI Agent。

5.1.1 Agent的定义

AI Agent是一个能够感知环境、做出决策并采取行动的自主实体。在AI和LLM的语境下,一个Agent通常包含以下核心组件:

  1. 感知模块:接收和处理环境信息
  2. 推理/决策模块:基于感知信息进行思考和决策
  3. 行动模块:执行具体的操作(如调用工具、生成文本)
  4. 记忆模块:存储历史信息和知识
  5. 通信模块:与其他Agent或系统进行交互
5.1.2 Agent的核心属性

一个好的Agent应该具备以下属性:

属性描述重要性
自主性能够在没有人为干预的情况下运行★★★★★
反应性能够感知环境并及时做出响应★★★★★
主动性能够主动设定和追求目标★★★★☆
社会性能够与其他Agent进行交互和协作★★★★★
适应性能够根据经验改进自身行为★★★☆☆
5.1.3 Agent vs 传统程序

Agent与传统程序有本质的区别:

维度传统程序AI Agent
控制方式预定义的控制流自主决策
输入输出明确的输入输出灵活的交互
环境交互有限或无持续感知和行动
适应性固定行为可学习和适应
复杂性处理难以处理模糊性善于处理不确定性

5.2 Agent编排的核心概念

5.2.1 什么是Agent编排?

Agent编排(Agent Orchestration)是指对多个Agent的任务分配、执行顺序、通信协作和资源管理进行统一协调和控制的过程。它就像一个乐队的指挥,确保每个演奏者(Agent)在正确的时间、以正确的方式演奏自己的部分,共同创造出和谐的乐章(完成复杂任务)。

5.2.2 编排的核心职责

一个编排系统通常承担以下职责:

  1. 任务分解与分配:将复杂任务分解为子任务,并分配给合适的Agent
  2. 流程控制:管理任务的执行顺序、条件分支和循环
  3. 数据路由:在Agent之间传递数据和信息
  4. 状态管理:跟踪系统状态和任务进度
  5. 错误处理:处理异常情况和错误恢复
  6. 资源管理:分配和管理计算资源
  7. 监控与可观察性:提供系统运行状态的可见性

5.3 常见的Agent架构模式

在设计多Agent系统时,有几种常见的架构模式:

5.3.1 分层架构(Hierarchical Architecture)

顶层协调Agent

中层管理Agent1

中层管理Agent2

执行Agent1

执行Agent2

执行Agent3

执行Agent4

特点:

  • 明确的上下级关系
  • 命令从上到下传递
  • 信息从下到上汇总
  • 适合有明确层级结构的任务

优点:

  • 结构清晰,易于理解
  • 责任明确
  • 适合大规模系统

缺点:

  • 灵活性较差
  • 顶层Agent可能成为瓶颈
  • 信息传递延迟
5.3.2 对等架构(Peer-to-Peer Architecture)

Agent1

Agent2

Agent3

Agent4

特点:

  • 所有Agent地位平等
  • 直接通信,无中间层
  • 分布式决策

优点:

  • 高度灵活
  • 没有单点故障
  • 适合需要紧密协作的场景

缺点:

  • 协调复杂
  • 可能出现冲突
  • 难以保证一致性
5.3.3 混合架构(Hybrid Architecture)

Agent组2

Agent3

Agent4

Agent组1

Agent1

Agent2

协调Agent

Agent组1

Agent组2

特点:

  • 结合了分层和对等的优点
  • 组内对等协作,组间分层协调
  • 平衡了灵活性和可控性

优点:

  • 兼顾灵活性和效率
  • 可扩展性好
  • 适合复杂的大规模系统
5.3.4 管道架构(Pipeline Architecture)

输入

Agent1:预处理

Agent2:分析

Agent3:生成

输出

特点:

  • 线性的任务流程
  • 数据依次经过各个Agent
  • 每个Agent负责特定步骤

优点:

  • 简单直观
  • 易于实现和调试
  • 适合有明确顺序的任务

缺点:

  • 灵活性有限
  • 难以处理分支和循环
  • 一个步骤阻塞会影响整体
5.3.5 中介架构(Mediator Architecture)

Agent1

中介/消息总线

Agent2

Agent3

Agent4

特点:

  • 通过中心中介进行通信
  • Agent之间不直接交互
  • 中介负责消息路由和协调

优点:

  • 解耦Agent
  • 易于扩展
  • 集中控制

缺点:

  • 中介可能成为瓶颈
  • 增加了系统复杂度
  • 中介故障影响全局

5.4 Agent通信机制

Agent之间的有效通信是协作的基础。常见的通信机制包括:

5.4.1 消息传递模式
Agent BAgent AAgent BAgent A处理请求请求消息响应消息
5.4.2 通信协议

Agent通信需要定义明确的协议,包括:

  1. 消息格式:JSON、XML、Protocol Buffers等
  2. 语义标准:消息类型、动作、参数等的定义
  3. 通信模式:请求-响应、发布-订阅、广播等
5.4.3 消息类型

常见的消息类型包括:

类型描述示例
请求请求其他Agent执行任务“请帮我分析这份数据”
响应对请求的回复“分析结果是…”
通知告知其他Agent事件“任务已完成”
询问获取信息“当前状态如何?”
提议提出协作建议“我们可以这样分工…”
承诺承诺执行任务“我会在明天前完成”

5.5 任务分配与调度算法

5.5.1 任务分配的考虑因素

在分配任务时,需要考虑多个因素:

  1. Agent能力:Agent是否有能力完成任务
  2. 当前负载:Agent的工作负荷
  3. 专业领域:Agent的专长领域
  4. 历史表现:Agent过去完成类似任务的表现
  5. 位置因素:数据局部性、网络延迟等
  6. 优先级:任务的紧急程度和重要性
5.5.2 常见调度算法
  1. 简单调度算法

    • 轮询(Round Robin):依次分配给各个Agent
    • 随机分配:随机选择Agent
    • 能力匹配:只分配给有能力的Agent
  2. 基于优先级的调度

    • 任务优先级排序
    • Agent优先级考虑
  3. 基于市场的调度

    • 类似拍卖的机制
    • Agent竞标任务
  4. 机器学习驱动的调度

    • 基于历史数据学习最优分配策略
    • 预测任务完成时间和资源需求

5.6 状态管理与一致性

在多Agent系统中,状态管理是一个复杂但关键的问题。

5.6.1 状态类型
  1. 本地状态:单个Agent的内部状态
  2. 共享状态:多个Agent共享的状态
  3. 全局状态:整个系统的状态
5.6.2 一致性模型
  1. 强一致性:所有节点看到相同的状态
  2. 最终一致性:经过一段时间后达到一致
  3. 会话一致性:在一个会话内保持一致
5.6.3 状态同步机制
  1. 事件溯源:记录所有状态变更事件
  2. 操作转换:处理并发操作
  3. CRDT(无冲突复制数据类型):支持分布式编辑

5.7 错误处理与容错

多Agent系统的错误处理比单一系统更复杂,因为错误可能发生在任何一个Agent或通信环节。

5.7.1 错误类型
  1. Agent错误:Agent崩溃、超时、返回错误结果
  2. 通信错误:消息丢失、延迟、乱序
  3. 数据错误:数据损坏、格式错误、不一致
  4. 逻辑错误:任务分配不当、协调逻辑错误
5.7.2 容错策略
  1. 重试机制:暂时性错误的重试
  2. 超时控制:避免无限等待
  3. 降级处理:在功能受限的情况下继续运行
  4. 故障转移:将任务转移到备用Agent
  5. ** checkpointing**:定期保存状态,便于恢复
  6. 自我修复:自动检测和修复问题

5.8 理论模型与数学基础

5.8.1 马尔可夫决策过程(MDP)

多Agent系统的决策过程可以用马尔可夫决策过程来建模:

(S,A,P,R,γ)(S, A, P, R, \gamma)(S,A,P,R,γ)

其中:

  • SSS:状态集合
  • AAA:动作集合
  • PPP:状态转移概率函数
  • RRR:奖励函数
  • γ\gammaγ:折扣因子
5.8.2 部分可观察马尔可夫决策过程(POMDP)

在更现实的场景中,Agent可能无法完全观察环境状态,这时可以用POMDP:

(S,A,P,R,Ω,O,γ)(S, A, P, R, \Omega, O, \gamma)(S,A,P,R,Ω,O,γ)

新增:

  • Ω\OmegaΩ:观察集合
  • OOO:观察概率函数
5.8.3 博弈论模型

多Agent之间的交互可以用博弈论来分析,特别是:

  • 合作博弈:Agent之间协作共同获益
  • 非合作博弈:Agent各自追求自身利益

纳什均衡是博弈论中的重要概念,描述的是一种策略组合,在这种组合中,没有任何一个Agent可以通过单独改变自己的策略来获得更好的结果。


6. 环境准备 (Environment Setup)

6.1 技术栈选择

在开始实现之前,我们需要选择合适的技术栈。对于AI Agent编排系统,我们需要考虑以下几个方面:

  1. 编程语言:Python是AI和LLM开发的事实标准
  2. LLM集成:需要支持多种LLM的集成
  3. 异步支持:处理并发任务需要良好的异步支持
  4. 状态管理:需要可靠的状态管理机制
  5. 可观测性:需要日志、监控和调试工具

6.2 核心依赖库

我们将使用以下核心库:

# requirements.txt

# LLM相关
openai>=1.0.0
langchain>=0.1.0

# 异步框架
asyncio>=3.4.3
aiohttp>=3.9.0

# 状态管理
redis>=5.0.0
celery>=5.3.0

# 数据处理
pydantic>=2.0.0
jsonpath-ng>=1.6.0

# 工具和辅助
python-dotenv>=1.0.0
loguru>=0.7.0
tenacity>=8.2.0

6.3 环境配置步骤

6.3.1 创建虚拟环境
# 创建项目目录
mkdir agent-orchestration
cd agent-orchestration

# 创建虚拟环境
python -m venv venv

# 激活虚拟环境
# Windows
venv\Scripts\activate
# Linux/Mac
source venv/bin/activate
6.3.2 安装依赖
# 升级pip
pip install --upgrade pip

# 安装依赖
pip install -r requirements.txt
6.3.3 配置环境变量

创建一个 .env 文件来存储配置:

# .env

# OpenAI配置
OPENAI_API_KEY=your_api_key_here
OPENAI_MODEL=gpt-4-turbo-preview

# Redis配置
REDIS_URL=redis://localhost:6379/0

# 应用配置
LOG_LEVEL=INFO
MAX_RETRIES=3
TIMEOUT=300
6.3.4 启动必要的服务
# 使用Docker启动Redis
docker run -d -p 6379:6379 --name orchestration-redis redis:latest

7. 分步实现 (Step-by-Step Implementation)

现在我们开始实现一个完整的AI Agent编排系统。我们将从基础组件开始,逐步构建整个系统。

7.1 项目结构设计

首先,让我们设计一个清晰的项目结构:

agent-orchestration/
├── src/
│   ├── __init__.py
│   ├── core/
│   │   ├── __init__.py
│   │   ├── agent.py          # Agent基类和实现
│   │   ├── orchestrator.py   # 编排器核心
│   │   ├── workflow.py       # 工作流定义
│   │   ├── message.py        # 消息定义
│   │   └── state.py          # 状态管理
│   ├── agents/
│   │   ├── __init__.py
│   │   ├── researcher.py     # 研究Agent
│   │   ├── writer.py         # 写作Agent
│   │   ├── reviewer.py       # 评审Agent
│   │   └── researcher.py     # 研究Agent
│   ├── tools/
│   │   ├── __init__.py
│   │   ├── search.py         # 搜索工具
│   │   └── document.py       # 文档处理工具
│   └── utils/
│       ├── __init__.py
│       ├── config.py         # 配置管理
│       └── logger.py         # 日志工具
├── examples/
│   └── research_report.py    # 研究报告示例
├── tests/
│   ├── __init__.py
│   ├── test_agent.py
│   └── test_orchestrator.py
├── requirements.txt
├── .env
└── README.md

7.2 基础组件实现

7.2.1 消息系统 (src/core/message.py)

消息是Agent之间通信的基础。让我们定义一个灵活的消息系统:

"""
消息系统模块 - 定义Agent之间通信的消息结构和协议
"""

from enum import Enum
from typing import Any, Dict, Optional, List
from datetime import datetime
from pydantic import BaseModel, Field
import json


class MessageType(Enum):
    """消息类型枚举"""
    REQUEST = "request"           # 请求消息
    RESPONSE = "response"         # 响应消息
    NOTIFICATION = "notification" # 通知消息
    QUERY = "query"               # 查询消息
    PROPOSAL = "proposal"         # 提议消息
    COMMITMENT = "commitment"     # 承诺消息
    ERROR = "error"               # 错误消息


class MessagePriority(Enum):
    """消息优先级"""
    LOW = 0
    NORMAL = 1
    HIGH = 2
    URGENT = 3


class Message(BaseModel):
    """消息基类"""
    
    message_id: str = Field(default_factory=lambda: datetime.now().strftime("%Y%m%d%H%M%S%f"))
    message_type: MessageType
    priority: MessagePriority = MessagePriority.NORMAL
    
    sender_id: str
    receiver_id: Optional[str] = None  # None表示广播
    conversation_id: Optional[str] = None
    
    timestamp: datetime = Field(default_factory=datetime.now)
    
    content: Dict[str, Any] = Field(default_factory=dict)
    metadata: Dict[str, Any] = Field(default_factory=dict)
    
    in_reply_to: Optional[str] = None  # 回复的消息ID
    
    class Config:
        """Pydantic配置"""
        arbitrary_types_allowed = True
        json_encoders = {
            datetime: lambda v: v.isoformat(),
            MessageType: lambda v: v.value,
            MessagePriority: lambda v: v.value
        }
    
    def to_dict(self) -> Dict[str, Any]:
        """转换为字典"""
        return self.model_dump()
    
    def to_json(self) -> str:
        """转换为JSON字符串"""
        return self.model_dump_json()
    
    @classmethod
    def from_dict(cls, data: Dict[str, Any]) -> "Message":
        """从字典创建消息"""
        # 处理枚举类型
        if isinstance(data.get("message_type"), str):
            data["message_type"] = MessageType(data["message_type"])
        if isinstance(data.get("priority"), str):
            data["priority"] = MessagePriority(data["priority"])
        if isinstance(data.get("timestamp"), str):
            data["timestamp"] = datetime.fromisoformat(data["timestamp"])
        
        return cls(**data)
    
    @classmethod
    def from_json(cls, json_str: str) -> "Message":
        """从JSON字符串创建消息"""
        return cls.from_dict(json.loads(json_str))
    
    def create_reply(self, content: Dict[str, Any], 
                     message_type: MessageType = MessageType.RESPONSE,
                     **kwargs) -> "Message":
        """创建回复消息"""
        return Message(
            message_type=message_type,
            sender_id=self.receiver_id or "unknown",
            receiver_id=self.sender_id,
            conversation_id=self.conversation_id,
            in_reply_to=self.message_id,
            content=content,
            **kwargs
        )


class MessageQueue:
    """消息队列 - 用于Agent之间的异步通信"""
    
    def __init__(self):
        self.queues: Dict[str, List[Message]] = {}
        self.subscribers: Dict[str, List[str]] = {}
    
    def send(self, message: Message) -> None:
        """发送消息"""
        # 如果指定了接收者,直接发送
        if message.receiver_id:
            if message.receiver_id not in self.queues:
                self.queues[message.receiver_id] = []
            self.queues[message.receiver_id].append(message)
        else:
            # 广播消息
            if message.conversation_id and message.conversation_id in self.subscribers:
                for subscriber in self.subscribers[message.conversation_id]:
                    if subscriber != message.sender_id:
                        if subscriber not in self.queues:
                            self.queues[subscriber] = []
                        self.queues[subscriber].append(message)
    
    def receive(self, agent_id: str, timeout: float = 0.0) -> Optional[Message]:
        """接收消息"""
        if agent_id not in self.queues:
            return None
        
        queue = self.queues[agent_id]
        if not queue:
            return None
        
        # 按优先级排序
        queue.sort(key=lambda m: (-m.priority.value, m.timestamp))
        
        return queue.pop(0)
    
    def subscribe(self, agent_id: str, conversation_id: str) -> None:
        """订阅会话"""
        if conversation_id not in self.subscribers:
            self.subscribers[conversation_id] = []
        if agent_id not in self.subscribers[conversation_id]:
            self.subscribers[conversation_id].append(agent_id)
    
    def unsubscribe(self, agent_id: str, conversation_id: str) -> None:
        """取消订阅"""
        if conversation_id in self.subscribers:
            if agent_id in self.subscribers[conversation_id]:
                self.subscribers[conversation_id].remove(agent_id)
    
    def get_queue_size(self, agent_id: str) -> int:
        """获取队列大小"""
        return len(self.queues.get(agent_id, []))
7.2.2 状态管理 (src/core/state.py)

状态管理对于多Agent系统至关重要:

"""
状态管理模块 - 管理Agent和系统的状态
"""

import json
import redis
import hashlib
from typing import Any, Dict, Optional, List
from datetime import datetime, timedelta
from abc import ABC, abstractmethod
from contextlib import contextmanager
from loguru import logger


class StateStore(ABC):
    """状态存储抽象基类"""
    
    @abstractmethod
    def get(self, key: str) -> Optional[Any]:
        """获取状态"""
        pass
    
    @abstractmethod
    def set(self, key: str, value: Any, ttl: Optional[int] = None) -> None:
        """设置状态"""
        pass
    
    @abstractmethod
    def delete(self, key: str) -> None:
        """删除状态"""
        pass
    
    @abstractmethod
    def exists(self, key: str) -> bool:
        """检查状态是否存在"""
        pass
    
    @abstractmethod
    def expire(self, key: str, ttl: int) -> None:
        """设置过期时间"""
        pass


class MemoryStateStore(StateStore):
    """内存状态存储 - 用于测试和简单场景"""
    
    def __init__(self):
        self._store: Dict[str, Dict[str, Any]] = {}
    
    def get(self, key: str) -> Optional[Any]:
        if key not in self._store:
            return None
        
        item = self._store[key]
        if item["expires_at"] and datetime.now() > item["expires_at"]:
            self.delete(key)
            return None
        
        return item["value"]
    
    def set(self, key: str, value: Any, ttl: Optional[int] = None) -> None:
        expires_at = None
        if ttl:
            expires_at = datetime.now() + timedelta(seconds=ttl)
        
        self._store[key] = {
            "value": value,
            "expires_at": expires_at,
            "created_at": datetime.now(),
            "updated_at": datetime.now()
        }
    
    def delete(self, key: str) -> None:
        if key in self._store:
            del self._store[key]
    
    def exists(self, key: str) -> bool:
        return key in self._store and self.get(key) is not None
    
    def expire(self, key: str, ttl: int) -> None:
        if key in self._store:
            self._store[key]["expires_at"] = datetime.now() + timedelta(seconds=ttl)
            self._store[key]["updated_at"] = datetime.now()


class RedisStateStore(StateStore):
    """Redis状态存储 - 用于生产环境"""
    
    def __init__(self, redis_url: str = "redis://localhost:6379/0"):
        self.redis_client = redis.from_url(redis_url)
        logger.info(f"Connected to Redis: {redis_url}")
    
    def get(self, key: str) -> Optional[Any]:
        try:
            data = self.redis_client.get(key)
            if data is None:
                return None
            return json.loads(data)
        except Exception as e:
            logger.error(f"Error getting key {key}: {e}")
            return None
    
    def set(self, key: str, value: Any, ttl: Optional[int] = None) -> None:
        try:
            data = json.dumps(value, default=str)
            if ttl:
                self.redis_client.setex(key, ttl, data)
            else:
                self.redis_client.set(key, data)
        except Exception as e:
            logger.error(f"Error setting key {key}: {e}")
    
    def delete(self, key: str) -> None:
        try:
            self.redis_client.delete(key)
        except Exception as e:
            logger.error(f"Error deleting key {key}: {e}")
    
    def exists(self, key: str) -> bool:
        try:
            return self.redis_client.exists(key) > 0
        except Exception as e:
            logger.error(f"Error checking key {key}: {e}")
            return False
    
    def expire(self, key: str, ttl: int) -> None:
        try:
            self.redis_client.expire(key, ttl)
        except Exception as e:
            logger.error(f"Error setting expire for key {key}: {e}")


class AgentState:
    """Agent状态管理"""
    
    def __init__(self, agent_id: str, store: StateStore):
        self.agent_id = agent_id
        self.store = store
        self._state_prefix = f"agent:{agent_id}"
    
    def _get_key(self, key: str) -> str:
        return f"{self._state_prefix}:{key}"
    
    def get(self, key: str, default: Any = None) -> Any:
        """获取状态"""
        value = self.store.get(self._get_key(key))
        return value if value is not None else default
    
    def set(self, key: str, value: Any, ttl: Optional[int] = None) -> None:
        """设置状态"""
        self.store.set(self._get_key(key), value, ttl)
    
    def delete(self, key: str) -> None:
        """删除状态"""
        self.store.delete(self._get_key(key))
    
    def exists(self, key: str) -> bool:
        """检查状态是否存在"""
        return self.store.exists(self._get_key(key))
    
    @contextmanager
    def transaction(self, keys: List[str]):
        """简单的事务上下文管理器"""
        # 对于Redis,可以使用pipeline
        # 这里简化处理
        yield
    
    def get_memory(self, memory_type: str = "short_term") -> Dict[str, Any]:
        """获取记忆"""
        return self.get(f"memory:{memory_type}", {})
    
    def add_to_memory(self, memory_type: str, key: str, value: Any) -> None:
        """添加到记忆"""
        memory = self.get_memory(memory_type)
        memory[key] = {
            "value": value,
            "timestamp": datetime.now().isoformat()
        }
        self.set(f"memory:{memory_type}", memory)
    
    def get_conversation_history(self, conversation_id: str) -> List[Dict[str, Any]]:
        """获取对话历史"""
        return self.get(f"conversation:{conversation_id}", [])
    
    def add_to_conversation(self, conversation_id: str, message: Dict[str, Any]) -> None:
        """添加到对话历史"""
        history = self.get_conversation_history(conversation_id)
        history.append({
            **message,
            "timestamp": datetime.now().isoformat()
        })
        self.set(f"conversation:{conversation_id}", history)


class WorkflowState:
    """工作流状态管理"""
    
    def __init__(self, workflow_id: str, store: StateStore):
        self.workflow_id = workflow_id
        self.store = store
        self._state_prefix = f"workflow:{workflow_id}"
    
    def _get_key(self, key: str) -> str:
        return f"{self._state_prefix}:{key}"
    
    def get_workflow_info(self) -> Dict[str, Any]:
        """获取工作流信息"""
        return self.get("info", {})
    
    def set_workflow_info(self, info: Dict[str, Any]) -> None:
        """设置工作流信息"""
        self.set("info", info)
    
    def get_step_state(self, step_id: str) -> Dict[str, Any]:
        """获取步骤状态"""
        return self.get(f"step:{step_id}", {})
    
    def set_step_state(self, step_id: str, state: Dict[str, Any]) -> None:
        """设置步骤状态"""
        self.set(f"step:{step_id}", state)
    
    def get_data(self, key: str, default: Any = None) -> Any:
        """获取工作流数据"""
        data = self.get("data", {})
        return data.get(key, default)
    
    def set_data(self, key: str, value: Any) -> None:
        """设置工作流数据"""
        data = self.get("data", {})
        data[key] = value
        self.set("data", data)
    
    def get(self, key: str, default: Any = None) -> Any:
        """获取状态"""
        value = self.store.get(self._get_key(key))
        return value if value is not None else default
    
    def set(self, key: str, value: Any, ttl: Optional[int] = None) -> None:
        """设置状态"""
        self.store.set(self._get_key(key), value, ttl)
7.2.3 Agent基类 (src/core/agent.py)

现在让我们定义Agent基类:

"""
Agent基类模块 - 定义所有Agent的基础功能和接口
"""

import asyncio
import time
import uuid
from abc import ABC, abstractmethod
from typing import Any, Dict, List, Optional, Callable, Set
from datetime import datetime
from enum import Enum
from loguru import logger

from .message import Message, MessageType, MessagePriority, MessageQueue
from .state import AgentState, StateStore


class AgentStatus(Enum):
    """Agent状态枚举"""
    IDLE = "idle"           # 空闲
    BUSY = "busy"           # 忙碌
    PAUSED = "paused"       # 暂停
    ERROR = "error"         # 错误
    STOPPED = "stopped"     # 停止


class AgentCapability:
    """Agent能力描述"""
    
    def __init__(self, name: str, description: str = "", 
                 parameters: Optional[Dict[str, Any]] = None,
                 metadata: Optional[Dict[str, Any]] = None):
        self.name = name
        self.description = description
        self.parameters = parameters or {}
        self.metadata = metadata or {}
    
    def to_dict(self) -> Dict[str, Any]:
        return {
            "name": self.name,
            "description": self.description,
            "parameters": self.parameters,
            "metadata": self.metadata
        }


class BaseAgent(ABC):
    """Agent基类 - 所有具体Agent的父类"""
    
    def __init__(self, agent_id: Optional[str] = None, 
                 name: Optional[str] = None,
                 description: str = "",
                 state_store: Optional[StateStore] = None,
                 message_queue: Optional[MessageQueue] = None,
                 config: Optional[Dict[str, Any]] = None):
        """
        初始化Agent
        
        Args:
            agent_id: Agent的唯一标识符
            name: Agent的名称
            description: Agent的描述
            state_store: 状态存储
            message_queue: 消息队列
            config: 配置信息
        """
        self.agent_id = agent_id or str(uuid.uuid4())
        self.name = name or self.__class__.__name__
        self.description = description
        
        self.config = config or {}
        self.status = AgentStatus.IDLE
        self.capabilities: List[AgentCapability] = []
        
        # 状态管理
        self.state_store = state_store
        self.state = AgentState(self.agent_id, state_store) if state_store else None
        
        # 消息队列
        self.message_queue = message_queue or MessageQueue()
        
        # 任务管理
        self.current_task: Optional[Dict[str, Any]] = None
        self.task_history: List[Dict[str, Any]] = []
        
        # 事件处理
        self.event_handlers: Dict[str, List[Callable]] = {}
        
        # 运行控制
        self._running = False
        self._task: Optional[asyncio.Task] = None
    
    def register_capability(self, capability: AgentCapability) -> None:
        """注册Agent能力"""
        self.capabilities.append(capability)
        logger.debug(f"Agent {self.name} registered capability: {capability.name}")
    
    def get_capabilities(self) -> List[Dict[str, Any]]:
        """获取Agent能力列表"""
        return [cap.to_dict() for cap in self.capabilities]
    
    def has_capability(self, capability_name: str) -> bool:
        """检查Agent是否具有特定能力"""
        return any(cap.name == capability_name for cap in self.capabilities)
    
    def on(self, event: str, handler: Callable) -> None:
        """注册事件处理器"""
        if event not in self.event_handlers:
            self.event_handlers[event] = []
        self.event_handlers[event].append(handler)
    
    def emit(self, event: str, *args, **kwargs) -> None:
        """触发事件"""
        if event in self.event_handlers:
            for handler in self.event_handlers[event]:
                try:
                    handler(*args, **kwargs)
                except Exception as e:
                    logger.error(f"Error in event handler for {event}: {e}")
    
    def send_message(self, message: Message) -> None:
        """发送消息"""
        message.sender_id = self.agent_id
        self.message_queue.send(message)
        logger.debug(f"Agent {self.name} sent message to {message.receiver_id}")
    
    def receive_message(self, timeout: float = 0.0) -> Optional[Message]:
        """接收消息"""
        return self.message_queue.receive(self.agent_id, timeout)
    
    async def process_message(self, message: Message) -> Optional[Message]:
        """
        处理接收到的消息
        
        Args:
            message: 接收到的消息
            
        Returns:
            可选的回复消息
        """
        logger.debug(f"Agent {self.name} processing message: {message.message_type}")
        
        # 根据消息类型分发处理
        if message.message_type == MessageType.REQUEST:
            return await self.handle_request(message)
        elif message.message_type == MessageType.RESPONSE:
            return await self.handle_response(message)
        elif message.message_type == MessageType.NOTIFICATION:
            return await self.handle_notification(message)
        elif message.message_type == MessageType.QUERY:
            return await self.handle_query(message)
        elif message.message_type == MessageType.ERROR:
            return await self.handle_error(message)
        else:
            logger.warning(f"Unknown message type: {message.message_type}")
            return None
    
    async def handle_request(self, message: Message) -> Optional[Message]:
        """处理请求消息"""
        # 默认实现,子类应该重写
        content = {"result": "Request received", "acknowledged": True}
Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐