把临时 Subagent 变成持久、可问责的 AI 团队:一套可落地的多智能体工程实践
最近在梳理多智能体系统时,发现团队里大量 Agent 还是“用完即走”的临时子任务:主 Agent 临时拉起一个 subagent,让它做一次代码评审、跑一轮测试、查一份文档,结果返回后就再也没有任何记录。这类 ad-hoc subagent 跑起来很爽,但在生产环境里却让人不踏实:任务是谁发起的?执行结果如何追溯?失败后能不能重试?权限边界有没有被绕过?
本文围绕“Turn ad-hoc subagents into durable, accountable AI teams”这一目标,完整拆解一套工程化的改造思路。我会先讲清楚 ad-hoc subagent 的局限与 durable、accountable 的含义,再给出一个基于 Python 标准库的最小可运行框架,包含任务注册表、状态机、重试机制、审计日志和简单权限边界。看完之后,你可以把这套设计迁移到自己的 Agent 编排系统里,把“随机拉起子代理”变成“有编号、有状态、有日志、有归属的 AI 团队协作”。
文章适合有 Python 基础、正在做 AI Agent 工程化或 LLM 应用落地的开发者;如果你刚接触 multi-agent,也能通过示例快速理解核心概念。
1. 背景与核心概念
1.1 什么是 ad-hoc subagent
在 LLM 应用开发中,subagent 是指由一个主 Agent(coordinator)按需创建的子执行单元。比如用户问“帮我把这个需求拆一下,顺便看看代码质量”,主 Agent 会在一次会话里临时生成一个“代码分析子代理”,等任务完成,它的上下文、中间结果、执行轨迹几乎全部消失。
这类做法我们称为 ad-hoc(临时专用)subagent。它的特点非常鲜明:
- 任务没有全局唯一编号,事后想查只能翻聊天记录。
- 状态保存在内存里,程序重启或网络抖动就丢了。
- 没有明确的权限边界,子代理和主代理共享一套 Prompt 和工具权限。
- 失败后没有统一的重试和补偿策略,只能靠用户再说一次。
- 结果没审计,谁发起的、用了哪些资源、为什么得出这个结论,全是黑盒。
在 demo 阶段,这些都不是问题;一旦进入生产环境,问题就会集中爆发。
1.2 durable:从“一次调用”到“可恢复执行”
durable 的核心是“持久化”和“可恢复”。我们希望每一次 subagent 执行都被记录为一个任务对象,这个对象拥有:
- 唯一 ID,例如
task_8f3a2c...; - 状态字段,例如
pending/running/succeeded/failed; - 输入载荷(payload),即这个任务要处理的数据;
- 输出结果或错误信息;
- 重试次数、所属负责人、父任务 ID 等元信息。
有了这些基础字段,系统就可以把任务放进注册表(Registry),执行一半崩了也能从快照或者数据库里恢复。任务不再是“调用之后就没影了”,而是变成一条可以查询、追踪、编排的记录。
1.3 accountable:可追溯、可审计、可回收
accountable 可以从三个层面理解:
- 可追溯:任务的发起人、时间、输入、输出都能对上,形成完整链路。
- 可审计:谁在什么时间段、以什么权限、调用了哪些 subagent,都有日志和审批记录。
- 可回收:如果某个子代理执行出错或越权,系统能快速定位到具体任务并终止或回滚。
换句话说,accountable 不仅是为了“出了问题能甩锅”,更是为了“出了问题能定位、能控制、能吸取教训”。对 AI 团队来说,缺少问责机制意味着错误会被无限放大:一个错误结论可能被多个 Agent 引用,最后生成一份看似合理实则错误的报告。
1.4 为什么需要一套“AI 团队”机制
单 Agent 的上下文窗口终究有限。真实业务里,一次开发交付往往需要多个角色配合:需求分析、技术方案、代码实现、测试验证、文档撰写。如果能把这些角色做成有序的 AI 团队,让它们在统一框架下协作,就能获得三个直接收益:
- 每个角色有独立职责和权限,避免“一个 Agent 一把梭”。
- 任务可以并行、排队、重试,交付过程更可控。
- 每个环节都有记录,管理者能了解进度和瓶颈。
这正是 ad-hoc subagent 转向 durable AI team 的根本动机。
2. 环境准备与版本说明
2.1 运行环境
本文的示例代码完全使用 Python 标准库实现,不依赖外部框架,方便你理解底层设计。建议环境如下:
| 项目 | 建议配置 |
|---|---|
| 操作系统 | Windows 10+ / macOS / Linux |
| Python | 3.10 或更高版本 |
| 依赖 | 仅标准库(dataclasses、enum、logging、uuid、json、threading、pathlib) |
| IDE | VS Code / PyCharm 均可 |
版本需要根据你的项目实际情况调整,本文示例以常见环境为例,重点演示配置思路。如果你的 Python 版本较低,部分类型注解语法(例如str | None)可能需要改为Optional[str]。
2.2 项目结构
ai-team-demo/ ├── ai_team/ │ ├── __init__.py │ ├── models.py │ ├── registry.py │ ├── agents.py │ └── orchestrator.py ├── main.py └── task_snapshot.json (运行后自动生成)下面我们逐文件实现。
3. 核心设计拆解
3.1 一次执行如何变成一条任务
在 ad-hoc 模式里,主 Agent 调用 subagent 是直接的函数调用,例如:
result = code_review_agent.run(code_snippet)这种写法的问题在于:调用本身没有任何中间产物,失败只能靠 try-except,成功只能靠返回值。要把它改造成 durable 执行,第一步就是把“调用”抽象成“任务”。
设计要点:
- 每个执行单元都封装为 Task 对象。
- Agent 只负责执行任务,不负责调度和存储。
- Orchestrator 负责提交、执行、重试、记录审计日志。
3.2 任务注册表(Task Registry)
任务注册表是 durable 的核心基础设施。它的职责包括:
- 保存所有任务的状态。
- 为任务生成全局唯一 ID。
- 提供状态流转方法。
- 可选持久化快照,让任务在进程重启后仍可恢复。
生产级实现通常会使用 Redis、MySQL 或专门的工作流引擎。本文用内存字典 + JSON 快照来演示核心逻辑。
3.3 状态机设计
一个任务的生命周期可以用下面这个简单状态机描述:
pending -> running -> succeeded \-> failed -> retrying -> running涉及重试时,需要注意:只有可重试的异常才允许进入 retrying,例如网络超时、临时服务不可用。业务逻辑错误(比如输入不合法)不应该无限重试。
3.4 责任归属与审计日志
accountable 体现在两个地方:
- 每个任务都带有
owner字段,标识发起人。 - Orchestrator 在任务提交、成功、失败、重试时都输出结构化日志。
在更大规模的系统里,审计日志还应包含调用链跟踪(trace_id)、资源消耗估算(token 用量、耗时)和审批记录。这里我们先实现最小版本。
4. 完整实战:实现一个可复用的 Subagent 运行框架
下面开始写代码。所有文件都放在ai-team-demo项目中。
4.1 定义任务模型
文件路径:ai-team-demo/ai_team/models.py
""" 任务模型定义:Task 是框架中最核心的数据结构。 一个 Task 表示一次 subagent 执行单元。 """ from __future__ import annotations import uuid from dataclasses import dataclass, field from datetime import datetime, timezone from enum import Enum from typing import Any, Dict, Optional class TaskStatus(str, Enum): """任务状态枚举""" PENDING = "pending" # 已提交,等待执行 RUNNING = "running" # 执行中 SUCCEEDED = "succeeded" # 执行成功 FAILED = "failed" # 执行失败(不重试) RETRYING = "retrying" # 执行失败,准备重试 CANCELLED = "cancelled" # 已取消 @dataclass class Task: """任务对象: - task_id: 全局唯一任务编号 - agent_name: 负责执行的 subagent 名称 - payload: 任务输入数据 - status: 当前状态 - owner: 发起人,默认为 system - attempts: 已执行次数 - max_attempts: 最大重试次数 - parent_task_id: 父任务编号,用于调用链追踪 - result: 成功后的结果 - error: 最后一次错误信息 """ task_id: str = field(default_factory=lambda: f"task_{uuid.uuid4().hex[:12]}") agent_name: str = "unknown" payload: Dict[str, Any] = field(default_factory=dict) status: TaskStatus = TaskStatus.PENDING owner: str = "system" attempts: int = 0 max_attempts: int = 3 parent_task_id: Optional[str] = None result: Optional[str] = None error: Optional[str] = None created_at: datetime = field( default_factory=lambda: datetime.now(timezone.utc) ) updated_at: datetime = field( default_factory=lambda: datetime.now(timezone.utc) ) def touch(self) -> None: """更新 updated_at 时间戳""" self.updated_at = datetime.now(timezone.utc) def to_dict(self) -> Dict[str, Any]: """转换为可 JSON 序列化的字典""" return { "task_id": self.task_id, "agent_name": self.agent_name, "payload": self.payload, "status": self.status.value, "owner": self.owner, "attempts": self.attempts, "max_attempts": self.max_attempts, "parent_task_id": self.parent_task_id, "result": self.result, "error": self.error, "created_at": self.created_at.isoformat(), "updated_at": self.updated_at.isoformat(), }设计解释:
task_id由 UUID 生成,避免并发冲突。parent_task_id支持父任务追踪,后续可以串联成调用链。attempts和max_attempts控制重试次数。- 所有时间使用 UTC,避免不同环境时区问题。
4.2 实现任务注册表
文件路径:ai-team-demo/ai_team/registry.py
""" 任务注册表:保存任务状态,并支持 JSON 快照持久化。 生产环境可替换为 Redis / MySQL 实现。 """ from __future__ import annotations import json import logging import threading from pathlib import Path from typing import Dict, List, Optional from .models import Task, TaskStatus logger = logging.getLogger(__name__) class TaskRegistry: """基于内存 + JSON 快照的任务注册表""" def __init__(self, snapshot_path: Optional[str] = None) -> None: self._tasks: Dict[str, Task] = {} self._lock = threading.RLock() self._snapshot_path = Path(snapshot_path) if snapshot_path else None self._load_snapshot() def add_task(self, task: Task) -> Task: """新增任务""" with self._lock: self._tasks[task.task_id] = task self._save_snapshot() logger.info("任务已注册: task_id=%s agent=%s", task.task_id, task.agent_name) return task def get_task(self, task_id: str) -> Optional[Task]: """查询任务""" with self._lock: return self._tasks.get(task_id) def update_status( self, task_id: str, status: TaskStatus, result: Optional[str] = None, error: Optional[str] = None, ) -> None: """更新任务状态""" with self._lock: task = self._tasks.get(task_id) if task is None: raise KeyError(f"任务不存在: {task_id}") task.status = status if result is not None: task.result = result if error is not None: task.error = error task.attempts += 1 if status == TaskStatus.RUNNING else 0 task.touch() self._save_snapshot() logger.info( "任务状态更新: task_id=%s status=%s attempts=%d", task_id, status.value, task.attempts, ) def list_tasks(self, owner: Optional[str] = None) -> List[Task]: """按负责人筛选任务列表""" with self._lock: tasks = list(self._tasks.values()) if owner: tasks = [t for t in tasks if t.owner == owner] return tasks def _save_snapshot(self) -> None: """保存快照,用于演示进程重启后恢复""" if self._snapshot_path is None: return try: payload = { task_id: task.to_dict() for task_id, task in self._tasks.items() } temp_path = self._snapshot_path.with_suffix(".tmp") temp_path.write_text( json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8", ) temp_path.replace(self._snapshot_path) except Exception: logger.exception("保存任务快照失败") def _load_snapshot(self) -> None: """加载快照""" if self._snapshot_path is None or not self._snapshot_path.exists(): return try: payload = json.loads(self._snapshot_path.read_text(encoding="utf-8")) for data in payload.values(): task = Task( task_id=data["task_id"], agent_name=data["agent_name"], payload=data["payload"], status=TaskStatus(data["status"]), owner=data["owner"], attempts=data["attempts"], max_attempts=data["max_attempts"], parent_task_id=data["parent_task_id"], result=data["result"], error=data["error"], ) self._tasks[task.task_id] = task logger.info("已从快照恢复 %d 个任务", len(self._tasks)) except Exception: logger.exception("加载任务快照失败,忽略")这里最大的亮点是线程锁 + 原子写临时文件再替换,避免并发写导致 JSON 损坏。生产环境可以直接把_tasks换成 MySQL 表或 Redis Hash。
4.3 定义 Subagent 基类与示例 Agent
文件路径:ai-team-demo/ai_team/agents.py
""" Subagent 定义: - SubAgent 是所有子代理的抽象基类 - 每个子代理拥有名称、描述、权限范围 - run() 方法接收一个 Task,返回执行结果字符串 """ from __future__ import annotations from abc import ABC, abstractmethod from .models import Task class SubAgent(ABC): """子代理基类""" name: str = "base" description: str = "基础子代理" permission_scope: str = "read-only" @abstractmethod def run(self, task: Task) -> str: """执行任务,返回结果""" raise NotImplementedError def __repr__(self) -> str: return f"<SubAgent name={self.name} scope={self.permission_scope}>" class CodeReviewAgent(SubAgent): """代码评审子代理:负责检查代码变更""" name = "code_review" description = "对提交的代码变更执行静态评审" permission_scope = "repo:read" def run(self, task: Task) -> str: changes = task.payload.get("changes", []) if not changes: return "评审结果:没有收到变更文件,跳过评审" issues = [] for file_path in changes: if file_path.endswith(".py"): # 这里可以接真实 LLM 或静态检查工具 issues.append(f"{file_path}: 建议补充函数 docstring") if issues: return "评审结果:发现 " + str(len(issues)) + " 个建议项\n" + "\n".join(issues) return "评审结果:未发现阻塞问题" class TestAgent(SubAgent): """测试执行子代理:负责触发测试流水线""" name = "test_runner" description = "在指定环境执行自动化测试" permission_scope = "ci:trigger" def run(self, task: Task) -> str: env = task.payload.get("env", "staging") # 模拟测试过程,真实项目中这里会调用 CI 平台 API if env == "prod": raise RuntimeError("禁止直接对生产环境执行测试,请先申请变更窗口") return f"测试结果:环境={env},全部通过(示例输出)" class DocWriterAgent(SubAgent): """文档编写子代理:负责根据上下文生成文档""" name = "doc_writer" description = "根据任务载荷生成技术文档" permission_scope = "doc:write" def run(self, task: Task) -> str: topic = task.payload.get("topic", "未命名主题") return f"文档输出:已完成关于《{topic}》的初稿,请在 docs/output 查看"可以看到,权限字段只是声明,真正执行时还需要另一个机制去校验。后面会在编排器里演示最简校验。
4.4 实现编排器
文件路径:ai-team-demo/ai_team/orchestrator.py
""" 编排器:负责任务提交、执行、重试和审计日志。 """ from __future__ import annotations import logging import time from typing import Any, Dict, List, Optional from .agents import SubAgent from .models import Task, TaskStatus from .registry import TaskRegistry logger = logging.getLogger(__name__) class TeamOrchestrator: """AI 团队编排器""" def __init__( self, registry: TaskRegistry, agents: List[SubAgent], ) -> None: self.registry = registry self.agents = {agent.name: agent for agent in agents} if len(self.agents) != len(agents): raise ValueError("存在重复的 agent_name") def submit( self, agent_name: str, payload: Dict[str, Any], owner: str = "system", parent_task_id: Optional[str] = None, max_attempts: int = 3, ) -> Task: """提交一个新任务""" if agent_name not in self.agents: raise KeyError(f"未知子代理: {agent_name}") task = Task( agent_name=agent_name, payload=payload, owner=owner, parent_task_id=parent_task_id, max_attempts=max_attempts, ) self.registry.add_task(task) logger.info( "提交任务: task_id=%s owner=%s agent=%s", task.task_id, owner, agent_name, ) return task def execute(self, task_id: str) -> Optional[str]: """执行指定任务,包含重试逻辑""" task = self.registry.get_task(task_id) if task is None: raise KeyError(f"任务不存在: {task_id}") agent = self.agents.get(task.agent_name) if agent is None: self.registry.update_status( task_id, TaskStatus.FAILED, error=f"找不到 subagent: {task.agent_name}" ) return None while task.attempts < task.max_attempts: try: self.registry.update_status(task_id, TaskStatus.RUNNING) result = agent.run(task) self.registry.update_status( task_id, TaskStatus.SUCCEEDED, result=result ) logger.info("任务成功: task_id=%s", task_id) return result except Exception as exc: logger.exception("任务执行异常: task_id=%s", task_id) task = self.registry.get_task(task_id) if task.attempts < task.max_attempts: self.registry.update_status( task_id, TaskStatus.RETRYING, error=str(exc), ) # 简单的退避策略 time.sleep(1) else: self.registry.update_status( task_id, TaskStatus.FAILED, error=str(exc), ) return None def execute_team(self, plan: List[Dict[str, Any]]) -> List[Dict[str, Any]]: """按计划顺序执行多次任务,返回结果列表""" results = [] for item in plan: agent_name = item["agent"] payload = item.get("payload", {}) owner = item.get("owner", "system") task = self.submit( agent_name=agent_name, payload=payload, owner=owner, max_attempts=item.get("max_attempts", 3), ) result = self.execute(task.task_id) results.append( { "task_id": task.task_id, "agent_name": agent_name, "owner": owner, "result": result, "status": task.status.value, } ) return results def list_tasks(self, owner: Optional[str] = None) -> List[Task]: """查询任务列表,供审计使用""" return self.registry.list_tasks(owner=owner)重试逻辑的要点:
- 每次循环重新从注册表读取 task,因为状态已经更新。
- 异常类型没有区分,实际工程建议只对可重试异常进行重试。
- 退避策略只用了固定 sleep,生产环境可以用指数退避。
4.5 编写入口主程序
文件路径:ai-team-demo/main.py
""" 演示入口:运行一个简单的 AI 团队协作计划 """ import logging from pathlib import Path from ai_team.agents import CodeReviewAgent, DocWriterAgent, TestAgent from ai_team.orchestrator import TeamOrchestrator from ai_team.registry import TaskRegistry LOGGING_FORMAT = "%(asctime)s | %(levelname)s | %(message)s" def build_orchestrator() -> TeamOrchestrator: """构建编排器,并加载本地的任务快照""" snapshot_path = Path(__file__).parent / "task_snapshot.json" registry = TaskRegistry(snapshot_path=str(snapshot_path)) orchestrator = TeamOrchestrator( registry=registry, agents=[ CodeReviewAgent(), TestAgent(), DocWriterAgent(), ], ) return orchestrator def run_demo_plan(orchestrator: TeamOrchestrator) -> None: """执行一个示例协作计划""" plan = [ { "agent": "code_review", "payload": {"changes": ["src/main.py", "tests/test_main.py"]}, "owner": "alice", }, { "agent": "test_runner", "payload": {"env": "staging"}, "owner": "alice", }, { "agent": "doc_writer", "payload": {"topic": "API 网关接入规范"}, "owner": "bob", }, ] print("\n" + "=" * 60) print("开始执行 AI 团队协作计划") print("=" * 60) results = orchestrator.execute_team(plan) print("\n" + "=" * 60) print("执行结果汇总") print("=" * 60) for item in results: print(f"- {item['agent_name']}: {item['status']}") print(f" owner={item['owner']}, task_id={item['task_id']}") print(f" result={item['result']}") print("\n" + "=" * 60) print("最近任务列表(审计视角)") print("=" * 60) for task in orchestrator.list_tasks(): print( f"{task.created_at.isoformat()} | " f"{task.task_id} | {task.agent_name} | " f"{task.status.value} | owner={task.owner}" ) def demo_retry() -> None: """演示任务失败重试机制""" orchestrator = build_orchestrator() try: task = orchestrator.submit( agent_name="test_runner", payload={"env": "prod"}, owner="alice", max_attempts=2, ) orchestrator.execute(task.task_id) except Exception as exc: print(f"提交时出现异常(按预期输出): {exc}") if __name__ == "__main__": logging.basicConfig(level=logging.INFO, format=LOGGING_FORMAT) demo_orchestrator = build_orchestrator() run_demo_plan(demo_orchestrator) demo_retry()说明:demo_retry中故意对prod环境执行测试,TestAgent会抛出异常,编排器会重试一次(max_attempts=2),最终标记为failed。这可以直观看到重试和审计日志。
4.6 运行与验证
在项目根目录执行:
python main.py预期输出(节选):
2025-xx-xx ... | INFO | 任务已注册: task_id=task_a1b2c3d4e5f6 agent=code_review 2025-xx-xx ... | INFO | 提交任务: task_id=task_a1b2c3d4e5f6 owner=alice agent=code_review 2025-xx-xx ... | INFO | 任务状态更新: task_id=task_a1b2c3d4e5f6 status=running attempts=1 2025-xx-xx ... | INFO | 任务成功: task_id=task_a1b2c3d4e5f6 执行结果汇总 - code_review: succeeded owner=alice, task_id=task_a1b2c3d4e5f6 result=评审结果:发现 2 个建议项... - test_runner: succeeded owner=alice, task_id=task_f7e8d9c0b1a2 result=测试结果:环境=staging,全部通过(示例输出) - doc_writer: succeeded owner=bob, task_id=task_112233445566 result=文档输出:已完成关于《API 网关接入规范》的初稿... 最近任务列表(审计视角) ...第二次运行时,如果task_snapshot.json存在,注册表会恢复上次任务列表,这正是“durable”的一种体现。
5. 常见问题与排查思路
5.1 任务悬挂在 running 状态
| 问题现象 | 常见原因 | 解决思路 |
|---|---|---|
任务长时间显示running | Agent 内部调用了同步阻塞 API,没有超时机制 | 给每个 Agent 执行加超时,例如asyncio.wait_for |
| 任务提交后没有执行 | 编排器没有调用execute | 检查调度流程,确认任务从队列进入执行器 |
| 进程重启后任务丢失 | 注册表没有持久化 | 接入数据库或工作流引擎 |
5.2 重试导致重复执行
| 问题现象 | 常见原因 | 解决思路 |
|---|---|---|
| Agent 执行了两次,产生重复数据 | 重试时没有幂等键 | 在 payload 中加入request_id,目标服务使用唯一约束 |
| 重试次数超过预期 | 可重试异常和业务异常没有区分 | 自定义异常类型,只对网络类异常重试 |
5.3 审计日志不完整
| 问题现象 | 常见原因 | 解决思路 |
|---|---|---|
| 查不到某个任务的发起人 | 调用submit时没传 owner | 强制要求所有任务带owner,默认值不建议使用 |
| 无法还原 agent 上下文 | 只记录结果,没记录输入和中间步骤 | Task 中应保存 payload 和最终结果,必要时保存关键中间状态 |
5.4 权限边界失效
| 问题现象 | 常见原因 | 解决思路 |
|---|---|---|
| 子代理访问了未授权资源 | 权限校验在代码里散落,没有统一入口 | 在 Orchestrator 中增加统一的权限拦截层 |
| 用户可以直接指定任意 agent_name | 没有白名单 | 在submit前校验 agent 是否在白名单中 |
5.5 排查清单
遇到问题可以按顺序检查:
- 任务是否具有唯一编号?状态是否在注册表中?
- 当前状态的迁移路径是否合法?例如
running不能直接跳到pending。 - 异常信息是否被记录到
task.error? - 审计日志的
owner、parent_task_id是否完整? - 重试逻辑是否会导致同一操作被执行多次?
6. 最佳实践与工程建议
6.1 任务 ID 与调用链设计
建议统一使用trace_id贯穿整个请求。一个用户请求可能产生多个主任务,主任务再派生多个 subagent 子任务。此时:
- 主任务持有
trace_id。 - 子任务持有
parent_task_id。 - 每个 subagent 的输出都包含
task_id和trace_id。
这样即使 Agent 数量很多,也能快速通过trace_id拉出完整执行链路。
6.2 状态持久化策略
示例里用 JSON 快照只是为了教学。生产环境建议:
- 任务表:
task_id、agent_name、status、owner、payload_json、result_json、error、created_at、updated_at。 - 用数据库事务更新状态,避免并发覆盖。
- 定期清理已成功且不再需要的任务,防止表无限膨胀。
6.3 幂等与重试
重试是把双刃剑。设计重试时必须回答:
- 失败后重试是否会重复调用第三方 API?
- 目标 API 是否支持幂等键?
- 重试的最大次数和退避策略是什么?
几个建议:
- 只对「临时性故障」重试,网络超时、HTTP 503、数据库连接断开可以重试。
- 业务校验错误(如参数不合法)直接失败,不重试。
- 使用指数退避 + 随机抖动,避免同时重试造成雪崩。
6.4 权限边界与最小权限
在 Accountable AI 团队里,权限应该默认拒绝:
- 每个 subagent 只拥有完成任务所需的最小权限。
- 主 Agent 不能随意授予新权限,统一由编排器管理。
- 涉及外部系统操作时,先走审批,再执行。
- 所有权限变更都写入审计日志。
6.5 可观测性建设
除了任务注册表,还应该暴露指标和日志:
- 指标:任务总数、成功率、平均执行时长、重试次数分布、Token 消耗估算。
- 日志:结构化 JSON 日志,字段包含 task_id、agent_name、owner、status、耗时。
- 追踪:用 OpenTelemetry 之类工具记录 span,跨子代理追踪调用链路。
6.6 从演示到生产:演进路线
本文的示例框架可以按以下路线演进:
- 将 TaskRegistry 换成 Redis/数据库实现,去掉 JSON 快照。
- 引入持久化队列,例如 Redis Stream 或 Kafka,让任务异步执行。
- 增加统一的 Agent 容器部署,让每个 subagent 运行在独立进程或容器中。
- 接入真正的 LLM 调用层,并把 token 用量和模型版本记录到 task 元信息。
- 增加人工审批节点,比如高权限任务必须等待负责人确认。
6.7 成本控制与安全边界
多 AI 协作场景下最容易失控的是成本:
- 每个子任务都会消耗 Token,建议在提交前估算任务成本。
- 设置每轮任务的总预算,超过预算自动暂停。
- 对子代理的模型选择做分级:简单任务用小模型,复杂推理才用大模型。
- 对涉及敏感数据的任务,禁止把真实数据拼进 Prompt,先做脱敏。
- 建立完善的审批和审计机制,确保每一步操作都有明确的责任人。
7. 总结与学习路线
本文的核心是从“临时调用 subagent”升级为“有状态、可追踪、可问责的 AI 团队协作”。我们实现了一个最小框架,包含四个关键部分:
- Task 模型:唯一 ID、状态、归属、重试字段。
- TaskRegistry:任务存储与快照恢复。
- SubAgent 基类:统一抽象,每个子代理有名称和权限声明。
- TeamOrchestrator:提交、执行、重试、审计日志。
如果想把这套逻辑落地到真实项目,建议按照 6.6 的路线逐步演进。优先做的三件事是:
- 把任务注册表换成数据库,保证持久化。
- 统一提交入口,强制 owner 和权限校验。
- 增加可观测性指标,让所有任务都可视化。
下一步可以继续学习:工作流引擎(如 Temporal)、Agent 通信协议、权限模型设计(RBAC / ABAC)、LLM Token 成本预算控制,以及人工审批与自动执行的混合编排模式。
建议你先把示例代码跑起来,再尝试把TaskRegistry换成 SQLite 或 MySQL 实现。动手改造之后,你对“durable、accountable”这两件事的理解会比只读文章深得多。如果本文对你有帮助,欢迎收藏备用,后续遇到 subagent 失控问题,可以直接回来翻排查清单。