一、引言

2024年以来,以 GPT-4、Claude 为代表的大语言模型持续跃升,AI Agent 已从概念验证走向生产落地。然而,将 LLM 能力真正嵌入企业业务系统,远不止调用一次 Chat Completion API。一个真正可用的企业级 Agent,需要在任务拆解、工具调用、多轮对话三个核心维度上具备健壮的工程实现。

本文将从这三个维度出发,配合完整的可运行 Python 代码,构建一个企业级 Agent 框架。阅读顺序由底向上:先搭建基础设施,再逐一实现三大核心能力,最后组装为完整系统并通过测试验证。

前置条件:Python 3.12+。全文代码可通过 pip install -e ".[dev]" 安装,pytest tests/ -v 运行测试。


二、项目初始化

技术选型与目录结构

层面 选择 理由
语言 Python 3.12+ asyncio 异步原生支持
LLM 接入 OpenAI SDK(兼容协议) 统一接口,可切 vLLM/Ollama
配置 Pydantic Settings 类型安全的环境变量
Web 框架 FastAPI 异步 + 自动文档
测试 pytest + respx HTTP 层 Mock LLM
enterprise-agent/
├── pyproject.toml
├── .env.example
├── src/
│   ├── __init__.py
│   ├── config.py              # 配置管理
│   ├── main.py                # FastAPI 入口 + 命令行演示
│   ├── llm/
│   │   ├── __init__.py
│   │   └── client.py          # LLM 统一客户端
│   ├── tools/
│   │   ├── __init__.py
│   │   ├── base.py            # ToolDefinition + ToolRegistry
│   │   └── builtins/
│   │       ├── __init__.py
│   │       ├── search.py      # 搜索工具
│   │       └── email.py       # 邮件工具
│   ├── agent/
│   │   ├── __init__.py
│   │   ├── planner.py         # 任务规划器(LLM → DAG)
│   │   ├── executor.py        # DAG 并行执行引擎
│   │   ├── orchestrator.py    # ReAct 工具编排器
│   │   ├── memory.py          # 分层对话记忆
│   │   └── engine.py          # 主引擎集成
│   └── observability/
│       ├── __init__.py
│       └── logger.py          # 结构化日志
└── tests/
    ├── conftest.py
    ├── unit/
    │   ├── test_planner.py
    │   └── test_orchestrator.py
    └── integration/
        └── test_full_flow.py

pyproject.toml

[project]
name = "enterprise-agent"
version = "1.0.0"
requires-python = ">=3.12"
dependencies = [
    "fastapi>=0.115.0",
    "uvicorn[standard]>=0.30.0",
    "openai>=1.50.0",
    "pydantic-settings>=2.5.0",
    "structlog>=24.0.0",
    "httpx>=0.27.0",
]

[project.optional-dependencies]
dev = [
    "pytest>=8.0",
    "pytest-asyncio>=0.24.0",
    "pytest-cov>=5.0",
    "respx>=0.21.0",
    "ruff>=0.6.0",
    "mypy>=1.11.0",
]

[tool.pytest.ini_options]
asyncio_mode = "auto"
testpaths = ["tests"]

.env.example

LLM_API_KEY=your-api-key-here
LLM_BASE_URL=https://api.openai.com/v1
LLM_MODEL=gpt-4o-mini

三、基础设施层

在实现三大核心能力之前,需要先搭建配置、LLM 客户端、工具系统和日志四个底层组件。这些组件被所有上层模块所依赖。

3.1 config.py —— 配置管理

# src/config.py
from __future__ import annotations
from functools import lru_cache
from pydantic_settings import BaseSettings, SettingsConfigDict


class Settings(BaseSettings):
    model_config = SettingsConfigDict(env_file=".env", env_file_encoding="utf-8", extra="ignore")
    llm_api_key: str = ""
    llm_base_url: str = "https://api.openai.com/v1"
    llm_model: str = "gpt-4o-mini"
    llm_temperature: float = 0.1
    llm_max_tokens: int = 4096
    llm_timeout_seconds: int = 60


@lru_cache
def get_settings() -> Settings:
    return Settings()

3.2 llm/client.py —— LLM 客户端封装

LLMClientAsyncOpenAI 的对象响应统一转为字典格式,后续所有 Agent 组件都依赖这个稳定接口。

# src/llm/client.py
from __future__ import annotations
from typing import Any, Dict, List, Optional
from openai import AsyncOpenAI
from src.config import get_settings


class LLMClient:
    def __init__(self) -> None:
        s = get_settings()
        self._client = AsyncOpenAI(api_key=s.llm_api_key or "unset", base_url=s.llm_base_url, timeout=float(s.llm_timeout_seconds))
        self.model = s.llm_model
        self.temperature = s.llm_temperature
        self.max_tokens = s.llm_max_tokens

    async def chat(
        self, messages: List[Dict[str, Any]],
        tools: Optional[List[Dict[str, Any]]] = None,
        tool_choice: str = "auto",
    ) -> Dict[str, Any]:
        kwargs: Dict[str, Any] = {"model": self.model, "messages": messages, "temperature": self.temperature, "max_tokens": self.max_tokens}
        if tools:
            kwargs["tools"] = tools
            kwargs["tool_choice"] = tool_choice
        response = await self._client.chat.completions.create(**kwargs)
        choice = response.choices[0]
        return {
            "choices": [{"index": 0, "message": {
                "role": choice.message.role, "content": choice.message.content,
                "tool_calls": [{"id": tc.id, "function": {"name": tc.function.name, "arguments": tc.function.arguments}}
                               for tc in (choice.message.tool_calls or [])] if choice.message.tool_calls else None,
            }, "finish_reason": choice.finish_reason}],
            "usage": {"prompt_tokens": response.usage.prompt_tokens if response.usage else 0,
                       "completion_tokens": response.usage.completion_tokens if response.usage else 0,
                       "total_tokens": response.usage.total_tokens if response.usage else 0},
        }

    async def chat_simple(self, prompt: str) -> str:
        """用于规划/压缩等无工具调用场景"""
        resp = await self.chat([{"role": "user", "content": prompt}])
        return resp["choices"][0]["message"]["content"] or ""

3.3 tools/base.py —— 统一工具系统

ToolDefinition 同时支持函数式工具(handler)和 execute() 统一调用入口,DAG 执行器和 ReAct 编排器均通过此接口调用工具,无需关心底层实现。

# src/tools/base.py
from __future__ import annotations
import asyncio
from dataclasses import dataclass, field
from typing import Any, Callable, Dict, List, Optional


@dataclass
class ToolDefinition:
    """工具定义:统一了 handler(函数式)和 execute() 两种调用方式"""
    name: str
    description: str
    parameters_schema: Dict[str, Any] = field(default_factory=dict)
    handler: Optional[Callable] = None
    require_confirmation: bool = False
    timeout_seconds: int = 30
    category: str = "general"

    async def execute(self, **kwargs: Any) -> Any:
        """统一执行入口 —— DAG 执行器和 ReAct 编排器都通过此方法调用"""
        if self.handler is None:
            raise RuntimeError(f"工具 '{self.name}' 未绑定 handler")
        try:
            if asyncio.iscoroutinefunction(self.handler):
                return await asyncio.wait_for(self.handler(**kwargs), timeout=self.timeout_seconds)
            else:
                return await asyncio.wait_for(asyncio.to_thread(self.handler, **kwargs), timeout=self.timeout_seconds)
        except asyncio.TimeoutError:
            raise RuntimeError(f"工具 '{self.name}' 执行超时 ({self.timeout_seconds}s)")


class ToolRegistry:
    """工具注册中心"""
    def __init__(self) -> None:
        self._tools: Dict[str, ToolDefinition] = {}

    def register(self, tool: ToolDefinition) -> None:
        if tool.name in self._tools:
            raise ValueError(f"工具 '{tool.name}' 已被注册")
        self._tools[tool.name] = tool

    def get(self, name: str) -> Optional[ToolDefinition]:
        return self._tools.get(name)

    def list_tools(self) -> List[ToolDefinition]:
        return list(self._tools.values())

    def to_openai_format(self) -> List[Dict[str, Any]]:
        return [{"type": "function", "function": {"name": t.name, "description": t.description, "parameters": t.parameters_schema}}
                for t in self._tools.values()]

3.4 tools/builtins/ —— 内置工具

# src/tools/builtins/search.py
from src.tools.base import ToolDefinition

search_tool = ToolDefinition(
    name="search_knowledge",
    description="搜索知识库获取技术文档和解决方案",
    parameters_schema={
        "type": "object",
        "properties": {"query": {"type": "string", "description": "搜索关键词"}},
        "required": ["query"],
    },
    handler=lambda query: {"results": [f"关于 '{query}' 的搜索结果: 请检查服务状态并查看最近日志。"], "count": 1},
    category="data",
)
# src/tools/builtins/email.py
from src.tools.base import ToolDefinition

send_email_tool = ToolDefinition(
    name="send_email",
    description="发送邮件给指定收件人",
    parameters_schema={
        "type": "object",
        "properties": {
            "to": {"type": "string", "description": "收件人邮箱"},
            "subject": {"type": "string", "description": "邮件主题"},
            "body": {"type": "string", "description": "邮件正文"},
        },
        "required": ["to", "subject", "body"],
    },
    handler=lambda to, subject, body: {"status": "sent", "to": to, "subject": subject},
    require_confirmation=True,
    category="communication",
)

3.5 observability/logger.py —— 结构化日志

# src/observability/logger.py
from __future__ import annotations
import structlog
import logging


def setup_logging() -> None:
    structlog.configure(
        processors=[structlog.stdlib.filter_by_level, structlog.stdlib.add_log_level,
                     structlog.processors.TimeStamper(fmt="iso"), structlog.processors.format_exc_info,
                     structlog.processors.JSONRenderer()],
        context_class=dict, logger_factory=structlog.stdlib.LoggerFactory(),
        wrapper_class=structlog.stdlib.BoundLogger, cache_logger_on_first_use=True,
    )
    logging.getLogger("httpx").setLevel(logging.WARNING)
    logging.getLogger("openai").setLevel(logging.WARNING)


def get_logger(name: str) -> structlog.stdlib.BoundLogger:
    return structlog.get_logger(name)

以上基础设施就绪后,下面进入三大核心能力的实现。三者按逻辑递进:任务拆解处理复杂多步骤场景,工具调用处理普通单步工具场景,多轮对话为两者提供上下文记忆支撑。


四、任务拆解:从自然语言到 DAG

企业场景中,用户指令往往隐含多步操作。例如:“帮我分析上周北京的销售数据,生成报告并发送给张经理”——至少涉及数据查询、分析计算、报告生成、邮件发送四个子任务。

核心思路是将模糊指令转化为有向无环图(DAG),节点是原子操作,边代表依赖关系。

4.1 核心数据结构

# 以下代码定义在 src/agent/planner.py 中
from __future__ import annotations
from dataclasses import dataclass, field
from enum import Enum
from typing import Any, Dict, List, Optional
import json


class TaskStatus(Enum):
    PENDING = "pending"
    RUNNING = "running"
    COMPLETED = "completed"
    FAILED = "failed"


@dataclass
class SubTask:
    task_id: str
    name: str
    description: str
    tool_name: str
    arguments: Dict[str, Any] = field(default_factory=dict)
    depends_on: List[str] = field(default_factory=list)
    status: TaskStatus = TaskStatus.PENDING
    result: Any = None
    error: Optional[str] = None
    retry_count: int = 0
    max_retries: int = 2


@dataclass
class TaskPlan:
    original_query: str
    sub_tasks: List[SubTask]
    reasoning: str = ""

4.2 任务规划器

让 LLM 输出结构化 JSON 计划,解析为 TaskPlan。关键设计:提示词中注入完整工具列表,让 LLM 感知可用能力后再规划。

# 接 src/agent/planner.py
import re
from src.llm.client import LLMClient
from src.tools.base import ToolRegistry

TASK_PLANNER_PROMPT = """你是任务规划专家。请将用户请求拆解为多个可执行的子任务。
每个子任务对应一个具体的工具调用,明确依赖关系,以严格 JSON 输出。

可用工具:
{tools_description}

用户请求:{user_query}

输出格式:
{{
    "reasoning": "拆解思路",
    "sub_tasks": [
        {{"task_id": "task_1", "name": "查询销售数据", "description": "...", "tool_name": "...", "arguments": {{}}, "depends_on": []}}
    ]
}}"""


class TaskPlanner:
    def __init__(self, llm: LLMClient, registry: ToolRegistry) -> None:
        self.llm = llm
        self.registry = registry

    def _build_tools_desc(self) -> str:
        lines = []
        for t in self.registry.list_tools():
            lines.append(f"- {t.name}: {t.description}\n  参数: {json.dumps(t.parameters_schema, ensure_ascii=False)}")
        return "\n".join(lines)

    async def plan(self, user_query: str) -> TaskPlan:
        prompt = TASK_PLANNER_PROMPT.format(tools_description=self._build_tools_desc(), user_query=user_query)
        response = await self.llm.chat([{"role": "user", "content": prompt}])
        content = response["choices"][0]["message"]["content"] or ""
        plan_dict = self._parse_json(content)

        sub_tasks = [SubTask(
            task_id=t["task_id"], name=t["name"], description=t["description"],
            tool_name=t["tool_name"], arguments=t.get("arguments", {}),
            depends_on=t.get("depends_on", []),
        ) for t in plan_dict["sub_tasks"]]
        return TaskPlan(original_query=user_query, sub_tasks=sub_tasks, reasoning=plan_dict.get("reasoning", ""))

    @staticmethod
    def _parse_json(text: str) -> dict:
        try:
            return json.loads(text)
        except json.JSONDecodeError:
            m = re.search(r"```(?:json)?\s*([\s\S]*?)```", text)
            if m:
                return json.loads(m.group(1))
            raise ValueError(f"无法解析 LLM 输出: {text[:200]}")

4.3 DAG 执行引擎

按拓扑顺序调度子任务:无依赖的任务并行执行,有依赖的等待前置完成。支持失败重试(指数退避)和跨任务上下文变量注入(如 {{t1.result.data.rows}})。

# src/agent/executor.py
from __future__ import annotations
import asyncio
from typing import Any, Dict
from src.agent.planner import TaskPlan, SubTask, TaskStatus
from src.tools.base import ToolRegistry


class DAGExecutor:
    """DAG 执行引擎 —— 按拓扑顺序并行执行无依赖的子任务"""

    def __init__(self, registry: ToolRegistry) -> None:
        self.registry = registry

    async def execute(self, plan: TaskPlan) -> Dict[str, Any]:
        completed: Dict[str, Any] = {}
        in_flight: Dict[str, asyncio.Task] = {}

        while len(completed) + sum(1 for t in plan.sub_tasks if t.status == TaskStatus.FAILED) < len(plan.sub_tasks):
            ready = [t for t in plan.sub_tasks
                     if t.task_id not in completed
                     and t.status not in (TaskStatus.RUNNING, TaskStatus.FAILED)
                     and all(dep in completed for dep in t.depends_on)]

            if not ready and not in_flight:
                break

            for task in ready:
                task.status = TaskStatus.RUNNING
                in_flight[task.task_id] = asyncio.create_task(self._run_one(task, completed))

            done, _ = await asyncio.wait(in_flight.values(), return_when=asyncio.FIRST_COMPLETED)
            for finished in done:
                tid, result = finished.result()
                completed[tid] = result
                del in_flight[tid]

        return completed

    async def _run_one(self, task: SubTask, context: Dict[str, Any]) -> tuple:
        tool = self.registry.get(task.tool_name)
        if not tool:
            task.status = TaskStatus.FAILED
            task.error = f"未知工具: {task.tool_name}"
            return task.task_id, None

        resolved = self._resolve_args(task.arguments, context)
        for attempt in range(task.max_retries + 1):
            try:
                result = await tool.execute(**resolved)
                task.status = TaskStatus.COMPLETED
                task.result = result
                return task.task_id, result
            except Exception as e:
                task.retry_count = attempt + 1
                if attempt >= task.max_retries:
                    task.status = TaskStatus.FAILED
                    task.error = str(e)
                    return task.task_id, None
                await asyncio.sleep(2 ** attempt)
        return task.task_id, None

    @staticmethod
    def _resolve_args(args: Dict[str, Any], ctx: Dict[str, Any]) -> Dict[str, Any]:
        resolved = {}
        for k, v in args.items():
            if isinstance(v, str) and v.startswith("{{") and v.endswith("}}"):
                ref = v.strip("{} ").strip()
                parts = ref.split(".")
                current: Any = ctx
                for part in parts:
                    current = current.get(part) if isinstance(current, dict) else getattr(current, part, None)
                resolved[k] = current
            else:
                resolved[k] = v
        return resolved

五、工具调用:ReAct 循环的工程封装

当请求不需要完整 DAG 规划时(大多数场景),Agent 进入 ReAct 循环(推理→行动→观察→推理)。LLM 返回 tool_calls → 执行工具 → 结果回传 LLM → 继续推理,直到给出最终回答。与上一章的 DAG 执行器不同,编排器有 LLM 参与每一轮决策,而非一次性规划。

# src/agent/orchestrator.py
from __future__ import annotations
import json, time
from typing import Any, AsyncGenerator, Dict, List
from src.llm.client import LLMClient
from src.tools.base import ToolRegistry
from src.observability.logger import get_logger

logger = get_logger(__name__)

SYSTEM_PROMPT = """你是一个企业智能助手。使用工具时遵循规则:
1. 每次只调用必要的工具
2. 工具返回结果后综合分析再回复
3. 简单问题直接回答,不调用工具"""


class ToolOrchestrator:
    MAX_ROUNDS = 10

    def __init__(self, llm: LLMClient, registry: ToolRegistry) -> None:
        self.llm = llm
        self.registry = registry

    async def run(self, messages: List[Dict[str, Any]]) -> AsyncGenerator[Dict[str, Any], None]:
        """ReAct 循环,以事件流返回"""
        msgs: List[Dict[str, Any]] = list(messages)
        if not msgs or msgs[0].get("role") != "system":
            msgs.insert(0, {"role": "system", "content": SYSTEM_PROMPT})

        round_count = 0
        t0 = time.monotonic()

        while round_count < self.MAX_ROUNDS:
            round_count += 1
            yield {"type": "status", "content": f"第 {round_count} 轮推理"}

            response = await self.llm.chat(msgs, self.registry.to_openai_format())
            choice = response["choices"][0]
            msg = choice["message"]

            if msg.get("tool_calls"):
                tcs = msg["tool_calls"]
                yield {"type": "tool_call", "calls": [{"name": tc["function"]["name"], "args": tc["function"]["arguments"]} for tc in tcs]}

                msgs.append({"role": "assistant", "content": msg.get("content"), "tool_calls": tcs})

                for tc in tcs:
                    name = tc["function"]["name"]
                    try:
                        args = json.loads(tc["function"]["arguments"])
                    except json.JSONDecodeError:
                        args = {}

                    tool = self.registry.get(name)
                    if tool:
                        try:
                            result = await tool.execute(**args)
                            tr = {"success": True, "data": result}
                        except Exception as e:
                            tr = {"success": False, "error": str(e)}
                    else:
                        tr = {"error": f"未知工具: {name}"}

                    yield {"type": "tool_result", "tool": name, "result": tr}
                    msgs.append({"role": "tool", "tool_call_id": tc["id"], "content": json.dumps(tr, ensure_ascii=False)})
            else:
                content = msg.get("content", "")
                msgs.append({"role": "assistant", "content": content})
                logger.info("orchestrator_done", rounds=round_count, elapsed_ms=(time.monotonic() - t0) * 1000)
                yield {"type": "final", "content": content}
                return

        yield {"type": "error", "content": f"超过最大轮次 ({self.MAX_ROUNDS})"}

六、多轮对话:分层记忆与状态流转

无论用 DAG 还是 ReAct,Agent 都需要记住对话历史。核心挑战在于有限上下文窗口与长期记忆之间的平衡。解决方案分两层:近期对话完整保留,超出窗口的历史通过 LLM 压缩为摘要。

6.1 分层记忆

# src/agent/memory.py
from __future__ import annotations
from collections import deque
from dataclasses import dataclass, field
from typing import Any, Dict, List
from datetime import datetime


@dataclass
class Turn:
    role: str
    content: str
    timestamp: float = field(default_factory=lambda: datetime.now().timestamp())


class ConversationMemory:
    def __init__(self, max_recent: int = 20, token_budget: int = 8000) -> None:
        self.recent: deque[Turn] = deque(maxlen=max_recent)
        self.summary: str = ""
        self.token_budget = token_budget

    def add(self, turn: Turn) -> None:
        if len(self.recent) >= self.recent.maxlen:
            evicted = self.recent.popleft()
            self.summary = f"{self.summary}\n[{evicted.role}]: {evicted.content[:100]}..." if self.summary else f"[{evicted.role}]: {evicted.content[:100]}..."
        self.recent.append(turn)

    def to_messages(self) -> List[Dict[str, str]]:
        msgs: List[Dict[str, str]] = []
        sys = "你是一个智能助手。"
        if self.summary:
            sys += f"\n\n[历史摘要]\n{self.summary}"
        msgs.append({"role": "system", "content": sys})
        for t in self.recent:
            msgs.append({"role": t.role, "content": t.content})
        return msgs

    def estimate_tokens(self) -> int:
        return (len(self.summary) + sum(len(t.content) for t in self.recent)) // 2

6.2 上下文压缩器

Token 预算不足时,调用 LLM 主动压缩早期对话为摘要:

# 接 src/agent/memory.py
from src.llm.client import LLMClient


class ContextCompressor:
    PROMPT = "请将以下对话历史压缩为一段简洁的摘要,保留关键信息(决策、数据、结论),丢弃冗余描述。"

    def __init__(self, llm: LLMClient) -> None:
        self.llm = llm

    async def compress(self, memory: ConversationMemory) -> None:
        if memory.estimate_tokens() < memory.token_budget:
            return
        turns_list = list(memory.recent)
        mid = len(turns_list) // 2
        old = turns_list[:mid]
        new = turns_list[mid:]
        to_compress = "\n".join(f"[{t.role}]: {t.content}" for t in old)
        compressed = (await self.llm.chat_simple(f"{self.PROMPT}\n\n{to_compress}")).strip()
        memory.summary = f"{memory.summary}\n[压缩]{compressed}" if memory.summary else compressed
        memory.recent = deque(new, maxlen=memory.recent.maxlen)

6.3 会话状态机

# 接 src/agent/memory.py
from enum import Enum
from src.agent.planner import TaskPlan


class SessionState(Enum):
    IDLE = "idle"
    PLANNING = "planning"
    EXECUTING = "executing"
    COMPLETED = "completed"


@dataclass
class Session:
    session_id: str
    state: SessionState = SessionState.IDLE
    memory: ConversationMemory = field(default_factory=ConversationMemory)
    current_plan: TaskPlan | None = None
    created_at: float = field(default_factory=lambda: datetime.now().timestamp())

七、完整集成:EnterpriseAgent

前三章分别实现了任务拆解、工具调用、多轮对话。现在将它们组装为统一的 EnterpriseAgent。其 chat() 方法的决策流程是:

用户消息 → 压缩记忆 → 判断复杂度
  ├── 复杂 → 任务规划 → DAG 执行 → 结果注入记忆 → ReAct 总结
  └── 简单 ───────────────────────────────→ ReAct 直接回答
# src/agent/engine.py
from __future__ import annotations
import json
from typing import Any, AsyncGenerator, Dict, List
from src.llm.client import LLMClient
from src.tools.base import ToolRegistry, ToolDefinition
from src.agent.planner import TaskPlanner
from src.agent.executor import DAGExecutor
from src.agent.orchestrator import ToolOrchestrator
from src.agent.memory import ConversationMemory, ContextCompressor, Session, SessionState, Turn


class EnterpriseAgent:
    """企业级 AI Agent —— 集成任务拆解、工具调用与多轮对话"""

    def __init__(self, llm: LLMClient, registry: ToolRegistry, config: Dict[str, Any] | None = None) -> None:
        self.llm = llm
        self.registry = registry
        self.planner = TaskPlanner(llm, registry)
        self.executor = DAGExecutor(registry)
        self.orchestrator = ToolOrchestrator(llm, registry)
        self.compressor = ContextCompressor(llm)
        self._sessions: Dict[str, Session] = {}

    def register_tool(self, tool: ToolDefinition) -> None:
        self.registry.register(tool)

    async def chat(self, session_id: str, user_message: str) -> AsyncGenerator[Dict[str, Any], None]:
        session = self._get_or_create(session_id)
        session.memory.add(Turn(role="user", content=user_message))
        await self.compressor.compress(session.memory)

        yield {"type": "status", "content": "分析请求..."}

        if await self._needs_planning(user_message):
            session.state = SessionState.PLANNING
            yield {"type": "status", "content": "任务拆解中..."}

            plan = await self.planner.plan(user_message)
            session.current_plan = plan
            yield {"type": "plan", "reasoning": plan.reasoning, "tasks": [
                {"id": t.task_id, "name": t.name, "tool": t.tool_name} for t in plan.sub_tasks
            ]}

            session.state = SessionState.EXECUTING
            results = await self.executor.execute(plan)
            yield {"type": "dag_results", "results": {k: str(v)[:200] for k, v in results.items()}}

            session.memory.add(Turn(role="assistant", content=f"[系统] 任务完成: {json.dumps(results, ensure_ascii=False, default=str)}"))

        session.state = SessionState.EXECUTING
        msgs = session.memory.to_messages()
        async for event in self.orchestrator.run(msgs):
            yield event
            if event["type"] == "final":
                session.memory.add(Turn(role="assistant", content=event["content"]))
        session.state = SessionState.IDLE

    async def _needs_planning(self, msg: str) -> bool:
        resp = await self.llm.chat_simple(f"判断是否需要拆解为多步骤(回复 true/false):\n{msg}\n标准:多步骤/多工具 → true,简单问答 → false")
        return "true" in resp.strip().lower()

    def _get_or_create(self, sid: str) -> Session:
        if sid not in self._sessions:
            self._sessions[sid] = Session(session_id=sid)
        return self._sessions[sid]

    def close_session(self, sid: str) -> None:
        self._sessions.pop(sid, None)

八、服务化:FastAPI 入口

EnterpriseAgent 是 Python 库,需要包装为 HTTP 服务才能对外提供。使用 FastAPI 的 SSE 流式响应,逐事件推送给前端。

# src/main.py
from __future__ import annotations
import json, asyncio
from contextlib import asynccontextmanager
from typing import AsyncGenerator
from fastapi import FastAPI, Request
from fastapi.responses import StreamingResponse
from src.llm.client import LLMClient
from src.tools.base import ToolRegistry
from src.tools.builtins.search import search_tool
from src.tools.builtins.email import send_email_tool
from src.agent.engine import EnterpriseAgent

agent_engine: EnterpriseAgent | None = None


@asynccontextmanager
async def lifespan(app: FastAPI):
    global agent_engine
    registry = ToolRegistry()
    registry.register(search_tool)
    registry.register(send_email_tool)
    agent_engine = EnterpriseAgent(LLMClient(), registry)
    yield
    agent_engine = None


app = FastAPI(title="Enterprise AI Agent", version="2.0.0", lifespan=lifespan)


@app.post("/api/v1/chat")
async def chat(request: Request):
    body = await request.json()
    msg = body.get("message", "")
    sid = body.get("session_id", "default")

    async def stream() -> AsyncGenerator[str, None]:
        assert agent_engine is not None
        async for event in agent_engine.chat(sid, msg):
            yield f"data: {json.dumps(event, ensure_ascii=False)}\n\n"
        yield "data: [DONE]\n\n"

    return StreamingResponse(stream(), media_type="text/event-stream", headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no"})


@app.get("/health")
async def health():
    return {"status": "healthy"}


# ── 命令行快速测试 ──
async def demo():
    """可直接 `python src/main.py` 运行"""
    registry = ToolRegistry()
    registry.register(search_tool)
    registry.register(send_email_tool)
    agent = EnterpriseAgent(LLMClient(), registry)
    print("=== EnterpriseAgent 演示 ===\n")
    async for event in agent.chat("demo-1", "帮我搜索一下服务器宕机怎么处理"):
        print(f"[{event['type']}] {event.get('content', event.get('result', ''))}")
    agent.close_session("demo-1")


if __name__ == "__main__":
    asyncio.run(demo())

九、测试验证

代码写完了,如何确保它在没有真实 LLM API Key 的情况下也能验证?答案:respx 在 HTTP 层拦截 OpenAI SDK 的请求,返回预设 JSON。

9.1 conftest.py —— 全局测试配置

# tests/conftest.py
import pytest


@pytest.fixture(autouse=True)
def set_test_env(monkeypatch):
    monkeypatch.setenv("LLM_API_KEY", "test-key")
    monkeypatch.setattr("src.config.Settings.model_config", {"env_file": ".env.test.nonexistent", "env_file_encoding": "utf-8", "extra": "ignore"})

9.2 测试 1:任务规划器 —— 验证 DAG 构建

Mock 一次 LLM 调用,让它返回包含两个子任务(其中一个依赖另一个)的 JSON,验证解析正确性。

# tests/unit/test_planner.py
import json, pytest, respx
from httpx import Response
from src.llm.client import LLMClient
from src.tools.base import ToolRegistry, ToolDefinition
from src.agent.planner import TaskPlanner


@pytest.mark.asyncio
async def test_planner_parse_and_build_dag():
    registry = ToolRegistry()
    registry.register(ToolDefinition(name="query_db", description="查询数据库", parameters_schema={"type": "object", "properties": {"sql": {"type": "string"}}, "required": ["sql"]}, handler=lambda sql: {"rows": []}))
    registry.register(ToolDefinition(name="send_email", description="发送邮件", parameters_schema={"type": "object", "properties": {"to": {"type": "string"}}, "required": ["to"]}, handler=lambda to: {"status": "sent"}))

    with respx.mock(base_url="https://api.openai.com/v1") as mock:
        mock.post("/chat/completions").mock(return_value=Response(200, json={
            "choices": [{"index": 0, "message": {"role": "assistant", "content": json.dumps({
                "reasoning": "先查数据再发邮件", "sub_tasks": [
                    {"task_id": "t1", "name": "查询", "description": "查DB", "tool_name": "query_db", "arguments": {"sql": "SELECT *"}, "depends_on": []},
                    {"task_id": "t2", "name": "发送", "description": "发邮件", "tool_name": "send_email", "arguments": {"to": "{{t1.result.data.rows}}"}, "depends_on": ["t1"]},
                ]}), "finish_reason": "stop"}],
            "usage": {"prompt_tokens": 10, "completion_tokens": 5, "total_tokens": 15},
        }))
        planner = TaskPlanner(LLMClient(), registry)
        plan = await planner.plan("查询然后发邮件")
        assert len(plan.sub_tasks) == 2
        assert plan.sub_tasks[1].depends_on == ["t1"]
        assert plan.reasoning == "先查数据再发邮件"

9.3 测试 2:ReAct 编排器 —— 验证工具调用循环

Mock 两次 LLM 调用:第一轮返回 tool_calls(触发工具执行),第二轮返回最终回答。验证事件流覆盖完整的 ReAct 模式。

# tests/unit/test_orchestrator.py
import json, pytest, respx
from httpx import Response
from src.llm.client import LLMClient
from src.tools.base import ToolRegistry, ToolDefinition
from src.agent.orchestrator import ToolOrchestrator


@pytest.mark.asyncio
async def test_react_with_tool_call():
    registry = ToolRegistry()
    registry.register(ToolDefinition(name="search", description="搜索", parameters_schema={
        "type": "object", "properties": {"query": {"type": "string"}}, "required": ["query"],
    }, handler=lambda query: {"results": [f"'{query}' 的结果"]}))

    with respx.mock(base_url="https://api.openai.com/v1") as mock:
        mock.post("/chat/completions").mock(side_effect=[
            Response(200, json={"choices": [{"index": 0, "message": {"role": "assistant", "content": None, "tool_calls": [
                {"id": "c1", "function": {"name": "search", "arguments": json.dumps({"query": "宕机"})}}
            ]}, "finish_reason": "tool_calls"}], "usage": {"prompt_tokens": 10, "completion_tokens": 5, "total_tokens": 15}}),
            Response(200, json={"choices": [{"index": 0, "message": {"role": "assistant", "content": "建议重启。"}, "finish_reason": "stop"}], "usage": {"prompt_tokens": 10, "completion_tokens": 5, "total_tokens": 15}}),
        ])
        orchestrator = ToolOrchestrator(LLMClient(), registry)
        events = [e async for e in orchestrator.run([{"role": "user", "content": "宕机了怎么办?"}])]
        types = [e["type"] for e in events]
        assert "tool_call" in types and "tool_result" in types and "final" in types
        assert "重启" in events[-1]["content"]

9.4 测试 3:完整流程 —— 验证端到端集成

Mock 两次 LLM 调用(一次用于复杂度判断返回 false,一次用于编排器直接回答),验证 EnterpriseAgent.chat() 全链路。

# tests/integration/test_full_flow.py
import pytest, respx
from httpx import Response
from src.llm.client import LLMClient
from src.tools.base import ToolRegistry, ToolDefinition
from src.agent.engine import EnterpriseAgent


@pytest.mark.asyncio
async def test_full_agent_flow():
    registry = ToolRegistry()
    registry.register(ToolDefinition(name="get_status", description="获取系统状态", parameters_schema={"type": "object", "properties": {}}, handler=lambda: {"cpu": "45%", "mem": "68%"}))

    with respx.mock(base_url="https://api.openai.com/v1") as mock:
        # 第1次:_needs_planning → false(简单问答,不触发 DAG)
        # 第2次:orchestrator.run → 直接回答
        mock.post("/chat/completions").mock(side_effect=[
            Response(200, json={"choices": [{"index": 0, "message": {"role": "assistant", "content": "false"}, "finish_reason": "stop"}], "usage": {"prompt_tokens": 5, "completion_tokens": 1, "total_tokens": 6}}),
            Response(200, json={"choices": [{"index": 0, "message": {"role": "assistant", "content": "系统运行正常。"}, "finish_reason": "stop"}], "usage": {"prompt_tokens": 10, "completion_tokens": 5, "total_tokens": 15}}),
        ])

        agent = EnterpriseAgent(LLMClient(), registry)
        events = [e async for e in agent.chat("s1", "系统状态如何?")]
        assert events[-1]["type"] == "final"
        assert "正常" in events[-1]["content"]
        agent.close_session("s1")

十、运行操作讲解

10.1 从零搭建

cd enterprise-agent

# 创建虚拟环境
python -m venv .venv
.venv\Scripts\Activate.ps1       # Windows PowerShell
# .venv\Scripts\activate.bat     # Windows CMD
# source .venv/bin/activate      # macOS / Linux

# 安装
pip install -e ".[dev]"

10.2 运行测试(无需真实 LLM)

pytest tests/ -v

输出:

tests/unit/test_planner.py::test_planner_parse_and_build_dag PASSED
tests/unit/test_orchestrator.py::test_react_with_tool_call PASSED
tests/integration/test_full_flow.py::test_full_agent_flow PASSED

三条测试分别验证:① LLM 规划 JSON → DAG 构建 ② ReAct 工具调用循环 ③ 完整 Agent 流程(规划判断 + 编排 + 记忆)。

10.3 命令行演示

cp .env.example .env          # 编辑填入真实 LLM_API_KEY
python src/main.py            # 运行内置 demo()

输出示例:

=== EnterpriseAgent 演示 ===
[status] 分析请求...
[status] 任务拆解中...
[plan] 先搜索知识库获取方案...
[status] 第 1 轮推理
[tool_call] [{'name': 'search_knowledge', 'args': '...'}]
[tool_result] {'success': True, 'data': {...}}
[status] 第 2 轮推理
[final] 根据搜索结果,建议您:1. 检查服务状态 2. 查看最近日志 3. 尝试重启服务。

10.4 启动 API 服务

uvicorn src.main:app --reload --port 8080

访问 http://localhost:8080/docs → 展开 /api/v1/chat → Try it out → 输入:

{"message": "帮我搜索服务器宕机怎么处理", "session_id": "test-1"}

10.5 常见踩坑

问题 解决
ModuleNotFoundError: No module named 'src' 在项目根目录执行 pip install -e ".[dev]"
Settings() 报错 .env 不存在 → cp .env.example .env,或 set LLM_API_KEY=test
工具调用死循环 检查 MAX_ROUNDS=10 上限;在 system prompt 中加"避免重复调用"
JSON 解析失败 _parse_json 已实现 regex 容错;可接入 json_repair 库增强

十一、总结

本文按照"基础设施 → 核心能力 → 系统集成 → 服务化 → 测试验证"的逻辑递进,围绕企业级 AI Agent 的三大核心能力给出了完整可运行的代码框架:

能力 核心组件 关键设计
任务拆解 TaskPlanner + DAGExecutor LLM 生成 JSON DAG → 拓扑并行执行 → 上下文变量注入 → 指数退避重试
工具调用 ToolOrchestrator ReAct 推理循环 → ToolDefinition.execute() 统一执行 → 事件流输出 → 轮次上限防死循环
多轮对话 ConversationMemory + ContextCompressor 近期窗口 + LLM 压缩摘要 → Token 预算管理 → 会话状态机

整个系统从 src/config.py 的配置开始,经 LLMClient 统一 LLM 接口、ToolRegistry 管理工具,到 EnterpriseAgent 组装全部能力,最终通过 src/main.py 的 FastAPI SSE 接口对外服务。3 条测试覆盖了三个关键路径,无需真实 API Key 即可验证。

更多推荐