实用笔记:代理架构——第6篇:多代理协调
《实用笔记:代理架构》操作指南——第6篇:多智能体协调——为采用该模式的团队提供的契约、校验机制及可直接插入的代码模块。
以下笔记为“代理架构——第6篇:多智能体协调模式”提供了一条实用的学习路径。重点在于契约、校验机制以及可直接插入的代码占位符,而非动机性阐述。 在完成概览阶段时,首先列出契约内容:所需输入、成功信号以及部分失败时的处理方式。这样的清单能确保后续的代码修改保持一致性。 在功能结果旁记录执行时间以及令牌或查询成本。提前了解成本情况,可避免在从演示环境过渡到共享环境时出现意外费用。
此处包含的内容
“你将发现什么”阶段若被视为可度量的界面,效果最佳。在扩大范围之前,先记录一份优秀的案例、一个失败案例以及回滚说明。 将配置置于应用程序代码之外。环境文件、密钥存储和功能标志应集中存放于一处,这样操作人员无需查看整个结构即可进行审计。 保持数据结构的层次简单且类型明确。嵌套的数据块会掩盖哪个节点修改了哪个字段的信息,还会在出现中断后导致流程无法继续。
为何单个代理会达到性能上限
将“单智能体”阶段视为可度量的模型时,其效果最佳。在扩大范围之前,先记录一个成功的用例、一个失败案例以及回滚说明。 同时记录正常流程与恢复流程。重试机制、人工审核环节以及死信处理都是产品本身的组成部分,而非后续需要补充的功能。 保持图结构的层次简单且类型明确。嵌套的数据块会掩盖哪个节点编写了哪个字段的信息,还会在流程中断后导致无法继续执行。
多智能体模式分类
将“多智能体分类”阶段视为可度量的对象来处理时,其效果最佳。在扩大范围之前,先记录一份理想的执行日志、一个故障案例以及回滚说明。 相较于庞大的脚本,应优先选择小型且可测试的单元。当某一步骤出现故障时,故障原因应能明确指向某个具体的责任主体,而非复杂的流程链。 保持图结构的状态简洁且具有类型定义。嵌套的数据结构会掩盖哪个节点修改了哪个字段的信息,还会在进程中断后导致无法继续执行。 将“多智能体分类”阶段视为可度量的对象来处理时,其效果最佳。在扩大范围之前,先记录一份理想的执行日志、一个故障案例以及回滚说明。 除了功能结果外,还需记录执行时间以及令牌或查询成本。提前了解这些成本信息,可避免在从演示环境过渡到共享环境时出现意外的费用支出。
+--------------------+--------------------------------------------------+------------------+
| Pattern | Structure | Best For |
+--------------------+--------------------------------------------------+------------------+
| Supervisor-Worker | One orchestrator decomposes + delegates | Complex tasks |
| | N workers execute specialized subtasks | needing expert |
| | | decomposition |
+--------------------+--------------------------------------------------+------------------+
| Pipeline | Agent A -> Agent B -> Agent C (sequential) | Transformation |
| | Each processes the previous output | chains, ETL-like |
| | | workflows |
+--------------------+--------------------------------------------------+------------------+
| Parallel Fan-out | Orchestrator sends same/related task to N agents | Research, |
| | Results aggregated into single output | analysis tasks |
| | | that parallelize |
+--------------------+--------------------------------------------------+------------------+
| Debate / Critique | Agent A produces solution, Agent B critiques | High-stakes |
| | Agent C synthesizes or adjudicates | outputs needing |
| | | adversarial QA |
+--------------------+--------------------------------------------------+------------------+
模式1:监督者-工作者
对于模式1的监管者-执行者阶段,在修改代码之前需明确输入参数、各步骤的负责人以及结束标准。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏状态。 配置信息应置于应用程序代码之外。环境文件、密钥存储以及功能标志应集中存放于一个位置,以便操作人员无需查看整个流程即可进行审计。 对于涉及资金支出或修改生产数据的节点,必须设置人工审批环节。编译时的连接方式并不等同于业务流程的完整性。
┌─────────────────────────────┐
│ SUPERVISOR AGENT │
│ (task decomposition + │
│ result synthesis) │
└──────┬───────────┬──────────┘
│ │
┌────────▼──┐ ┌────▼──────────┐ ┌──────────────┐
│ Worker A │ │ Worker B │ │ Worker C │
│ (code │ │ (security │ │ (test │
│ analysis)│ │ review) │ │ coverage) │
└────────────┘ └───────────────┘ └──────────────┘
# harness/multi_agent/supervisor.py
from typing import Literal, TypedDict, Annotated, List
from langgraph.graph import StateGraph, END, START
from langgraph.graph.message import add_messages
from langchain_core.messages import BaseMessage, HumanMessage, SystemMessage
from langchain_aws import ChatBedrock
from pydantic import BaseModel
import boto3
class SupervisorState(TypedDict):
messages: Annotated[List[BaseMessage], add_messages]
task_spec: str
subtasks: List[dict] # decomposed work items
worker_results: dict # keyed by subtask ID
current_worker: str # which worker is active
synthesis_complete: bool
agent_run_id: str
class SubtaskAssignment(BaseModel):
"""Structured output from supervisor decomposition."""
subtask_id: str
worker_type: Literal["code_analyst", "security_reviewer", "test_evaluator"]
description: str
depends_on: List[str] # subtask IDs this depends on
priority: int
SUPERVISOR_SYSTEM_PROMPT = """
You are an orchestrator agent. You do not write code or perform analysis yourself.
Your job is to:
1. Break the task into discrete subtasks
2. Assign each subtask to the correct specialist worker
3. Track dependencies between subtasks
4. Synthesize worker outputs into a coherent final result
Available workers:
- code_analyst: Reads and analyzes code structure, dependencies, patterns
- security_reviewer: Evaluates security implications, checks against CVEs
- test_evaluator: Assesses test coverage, identifies gaps
When decomposing, be specific. A subtask description like "analyze the auth module"
is useful. "analyze the code" is not.
Respond in JSON when asked to decompose. Respond in prose when asked to synthesize.
"""
def build_supervisor(region: str = "us-east-1") -> callable:
bedrock = boto3.client("bedrock-runtime", region_name=region)
# Supervisor uses the heavier model — it's doing strategic reasoning
supervisor_model = ChatBedrock(
client=bedrock,
model_id="anthropic.claude-3-7-sonnet-20250219-v1:0",
model_kwargs={
"temperature": 0.2,
"max_tokens": 8000,
"thinking": {"type": "enabled", "budget_tokens": 5000}
}
)
def supervisor_node(state: SupervisorState) -> SupervisorState:
if not state.get("subtasks"):
# First pass: decompose the task
response = supervisor_model.invoke([
SystemMessage(content=SUPERVISOR_SYSTEM_PROMPT),
HumanMessage(content=f"Decompose this task into subtasks:\n{state['task_spec']}")
])
import json
try:
subtasks = json.loads(response.content)
if isinstance(subtasks, dict) and "subtasks" in subtasks:
subtasks = subtasks["subtasks"]
except json.JSONDecodeError:
subtasks = []
return {**state, "subtasks": subtasks}
# All workers done: synthesize
results_summary = "\n\n".join([
f"=== {worker} ===\n{result}"
for worker, result in state["worker_results"].items()
])
synthesis_prompt = f"""
Original task: {state['task_spec']}
Worker results:
{results_summary}
Synthesize these into a coherent final report. Highlight conflicts between
worker findings and make clear recommendations.
"""
response = supervisor_model.invoke([
SystemMessage(content=SUPERVISOR_SYSTEM_PROMPT),
HumanMessage(content=synthesis_prompt)
])
return {
**state,
"messages": state["messages"] + [response],
"synthesis_complete": True,
}
return supervisor_node
def route_to_worker(state: SupervisorState) -> str:
"""
Routes to the next worker with unfinished subtasks.
Returns END when all subtasks are complete and synthesis is done.
"""
if state.get("synthesis_complete"):
return END
# Find next unfinished subtask whose dependencies are met
completed = set(state.get("worker_results", {}).keys())
for subtask in state.get("subtasks", []):
sid = subtask["subtask_id"]
if sid in completed:
continue
deps = set(subtask.get("depends_on", []))
if deps.issubset(completed):
return subtask["worker_type"] # route to this worker
# All subtasks done, back to supervisor for synthesis
return "supervisor"
# harness/multi_agent/workers.py
from langchain_aws import ChatBedrock
from langchain_core.messages import SystemMessage, HumanMessage
import boto3
CODE_ANALYST_PROMPT = """
You are a code analysis specialist. You receive specific, bounded analysis tasks.
Focus only on what you were asked to analyze. Do not expand scope.
Return structured findings: what you found, confidence level, and specific evidence.
"""
SECURITY_REVIEWER_PROMPT = """
You are a security review specialist. You look for vulnerabilities, insecure patterns,
and CVE-relevant code. Reference specific CWE numbers when applicable.
Return structured findings with severity levels (CRITICAL, HIGH, MEDIUM, LOW).
"""
TEST_EVALUATOR_PROMPT = """
You are a test coverage specialist. You assess test quality and identify gaps.
Focus on: coverage percentage where available, missing edge cases, and untested paths.
Return structured findings with specific test gaps and suggested test cases.
"""
WORKER_PROMPTS = {
"code_analyst": CODE_ANALYST_PROMPT,
"security_reviewer": SECURITY_REVIEWER_PROMPT,
"test_evaluator": TEST_EVALUATOR_PROMPT,
}
def build_worker(worker_type: str, region: str = "us-east-1") -> callable:
bedrock = boto3.client("bedrock-runtime", region_name=region)
# Workers use the faster model — they execute a specific, bounded task
worker_model = ChatBedrock(
client=bedrock,
model_id="anthropic.claude-3-5-sonnet-20241022-v2:0",
model_kwargs={"temperature": 0, "max_tokens": 4000}
)
system_prompt = WORKER_PROMPTS[worker_type]
def worker_node(state: SupervisorState) -> SupervisorState:
# Find the subtask assigned to this worker type
completed = set(state.get("worker_results", {}).keys())
current_subtask = None
for subtask in state["subtasks"]:
if subtask["worker_type"] == worker_type and subtask["subtask_id"] not in completed:
deps = set(subtask.get("depends_on", []))
if deps.issubset(completed):
current_subtask = subtask
break
if not current_subtask:
return state
# Pass relevant prior results as context if this task has dependencies
context = ""
if current_subtask.get("depends_on"):
for dep_id in current_subtask["depends_on"]:
if dep_id in state.get("worker_results", {}):
context += f"\nPrevious analysis ({dep_id}):\n{state['worker_results'][dep_id]}\n"
prompt = f"Task: {current_subtask['description']}"
if context:
prompt = f"Prior context:{context}\n\n{prompt}"
response = worker_model.invoke([
SystemMessage(content=system_prompt),
HumanMessage(content=prompt)
])
updated_results = {**state.get("worker_results", {})}
updated_results[current_subtask["subtask_id"]] = response.content
return {**state, "worker_results": updated_results}
return worker_node
# harness/multi_agent/supervisor_graph.py
from langgraph.graph import StateGraph, END, START
from langgraph.checkpoint.memory import MemorySaver
from harness.multi_agent.supervisor import SupervisorState, build_supervisor, route_to_worker
from harness.multi_agent.workers import build_worker
def build_supervisor_graph(region: str = "us-east-1"):
graph = StateGraph(SupervisorState)
graph.add_node("supervisor", build_supervisor(region))
graph.add_node("code_analyst", build_worker("code_analyst", region))
graph.add_node("security_reviewer", build_worker("security_reviewer", region))
graph.add_node("test_evaluator", build_worker("test_evaluator", region))
graph.add_edge(START, "supervisor")
graph.add_conditional_edges(
"supervisor",
route_to_worker,
{
"code_analyst": "code_analyst",
"security_reviewer": "security_reviewer",
"test_evaluator": "test_evaluator",
END: END,
}
)
# All workers route back to supervisor after completing their subtask
for worker in ["code_analyst", "security_reviewer", "test_evaluator"]:
graph.add_edge(worker, "supervisor")
return graph.compile(checkpointer=MemorySaver())
模式2:管道编排
在模式2的管道编排阶段,应在修改代码之前明确输入参数、各步骤的负责人以及终止条件。操作人员应能够从已知的检查点重新运行相应步骤,而无需猜测隐藏状态。 需同时记录正常流程与故障恢复流程。重试机制、人工审核环节以及死信处理都是产品功能的一部分,而非后续需要补充的内容。 对于涉及资金支出或修改生产数据的操作,必须经过人工审批。编译时的连接方式并不能保证业务的完整性。
┌───────────────┐ ┌───────────────┐ ┌───────────────┐
│ Agent A │ │ Agent B │ │ Agent C │
│ (extraction) │────>│ (enrichment) │────>│ (validation) │
└───────────────┘ └───────────────┘ └───────────────┘
output A output B output C
becomes input B becomes input C final result
# harness/multi_agent/pipeline.py
from typing import TypedDict, Annotated, List, Optional, Any
from langgraph.graph import StateGraph, END, START
from langchain_core.messages import BaseMessage, HumanMessage, SystemMessage
from langchain_aws import ChatBedrock
import boto3
class PipelineState(TypedDict):
original_input: str
stage_outputs: List[dict] # accumulates each stage's compressed result
current_stage: int
final_output: Optional[str]
agent_run_id: str
def compress_for_handoff(full_output: str, model: ChatBedrock) -> str:
"""
Compresses a stage's full output into a structured handoff summary.
This is the core of pipeline context management — the next agent gets
the substance, not the reasoning trace.
"""
response = model.invoke([
SystemMessage(content="""
Compress the following agent output into a structured handoff summary.
Include: key findings, decisions made, artifacts produced, and what the
next stage needs to know. Discard reasoning traces and intermediate steps.
Target length: 20% of original. Use bullet points for clarity.
"""),
HumanMessage(content=f"Compress this:\n\n{full_output}")
])
return response.content
def build_pipeline_stage(
stage_name: str,
system_prompt: str,
region: str = "us-east-1"
) -> callable:
bedrock = boto3.client("bedrock-runtime", region_name=region)
model = ChatBedrock(
client=bedrock,
model_id="anthropic.claude-3-5-sonnet-20241022-v2:0",
model_kwargs={"temperature": 0, "max_tokens": 6000}
)
# Cheaper model for compression — this is mechanical, not creative
compressor = ChatBedrock(
client=bedrock,
model_id="anthropic.claude-haiku-3-5",
model_kwargs={"temperature": 0, "max_tokens": 2000}
)
def stage_node(state: PipelineState) -> PipelineState:
# Build context from compressed prior stage outputs only
prior_context = ""
for past_stage in state.get("stage_outputs", []):
prior_context += f"\n### {past_stage['stage']} output:\n{past_stage['compressed_output']}\n"
prompt = f"Original task: {state['original_input']}\n"
if prior_context:
prompt += f"\nPrior stage results:\n{prior_context}\n"
prompt += f"\nNow perform your stage: {stage_name}"
response = model.invoke([
SystemMessage(content=system_prompt),
HumanMessage(content=prompt)
])
full_output = response.content
# Compress before storing — next stage won't see raw output
compressed = compress_for_handoff(full_output, compressor)
updated_outputs = list(state.get("stage_outputs", []))
updated_outputs.append({
"stage": stage_name,
"full_output": full_output,
"compressed_output": compressed,
})
return {
**state,
"stage_outputs": updated_outputs,
"current_stage": state.get("current_stage", 0) + 1,
}
return stage_node
模式3:并行扇出/扇入
在模式3的并行分发阶段,修改代码之前需明确输入参数、该步骤的负责人以及终止条件。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏状态。 优先选择小型、可测试的单元,而非庞大的脚本。当某个步骤失败时,故障应指向单一责任点,而非复杂的流程链。 对于涉及资金支出或修改生产数据的操作,必须经过人工审批。编译时的逻辑连接并不等同于业务功能的完整性。 在模式3的并行分发阶段,修改代码之前需明确输入参数、该步骤的负责人以及终止条件。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏状态。 在功能结果之外,还需记录执行时间以及令牌或查询成本。提前了解成本情况,可避免在流程从演示环境转向实际生产环境时出现意外账单。
模式4:辩论/批评。
┌──────────────────┐
│ ORCHESTRATOR │
│ (task splitter) │
└──┬───┬───┬───┬──┘
│ │ │ │
┌──────────▼┐ ┌▼─┐ ┌▼──┐ ┌▼──────────┐
│ Agent 1 │ │A2│ │A3 │ │ Agent 4 │
│ (region A)│ │ │ │ │ │ (region D)│
└──────────┬┘ └┬─┘ └┬──┘ └┬──────────┘
│ │ │ │
┌──▼───▼────▼─────▼──┐
│ AGGREGATOR │
│ (result merger) │
└────────────────────┘
# harness/multi_agent/fanout.py
import asyncio
from typing import List, TypedDict, Annotated, Optional
from langchain_aws import ChatBedrock
from langchain_core.messages import SystemMessage, HumanMessage
import boto3
class FanoutResult(TypedDict):
agent_id: str
input_slice: str
output: str
success: bool
error: Optional[str]
async def run_agent_async(
agent_id: str,
input_slice: str,
system_prompt: str,
model: ChatBedrock,
) -> FanoutResult:
"""Run a single agent asynchronously."""
try:
response = await model.ainvoke([
SystemMessage(content=system_prompt),
HumanMessage(content=input_slice)
])
return FanoutResult(
agent_id=agent_id,
input_slice=input_slice,
output=response.content,
success=True,
error=None,
)
except Exception as e:
return FanoutResult(
agent_id=agent_id,
input_slice=input_slice,
output="",
success=False,
error=str(e),
)
async def fan_out(
task_slices: List[str],
system_prompt: str,
region: str = "us-east-1",
max_concurrent: int = 5, # don't hammer Bedrock rate limits
) -> List[FanoutResult]:
"""
Runs agents concurrently with a semaphore to cap parallelism.
max_concurrent protects against Bedrock throttling — you
will hit rate limits if you fire 20 concurrent requests.
"""
bedrock = boto3.client("bedrock-runtime", region_name=region)
model = ChatBedrock(
client=bedrock,
model_id="anthropic.claude-3-5-sonnet-20241022-v2:0",
model_kwargs={"temperature": 0, "max_tokens": 4000}
)
semaphore = asyncio.Semaphore(max_concurrent)
async def bounded_run(agent_id, slice_content):
async with semaphore:
return await run_agent_async(agent_id, slice_content, system_prompt, model)
tasks = [
bounded_run(f"agent_{i}", slice_content)
for i, slice_content in enumerate(task_slices)
]
return await asyncio.gather(*tasks)
def aggregate_results(
results: List[FanoutResult],
aggregator_model: ChatBedrock,
aggregation_strategy: str = "synthesize",
) -> str:
"""
Merges parallel agent outputs.
aggregation_strategy options:
- "synthesize": ask a model to merge findings coherently
- "concat": simple concatenation (fast, no model call needed)
- "vote": majority-vote for classification tasks
"""
if aggregation_strategy == "concat":
successful = [r for r in results if r["success"]]
return "\n\n---\n\n".join(r["output"] for r in successful)
failed = [r for r in results if not r["success"]]
successful = [r for r in results if r["success"]]
if failed:
# Log partial failures but don't crash — partial results are usually useful
for f in failed:
print(f"Agent {f['agent_id']} failed: {f['error']}")
outputs_for_synthesis = "\n\n".join([
f"[Agent {r['agent_id']}]:\n{r['output']}"
for r in successful
])
response = aggregator_model.invoke([
SystemMessage(content="""
You are a results aggregator. You receive outputs from multiple parallel agents
that each analyzed a different slice of the same problem. Your job is to:
1. Identify common findings across agents
2. Surface unique findings from individual agents
3. Flag any contradictions between agents and explain which to trust
4. Produce a single coherent output as if one expert had analyzed everything
Do not simply concatenate. Actively synthesize.
"""),
HumanMessage(content=f"Synthesize these {len(successful)} agent outputs:\n\n{outputs_for_synthesis}")
])
return response.content
# Convenience wrapper for synchronous callers
def run_parallel_analysis(
task_slices: List[str],
system_prompt: str,
region: str = "us-east-1",
) -> str:
results = asyncio.run(fan_out(task_slices, system_prompt, region))
bedrock = boto3.client("bedrock-runtime", region_name=region)
aggregator = ChatBedrock(
client=bedrock,
model_id="anthropic.claude-3-7-sonnet-20250219-v1:0",
model_kwargs={"temperature": 0, "max_tokens": 8000}
)
return aggregate_results(results, aggregator)
模式4:辩论/批评
在处理模式4的辩论与批评阶段时,首先写下相关约定:所需的输入、成功标志以及部分失败时的处理方式。这样的清单能确保后续的代码修改保持一致性。 将配置信息置于应用程序代码之外。环境文件、密钥存储以及功能开关应集中存放于一个位置,这样操作人员无需查看整个系统结构即可进行审计。 在耗时较高的步骤之后设置检查点。当操作人员重新尝试后续节点时,恢复流程不应再次调用相同的大型语言模型。
┌─────────────┐ ┌─────────────┐
│ Proposer │ │ Challenger │
│ (solution │ │ (solution │
│ attempt 1)│ │ attempt 2)│
└──────┬──────┘ └──────┬──────┘
│ │
└──────────┬─────────────┘
▼
┌─────────────────┐
│ ADJUDICATOR │
│ (critique + │
│ synthesis) │
└─────────────────┘
# harness/multi_agent/debate.py
from dataclasses import dataclass
from typing import Optional
from langchain_aws import ChatBedrock
from langchain_core.messages import SystemMessage, HumanMessage
import boto3
@dataclass
class DebateResult:
proposer_solution: str
challenger_solution: str
adjudication: str
final_recommendation: str
agreement_level: str # HIGH | MEDIUM | LOW | CONTRADICTION
PROPOSER_PROMPT = """
You are a solution proposer. Approach the problem carefully and produce your best
solution. Explain your reasoning. Do not hedge excessively — commit to a specific answer.
"""
CHALLENGER_PROMPT = """
You are a solution challenger. You will receive a problem that another agent has
already attempted. Produce your own independent solution WITHOUT seeing their work.
Approach this fresh. Your goal is not to contradict — it's to find the best solution.
"""
ADJUDICATOR_PROMPT = """
You are an adjudicator reviewing two independent solutions to the same problem.
Your job:
1. Identify where the two solutions agree (these are likely correct)
2. Identify where they diverge (these need careful evaluation)
3. For each divergence, evaluate which solution is better and why
4. Produce a final synthesis that takes the best of both
Be direct about contradictions. Do not smooth over genuine disagreements —
surface them clearly so the human reviewer can make a judgment call.
Rate the agreement level: HIGH (minor differences), MEDIUM (some significant
divergences), LOW (fundamentally different approaches), or CONTRADICTION
(mutually exclusive conclusions).
"""
def run_debate(
problem: str,
region: str = "us-east-1",
) -> DebateResult:
bedrock = boto3.client("bedrock-runtime", region_name=region)
heavy_model = ChatBedrock(
client=bedrock,
model_id="anthropic.claude-3-7-sonnet-20250219-v1:0",
model_kwargs={
"temperature": 0.3, # slight temperature for independent solutions
"max_tokens": 6000,
"thinking": {"type": "enabled", "budget_tokens": 4000}
}
)
# Proposer works the problem
proposer_response = heavy_model.invoke([
SystemMessage(content=PROPOSER_PROMPT),
HumanMessage(content=problem)
])
proposer_solution = proposer_response.content
# Challenger works the same problem independently
# Note: Challenger does NOT see Proposer's solution
challenger_response = heavy_model.invoke([
SystemMessage(content=CHALLENGER_PROMPT),
HumanMessage(content=problem)
])
challenger_solution = challenger_response.content
# Adjudicator sees both and synthesizes
adjudicator_response = heavy_model.invoke([
SystemMessage(content=ADJUDICATOR_PROMPT),
HumanMessage(content=f"""
Problem: {problem}
Solution A (Proposer):
{proposer_solution}
Solution B (Challenger):
{challenger_solution}
Adjudicate and synthesize.
""")
])
adjudication = adjudicator_response.content
# Extract agreement level from adjudication
agreement_level = "MEDIUM"
for level in ["CONTRADICTION", "LOW", "HIGH", "MEDIUM"]:
if level in adjudication.upper():
agreement_level = level
break
return DebateResult(
proposer_solution=proposer_solution,
challenger_solution=challenger_solution,
adjudication=adjudication,
final_recommendation=adjudication,
agreement_level=agreement_level,
)
智能体间交接与上下文传递
在处理“智能体间交接与上下文传递”阶段时,首先写下相关契约:所需的输入参数、成功信号以及部分失败时的处理方式。这样的检查清单能确保后续的代码修改保持一致性。 同时记录正常流程和异常恢复流程。重试机制、人工审核环节以及死信处理都是产品功能的一部分,而非后续需要补充的内容。 在成本较高的操作之后设置检查点。当操作员重新尝试某个节点时,恢复流程不应再次调用相同的大型语言模型。
# harness/multi_agent/handoff.py
from dataclasses import dataclass, field
from typing import Any, Optional
from langchain_core.messages import HumanMessage
@dataclass
class AgentHandoff:
"""
Structured context passed between agents.
The split between result and trace is deliberate: downstream agents
need the result, not the full reasoning history. Keeping them separate
lets each agent decide how much context it wants to consume.
"""
source_agent: str
task_completed: str
result_summary: str # compressed, structured result
artifacts: dict = field(default_factory=dict) # files, code, structured data
reasoning_trace: Optional[str] = None # full trace, passed only if downstream needs it
confidence: str = "MEDIUM" # HIGH | MEDIUM | LOW
flags: list = field(default_factory=list) # NEEDS_REVIEW, PARTIAL_RESULT, etc.
def to_context_message(self, include_trace: bool = False) -> HumanMessage:
"""
Converts handoff to a HumanMessage for injection into next agent's context.
include_trace=True only when the downstream agent genuinely needs the reasoning.
"""
content = f"""
[HANDOFF FROM: {self.source_agent}]
Task completed: {self.task_completed}
Confidence: {self.confidence}
Flags: {', '.join(self.flags) if self.flags else 'none'}
Result summary:
{self.result_summary}
"""
if self.artifacts:
content += f"\nArtifacts available:\n"
for key, value in self.artifacts.items():
if isinstance(value, str) and len(value) < 500:
content += f" {key}: {value}\n"
else:
content += f" {key}: [available, {type(value).__name__}]\n"
if include_trace and self.reasoning_trace:
content += f"\nFull reasoning trace:\n{self.reasoning_trace}"
return HumanMessage(content=content)
生产环境强化措施
在开展生产环境强化阶段时,首先需明确相关规范:所需输入、成功标志以及部分失败时的处理方式。这份清单能确保后续的代码修改始终符合要求。 建议使用小型、可测试的单元,而非庞大的脚本。当某个步骤失败时,故障应指向单一责任模块,而非复杂的流程链。 在成本较高的步骤之后设置检查点。当操作员重新尝试后续节点时,恢复流程不应再次调用相同的大型语言模型。 在开展生产环境强化阶段时,首先需明确相关规范:所需输入、成功标志以及部分失败时的处理方式。这份清单能确保后续的代码修改始终符合要求。 在功能结果旁记录执行时间以及Token或查询成本。提前了解成本情况,可避免从演示环境过渡到共享环境时出现意外收费。
防止成本激增
将“成本激增预防”阶段视为可度量的对象来处理时,其效果最佳。在扩大范围之前,先收集一份典型的成功案例、一个故障实例以及回滚说明。 将配置与应用程序代码分开。环境文件、密钥存储和功能标志应集中存放于一处,这样操作人员无需查看整个结构就能进行审计。 保持图结构的扁平化与类型化。嵌套的数据块会掩盖哪个节点修改了哪个字段的信息,且在中断后会导致恢复失败。
# harness/multi_agent/budget.py
import boto3
import time
from decimal import Decimal
class AgentBudgetGuard:
"""
Tracks cumulative cost and agent count per run.
Hard-stops execution when limits are exceeded.
"""
def __init__(
self,
max_agents_per_run: int = 10,
max_total_tokens: int = 500_000,
table_name: str = "agent-budget-state",
region: str = "us-east-1",
):
self.max_agents = max_agents_per_run
self.max_tokens = max_total_tokens
self.table = boto3.resource("dynamodb", region_name=region).Table(table_name)
def register_agent_spawn(self, run_id: str, agent_id: str) -> bool:
"""
Returns True if spawn is allowed, False if budget exceeded.
Call this before spawning any sub-agent.
"""
response = self.table.update_item(
Key={"run_id": run_id},
UpdateExpression="SET agent_count = if_not_exists(agent_count, :z) + :inc",
ExpressionAttributeValues={":z": 0, ":inc": 1},
ReturnValues="UPDATED_NEW",
)
new_count = int(response["Attributes"]["agent_count"])
if new_count > self.max_agents:
raise AgentBudgetExceededError(
f"Run {run_id} attempted to spawn agent #{new_count}, "
f"but max_agents_per_run is {self.max_agents}. "
f"Either the supervisor is over-decomposing, or there is a spawn loop."
)
return True
def record_token_usage(self, run_id: str, tokens_used: int):
response = self.table.update_item(
Key={"run_id": run_id},
UpdateExpression="SET total_tokens = if_not_exists(total_tokens, :z) + :inc",
ExpressionAttributeValues={":z": 0, ":inc": tokens_used},
ReturnValues="UPDATED_NEW",
)
total = int(response["Attributes"]["total_tokens"])
if total > self.max_tokens:
raise AgentBudgetExceededError(
f"Run {run_id} consumed {total:,} tokens, exceeding limit of {self.max_tokens:,}."
)
class AgentBudgetExceededError(Exception):
pass
死锁检测
将死锁检测阶段视为可度量的对象时,其效果最佳。在扩大范围之前,先记录一份理想的操作日志、一个故障案例以及回滚说明。 同时记录正常流程和恢复流程。重试机制、人工干预环节以及死信处理都是产品本身的组成部分,而非后续需要补充的功能。 保持图结构简洁且类型明确。嵌套的数据块会掩盖哪个节点修改了哪个字段的信息,还会在中断后导致流程无法继续。
# harness/multi_agent/deadlock.py
import boto3
import time
from typing import List
class DeadlockDetector:
"""
Tracks the delegation chain per run and detects cycles.
Stored in DynamoDB so it works across parallel agent branches.
"""
def __init__(self, table_name: str = "agent-delegation-chain", region: str = "us-east-1"):
self.table = boto3.resource("dynamodb", region_name=region).Table(table_name)
def record_delegation(self, run_id: str, from_agent: str, to_agent: str):
"""
Records a delegation event and checks for cycles.
Raises DeadlockDetectedError if a cycle is found.
"""
self.table.put_item(Item={
"run_id": run_id,
"delegation_id": f"{from_agent}->{to_agent}-{int(time.time())}",
"from_agent": from_agent,
"to_agent": to_agent,
"timestamp": int(time.time()),
})
chain = self._get_delegation_chain(run_id)
if self._has_cycle(chain):
cycle_description = self._describe_cycle(chain)
raise DeadlockDetectedError(
f"Delegation cycle detected in run {run_id}: {cycle_description}. "
f"Check supervisor decomposition logic for circular dependencies."
)
def _get_delegation_chain(self, run_id: str) -> List[tuple]:
response = self.table.query(
KeyConditionExpression="run_id = :rid",
ExpressionAttributeValues={":rid": run_id}
)
return [(item["from_agent"], item["to_agent"]) for item in response.get("Items", [])]
def _has_cycle(self, chain: List[tuple]) -> bool:
graph = {}
for from_a, to_a in chain:
graph.setdefault(from_a, set()).add(to_a)
visited, rec_stack = set(), set()
def dfs(node):
visited.add(node)
rec_stack.add(node)
for neighbor in graph.get(node, []):
if neighbor not in visited:
if dfs(neighbor): return True
elif neighbor in rec_stack:
return True
rec_stack.discard(node)
return False
return any(dfs(node) for node in graph if node not in visited)
def _describe_cycle(self, chain: List[tuple]) -> str:
return " -> ".join(f"{f}->{t}" for f, t in chain[-5:])
class DeadlockDetectedError(Exception):
pass
跨代理可观测性
将跨代理可观测性阶段视为一个可度量的对象来处理,效果最佳。在扩大范围之前,先记录一份理想的日志、一个故障案例以及回滚说明。 优先选择小型且可测试的单元,而非庞大的脚本。当某个步骤出现故障时,故障应指向单一责任主体,而非复杂的流程链。 保持图结构简洁且类型明确。嵌套的数据块会掩盖哪个节点编写了哪个字段的信息,还会在流程中断后导致无法继续执行。 将跨代理可观测性阶段视为一个可度量的对象来处理,效果最佳。在扩大范围之前,先记录一份理想的日志、一个故障案例以及回滚说明。 在功能结果之外,还需记录执行时间以及令牌或查询成本。提前了解成本情况,可避免在从演示环境过渡到共享环境时出现意外费用。
# harness/multi_agent/observability.py
import os
from contextlib import contextmanager
from langsmith import Client
from langsmith.run_trees import RunTree
class MultiAgentTracer:
"""
Maintains a run tree across all agents in a multi-agent system.
Pass the parent_run_id to each agent so their traces nest correctly.
"""
def __init__(self):
self.client = Client()
@contextmanager
def agent_span(self, parent_run_id: str, agent_name: str, inputs: dict):
"""
Context manager for an individual agent's trace span.
Usage:
with tracer.agent_span(parent_run_id, "security_reviewer", {...}) as span:
result = run_security_review(...)
span.end(outputs={"result": result})
"""
run = self.client.create_run(
name=agent_name,
run_type="chain",
inputs=inputs,
parent_run_id=parent_run_id,
)
try:
yield run
except Exception as e:
self.client.update_run(run.id, error=str(e))
raise
finally:
self.client.update_run(run.id, end_time=None) # auto-sets end time
生产环境现实检验
在“生产环境现实检验”阶段,应在修改代码之前明确输入参数、该步骤的负责人以及结束标准。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏状态。 配置信息应置于应用程序代码之外。环境文件、密钥存储以及功能开关应集中存放于一个位置,以便操作人员无需查看整个系统结构即可进行审计。 对于涉及资金支出或修改生产数据的操作,必须经过人工审批。编译时的连接方式并不等同于业务功能的完整性。
参考架构
在参考架构阶段,应在修改代码之前明确输入内容、各步骤的负责人以及结束标准。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏状态。 需同时记录正常流程和故障恢复流程。重试机制、人工审核环节以及错误处理都是产品不可或缺的部分,而非后续需要补充的内容。 对于涉及资金支出或修改生产数据的操作,必须经过人工审批。编译时的配置并不等同于业务功能的完整性。
┌───────────────────────────────────────────┐
│ User Request │
└──────────────────┬────────────────────────┘
│
┌──────────────────▼────────────────────────┐
│ AgentHarness Runtime │
│ (budget guard, deadlock detector, tracer)│
└──────────────────┬────────────────────────┘
│
┌──────────────────▼────────────────────────┐
│ SUPERVISOR / ORCHESTRATOR │
│ Claude 3.7 + extended thinking │
│ Task decomposition + result synthesis │
└────┬──────────────┬──────────────┬────────┘
│ │ │
┌───────────▼──┐ ┌────────▼──┐ ┌───────▼───────┐
│ Worker A │ │ Worker B │ │ Worker C │
│ Claude 3.5 │ │ Claude 3.5 │ │ Claude 3.5 │
│ specialist │ │ specialist │ │ specialist │
└───────┬──────┘ └─────┬─────┘ └──────┬────────┘
│ │ │
┌───────▼───────────────▼────────────────▼────────┐
│ Tool Execution Layer │
│ (auth, retry, circuit breaker from Art.5) │
└───────────────────────┬─────────────────────────┘
│
┌───────────────────────▼──────────────────────────┐
│ AWS Services │
│ Bedrock │ DynamoDB (budget+deadlock+loop+circ) │
│ Secrets Manager │ Knowledge Bases │
└──────────────────────────────────────────────────┘
Observability: LangSmith run trees — full hierarchy per user request
参考基础设施栈
在参考基础设施堆栈阶段,应在修改代码之前明确输入参数、该步骤的负责人以及完成标准。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏状态。相比庞大的脚本,更应采用小型且可测试的单元。当某个步骤失败时,故障原因应能指向单一责任主体,而非复杂的流程链。对于涉及资金支出或修改生产数据的操作,必须经过人工审批。编译时的配置并不等同于业务功能的完整性。
+-----------------------------+---------------------+------------------------------+
| Component | Technology | Role |
+-----------------------------+---------------------+------------------------------+
| Orchestration | LangGraph 0.2+ | Multi-agent graph, routing, |
| | | supervisor-worker topology |
+-----------------------------+---------------------+------------------------------+
| Supervisor Model | Claude 3.7 Sonnet | Task decomposition, |
| | (extended thinking) | result synthesis |
+-----------------------------+---------------------+------------------------------+
| Worker Models | Claude 3.5 Sonnet | Specialized execution, |
| | | bounded tasks |
+-----------------------------+---------------------+------------------------------+
| Parallel Execution | asyncio + Bedrock | Concurrent agent runs with |
| | | semaphore-gated concurrency |
+-----------------------------+---------------------+------------------------------+
| Context Compression | Claude Haiku 3.5 | Pipeline stage handoffs, |
| | | summary generation |
+-----------------------------+---------------------+------------------------------+
| Budget Guard | DynamoDB | Agent count + token limits |
| | | per run |
+-----------------------------+---------------------+------------------------------+
| Deadlock Detection | DynamoDB | Delegation cycle detection |
+-----------------------------+---------------------+------------------------------+
| Loop Detection | DynamoDB (Art. 5) | Per-resource edit tracking |
+-----------------------------+---------------------+------------------------------+
| Circuit Breaker State | DynamoDB (Art. 5) | Shared across all agents |
| | | in a run |
+-----------------------------+---------------------+------------------------------+
| Cross-Agent Observability | LangSmith run trees | Full hierarchy per request |
+-----------------------------+---------------------+------------------------------+
| Auth Propagation | CredentialManager | JWT passed to all workers |
| | (Art. 5) | via execution context |
+-----------------------------+---------------------+------------------------------+
| Local Dev Alternative | Ollama + Docker | All patterns testable |
| | Compose | without Bedrock costs |
+-----------------------------+---------------------+------------------------------+
| Infrastructure as Code | Terraform | DynamoDB tables, IAM roles |
+-----------------------------+---------------------+------------------------------+
关于我们未来发展方向的一些说明
在“A阶段说明”中,应在修改代码之前明确输入参数、该步骤的负责人以及结束标准。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏状态。 将此阶段视为输入与已验证输出之间的契约。为相关成果命名,定义成功判定标准,并拒绝默许的半完成状态。 对于涉及资金支出或修改生产数据的操作,必须经过人工审批。编译时的连接方式并不等同于业务上的完整性。
操作检查清单
在“操作检查清单”阶段,应在修改代码之前明确输入参数、该步骤的负责人以及结束标准。操作人员应能够从已知的检查点重新运行该步骤,而无需猜测隐藏状态。
将此阶段视为输入与经过验证的输出之间的契约。为相关成果命名,明确成功标准,绝不允许默默地仅完成部分工作。
对于会耗费资金或修改生产数据的操作,必须经过人工审批。编译时的连接方式并不等同于业务上的完整性。
编写简短的操作手册:说明如何轮换密钥、如何清空队列、以及如何回滚上一次的导入操作。
在功能结果旁记录处理时间以及令牌或查询成本。提前了解成本情况,可避免在系统从演示环境过渡到共享环境时出现意外账单。
对于会耗费资金或修改生产数据的操作,必须经过人工审批。编译时的连接方式并不等同于业务上的完整性。
在推广该技术栈之前,应先冻结版本,为关键流程记录标准输出日志,并明确回滚步骤。共享环境需要设置速率限制、租户验证机制,以及负责密钥轮换的明确责任人。与其展示花哨的一次性演示,不如注重扎实的可靠性。
a0dc7ff1211b的批处理说明:不要将提供商密钥放入代码仓库,为每个会话设置令牌使用上限,并将日志存储在评估用示例文件旁边,以便后续模型更换时仍能保持对比性。