parallel-web-research¶
第10章 · 多 Agent 协作 · 配套项目
chapter10/parallel-web-research
项目说明¶
实验 10-6 · 同时从多个网站搜集信息的 Agent(★★)¶
《深入理解 AI Agent》配套实验。演示多个同构 Agent 的并行搜索 + 中心协调: 主协调器同时启动 N 个子 Agent,每个子 Agent 访问一个"网站/来源"找答案; 一旦某个子 Agent 命中目标,其余立即优雅停止。
书中的原型是"10 个并行的 Computer Use Agent 同时访问不同网站找信息"。为便于 自动验证,本实验不启动真实浏览器,而是用一批可控的模拟信息源代替;把重点 完整放在协调机制上——消息总线、并行派发、实时监控、级联终止、竞态处理均为真实实现。
目录结构¶
| 文件 | 作用 |
|---|---|
message_bus.py |
进程内异步消息总线(Redis Pub/Sub 风格),带 Envelope 信封与订阅机制 |
sources.py |
模拟的 10 个"网站/来源",各有不同延迟;可控的关键词命中判断;build_sources(n) 支持任意并行度 |
llm.py |
可选的 LLM 判断层(默认离线关键词判断,配了 key 则用真实大模型) |
agents.py |
主协调器 Coordinator 与子 Agent WorkerAgent(核心协调逻辑);run_sequential 串行基线 |
demo.py |
一条命令的演示入口,带 argparse CLI 与末尾断言式自检 |
架构与机制¶
┌────────────────────────────┐
│ Coordinator(主协调器) │
│ · 并行派发 task_assigned │
│ · 维护任务状态表(状态机) │
│ · 首个命中→加锁结算(幂等) │
│ · 广播一轮 terminate │
└───────────┬────────────────┘
│
┌────────────────┴─────── MessageBus(异步消息总线)───────┐
│ Envelope{ sender_id, target, type, payload, seq, ts } │
│ type: task_assigned/status_update/result/terminate/ack │
└──┬─────────┬─────────┬─────────┬──────────────┬──────────┘
│ │ │ │ │
┌──▼──┐ ┌──▼──┐ ┌──▼──┐ ┌──▼──┐ ... ┌──▼──┐
│W-00 │ │W-01 │ │W-02 │ │W-03 │ │W-09 │ 子 Agent
│网站A│ │网站B│ │网站C│ │网站D│ │网站J│ (同构·并行)
└─────┘ └─────┘ └─────┘ └─────┘ └─────┘
对应书中强调的五个机制:
- 消息总线(Message Bus):所有通信都是发布到总线的带信封消息,按
type订阅。日志中每条BUS ...行就是一次发布/投递(Redis Pub/Sub 语义的进程内实现)。 - 并行派发:
Coordinator同时给 10 个子 Agent 发task_assigned并asyncio.create_task并发执行。 - 实时监控(push 范式):子 Agent 执行中主动
status_update上报进度, 主 Agent 维护任务状态表并在状态机跳变时实时刷新打印。 状态机:已提交 → 执行中 →(需要输入)→ 已完成 / 失败 / 已终止。 - 级联终止:某子 Agent 命中后,主 Agent 广播
terminate;其余子 Agent 在 循环的安全检查点发现信号后回ack并优雅退出(状态置为"已终止")。 - 竞态处理:多个子 Agent 可能几乎同时命中,主 Agent 用
asyncio.Lock+ 幂等标志_settled保证只结算一次、只广播一轮终止;迟到的命中被记录并忽略。
为让"竞态""级联终止"可复现,各来源被赋予不同的模拟延迟,其中
geo-journal与forum-qa两个正确源被设成相同延迟,从而稳定地在同一时刻命中、触发竞态。
运行¶
cd chapter10/parallel-web-research
pip install -r requirements.txt # 仅离线演示的话可跳过,纯标准库即可运行
python demo.py
默认走离线关键词判断(无需联网、结果可复现)。若要让子 Agent 用真实 LLM 判断:
cp env.example .env
# 在 .env 填入 OPENAI_API_KEY(也支持 Moonshot / 火山方舟 ARK 的 OpenAI 兼容网关)
python demo.py
# 或不改 .env,直接用命令行开关(仍需配置 key 才会真正生效):
python demo.py --use-llm --model gpt-5.6-luna
可用 key:OPENAI_API_KEY(默认模型 gpt-5.6-luna)/ MOONSHOT_API_KEY / ARK_API_KEY
(填到 OPENAI_API_KEY 并按需设置 OPENAI_BASE_URL、OPENAI_MODEL)。
通用回退:若未设置 OPENAI_API_KEY 但设了 OPENROUTER_API_KEY,则真实 LLM 判断
自动改走 OpenRouter,并把模型名映射到其命名空间(gpt-5.6-luna → openai/gpt-5.6-luna)。
不影响协调机制,仅改变"是否命中"的判断。
命令行参数¶
python demo.py --help 可查看完整帮助。所有参数都不改变默认行为——不传任何参数即为
原有的「10 个 Agent + 内置问题 + 离线可复现 + 详细 BUS 日志」演示。
| 参数 | 作用 | 默认 |
|---|---|---|
-q, --query 问题 |
研究问题(离线关键词判断是针对内置来源调校的,自定义问题一般需搭配 --use-llm) |
内置问题 |
-n, --agents N |
并行子 Agent 数量(N≥2 时始终包含两个含答案的源以稳定演示竞态/级联终止) | 10 |
--model MODEL |
LLM 模型名(等价于设 OPENAI_MODEL,仅 --use-llm 且配置 key 时生效) |
环境变量 |
-o, --output PATH |
把结论(含并行/串行耗时、winner、竞态统计)写入 JSON 文件 | 不写 |
--compare |
并行跑完后再实测一遍串行基线,打印墙钟耗时对比 | 关闭 |
--use-llm |
强制真实 LLM 判断(仍需配 OPENAI_API_KEY 或 OPENROUTER_API_KEY,否则自动回退离线判断) |
关闭 |
--quiet |
减少逐条 BUS 日志(状态表/结论/自检不受影响) | 关闭 |
python demo.py --agents 6 --compare # 6 个并行 Agent,并对比串行墙钟耗时
python demo.py --output result.json # 结论落盘为 JSON
并行 vs 串行的墙钟收益(--compare)¶
对应书中实验要求「记录并对比并行/串行时间差异」。--compare 会在并行演示之后,用完全相同的
来源集合再跑一遍串行基线(逐个 await source.fetch() + 判断,命中即止),耗时是实测而非
估算。示例输出(默认 10 源,离线判断):
5) 并行执行墙钟耗时:1.57s(含收敛静默期)
------------------------------------------------------------------------------
并行 vs 串行 墙钟对比(--compare,串行基线为实测)
------------------------------------------------------------------------------
串行:命中前逐个抓取了 3/10 个源,墙钟耗时 2.60s,winner=geo-journal
并行:墙钟耗时 1.57s,winner=worker-02
加速比 ≈ 1.66×,节省约 1.03s(并行让最快的源立即结束全局搜索)。
串行必须依次抓完 baike-wiki→news-portal 才轮到最快命中的 geo-journal(累计 2.6s);并行则让 所有源同时开跑,最快的源一命中就触发级联终止、立刻结束全局搜索。并行墙钟包含级联终止的收敛 开销,因此不是理想的「首个源延迟」,但仍显著快于串行——这正是并行 + 级联终止的价值所在。
演示说明什么(真实运行输出关键片段)¶
(a) 消息总线的发布/订阅在工作——每条带信封的消息都打印出来:
BUS [t= 0.00s #3 ] coordinator -> worker-02 | task_assigned | {"question": "...", "source": "geo-journal"}
BUS [t= 0.00s #13 ] worker-02 -> coordinator | status_update | {"state": "执行中", ...}
(b) N 个子 Agent 并行执行 + 主 Agent 实时刷新状态表:
── 任务状态表(worker-02 -> 执行中) ──
worker-00 源=baike-wiki 状态=执行中 | 开始抓取来源
worker-02 源=geo-journal 状态=执行中 | 开始抓取来源
...
(c) 级联终止——命中后广播 terminate,其余子 Agent ack 并优雅退出:
BUS [t=0.60s #41 ] coordinator -> ALL | terminate | {"reason":"answer_found","winner":"worker-02"}
BUS [t=0.67s #43 ] worker-09 -> coordinator | ack | {"acked":"terminate","source":"map-service"}
[ack] worker-09 已确认终止(1 个已 ack)
...最终 8 个未命中的 Worker 全部状态=已终止
(d) 竞态:即使几乎同时命中,也只结算一次、只广播一轮终止:
BUS [t=0.60s #37 ] worker-02 -> coordinator | result | {"found":true, "answer":"...珠穆朗玛峰...8848 米..."}
BUS [t=0.60s #38 ] worker-04 -> coordinator | result | {"found":true, "answer":"...珠穆朗玛峰...8848.86 米..."}
[结算] 首个命中来自 worker-02 —— 加锁结算,广播一轮 terminate。
[竞态] worker-04 也命中,但已由 worker-02 结算 —— 忽略此次命中,不重复广播终止。
demo.py 末尾有断言式自检:terminate 广播轮数 == 1、只结算一次 == True、
winner 非空。跑通即证明机制正确:
局限与注意事项¶
本实验把重点放在协调机制(消息总线/并行派发/级联终止/竞态处理)上,这些均为真实 实现;但为了可离线运行与自动验证,以下三处做了简化,是已知局限:
- 局限·模拟源非真实浏览器:不启动真实浏览器,来源是可控的模拟数据 + 延迟。若要接
真实 Computer Use,只需把
WorkerAgent.run()里的"抓取一步 + 判断"换成真实浏览器 操作,协调层无需改动。 - 局限·竞态靠相同延迟稳定复现:
geo-journal与forum-qa两个正确源被人为设成 相同延迟,才能稳定触发"同时命中";真实环境里竞态是偶发的,但加锁 + 幂等的结算逻辑 对偶发竞态同样成立,不依赖这个人为设定。 - 局限·进程内总线非真实 Redis:
MessageBus用进程内 async 队列模拟 Redis Pub/Sub,无需真装 Redis;语义(信封、按类型订阅、点对点/广播投递)一致,但不具备 跨进程/跨机器能力,便于单机复现与自动验证。 - 子 Agent 在循环里定期检查终止信号(每步抓取前后),因此终止是"安全点响应"而非 强杀,能保证资源被正常收尾。
源代码¶
agents.py¶
"""
主协调器(Coordinator)与子 Agent(Worker)
===========================================
实现实验 10-6 的核心协调机制:
1. 并行派发:Coordinator 同时启动 N 个同构 Worker,各搜一个"网站/来源"。
2. 消息总线:Worker 与 Coordinator 全部通过 ``MessageBus`` 用信封通信。
3. 实时监控(push 范式):Worker 执行中主动 ``status_update`` 上报,
Coordinator 维护任务状态表并实时刷新打印。
4. 级联终止:某 Worker 命中目标后,Coordinator 广播 ``terminate``,
其余 Worker 在循环安全点检查到信号后 ack 并优雅退出。
5. 竞态处理:多个 Worker 可能几乎同时命中,Coordinator 用 ``asyncio.Lock``
+ 幂等标志保证**只结算一次、只广播一轮终止**。
状态机:submitted -> running -> (needs_input) -> succeeded / failed / terminated
"""
from __future__ import annotations
import asyncio
import time
from dataclasses import dataclass, field
from enum import Enum
from typing import Dict, List, Optional
from llm import judge_answer, llm_available
from message_bus import BROADCAST, MessageBus, _now
from sources import Source
class TaskState(str, Enum):
SUBMITTED = "已提交"
RUNNING = "执行中"
NEEDS_INPUT = "需要输入"
SUCCEEDED = "已完成"
FAILED = "失败"
TERMINATED = "已终止"
@dataclass
class TaskRecord:
"""Coordinator 状态表里的一行:一个 Worker 的实时状态。"""
worker_id: str
source_name: str
state: TaskState = TaskState.SUBMITTED
note: str = ""
updated: float = field(default_factory=_now)
# ————————————————————————————— 子 Agent —————————————————————————————
class WorkerAgent:
"""
一个同构子 Agent:负责抓取并搜索单个来源。
通过消息总线接收 task_assigned / terminate,上报 status_update / result / ack。
"""
def __init__(self, worker_id: str, source: Source, bus: MessageBus, question: str):
self.id = worker_id
self.source = source
self.bus = bus
self.question = question
# 订阅:只关心发给自己或广播的 task_assigned 与 terminate
self.sub = bus.subscribe(worker_id, types=["task_assigned", "terminate"])
self._terminated = asyncio.Event()
async def _report(self, state: TaskState, note: str = "", type: str = "status_update"):
"""向 Coordinator 推送一条状态更新(实时监控的 push 范式)。"""
await self.bus.send(
self.id,
"coordinator",
type,
{"state": state.value, "note": note, "source": self.source.name},
)
async def _drain_signals(self) -> bool:
"""
在安全点检查是否收到 terminate 信号(非阻塞)。
收到则 ack 并返回 True,调用方据此优雅退出。
"""
while not self.sub.inbox.empty():
env = self.sub.inbox.get_nowait()
if env.type == "terminate":
self._terminated.set()
if self._terminated.is_set():
await self.bus.send(
self.id, "coordinator", "ack",
{"acked": "terminate", "source": self.source.name},
)
await self._report(TaskState.TERMINATED, "收到终止信号,安全退出")
return True
return False
async def run(self):
"""
Worker 主循环:分多步"抓取 + 搜索",每步之间都检查终止信号。
把单次抓取切成多步,是为了模拟真实 Computer Use Agent 的多轮操作,
也给级联终止提供"安全检查点"。
"""
# 等待 Coordinator 派发任务(task_assigned)
env = await self.sub.get()
while env.type != "task_assigned":
env = await self.sub.get()
await self._report(TaskState.RUNNING, "开始抓取来源")
steps = 3 # 把抓取拆成 3 步,制造可被中断的检查点
per_step = max(self.source.latency / steps, 0.05)
collected = ""
for i in range(1, steps + 1):
# —— 安全检查点:先看有没有被要求终止 ——
if await self._drain_signals():
return
# —— 执行一步抓取(模拟 Computer Use 的一轮操作耗时)——
await asyncio.sleep(per_step)
collected = self.source.content # 抓到的文本,最后一步再判断是否命中
await self._report(TaskState.RUNNING, f"抓取进度 {i}/{steps}")
# 再次检查终止(可能在最后一步耗时里收到)
if await self._drain_signals():
return
# —— 用(可选)LLM 或关键词判断是否命中答案 ——
answer = await judge_answer(self.question, collected)
if answer:
# 命中:把结果发回 Coordinator(可能与别的 Worker 竞态)
await self.bus.send(
self.id, "coordinator", "result",
{"found": True, "answer": answer, "source": self.source.name},
)
await self._report(TaskState.SUCCEEDED, f"命中:{answer}")
else:
await self.bus.send(
self.id, "coordinator", "result",
{"found": False, "source": self.source.name},
)
await self._report(TaskState.FAILED, "该来源未找到答案")
# ————————————————————————————— 主协调器 —————————————————————————————
class Coordinator:
"""
中心协调器:并行派发子 Agent、维护状态表、结算首个命中、广播级联终止。
"""
def __init__(self, bus: MessageBus, question: str):
self.bus = bus
self.question = question
# 订阅所有子 Agent 上报的消息类型
self.sub = bus.subscribe("coordinator", types=["status_update", "result", "ack"])
self.table: Dict[str, TaskRecord] = {}
self.workers: List[WorkerAgent] = []
# —— 竞态处理的关键状态 ——
self._settle_lock = asyncio.Lock() # 保证结算与终止广播互斥
self._settled = False # 幂等标志:是否已结算过
self.winner: Optional[str] = None # 第一个命中的 Worker
self.answer: Optional[str] = None
self.duplicate_hits: List[str] = [] # 记录"迟到的命中",证明竞态被正确忽略
self._acks: set[str] = set()
self._expected_workers = 0
def add_worker(self, worker: WorkerAgent):
self.workers.append(worker)
self.table[worker.id] = TaskRecord(worker.id, worker.source.name)
def _print_table(self, reason: str):
"""实时刷新并打印任务状态表。"""
print(f"\n ── 任务状态表({reason}) ──")
for rec in self.table.values():
print(
f" {rec.worker_id:<10} 源={rec.source_name:<12} "
f"状态={rec.state.value:<5} {('| ' + rec.note) if rec.note else ''}"
)
print()
async def _dispatch(self):
"""并行派发:给每个 Worker 发 task_assigned。"""
self._expected_workers = len(self.workers)
for w in self.workers:
self.table[w.id].state = TaskState.SUBMITTED
await self.bus.send(
"coordinator", w.id, "task_assigned",
{"question": self.question, "source": w.source.name},
)
self._print_table("已派发全部子 Agent")
async def _settle_if_first(self, worker_id: str, answer: str) -> bool:
"""
竞态处理核心:用锁 + 幂等标志保证只结算一次、只广播一轮终止。
返回 True 表示本次是"首个有效命中"。
"""
async with self._settle_lock:
if self._settled:
# 迟到的命中:已结算过,直接忽略(幂等)
self.duplicate_hits.append(worker_id)
print(
f" [竞态] {worker_id} 也命中,但已由 {self.winner} 结算 —— "
f"忽略此次命中,不重复广播终止。"
)
return False
# 首个命中:结算并广播一轮终止
self._settled = True
self.winner = worker_id
self.answer = answer
print(f" [结算] 首个命中来自 {worker_id} —— 加锁结算,广播一轮 terminate。")
await self.bus.send(
"coordinator", BROADCAST, "terminate",
{"reason": "answer_found", "winner": worker_id},
)
return True
async def run(self, quiet_period: float = 2.5) -> dict:
"""
协调主循环:派发 -> 监听上报 -> 结算首个命中 -> 收集 ack -> 汇总。
"""
mode = "LLM 判断" if llm_available() else "关键词判断(离线可复现)"
print(f" [协调器] 判断模式:{mode}\n")
t0 = time.monotonic() # 记录并行执行的墙钟起点
await self._dispatch()
# 启动所有 Worker 协程
worker_tasks = [asyncio.create_task(w.run()) for w in self.workers]
done_states = {TaskState.SUCCEEDED, TaskState.FAILED, TaskState.TERMINATED}
last_terminate_time: Optional[float] = None
while True:
env = await self.sub.get_nowait_or_wait(timeout=0.5)
if env is None:
# 没有新消息:若已结算且过了静默期,认为收敛,退出
if self._settled and last_terminate_time is not None:
if _now() - last_terminate_time > quiet_period:
break
# 全部 Worker 进入终态也退出
if all(r.state in done_states for r in self.table.values()):
break
continue
rec = self.table.get(env.sender_id)
if env.type == "status_update" and rec:
prev = rec.state
rec.state = TaskState(env.payload["state"])
rec.note = env.payload.get("note", "")
rec.updated = _now()
# 仅在"状态机跳变"时刷新状态表,避免每一步抓取都刷屏
if rec.state != prev:
self._print_table(f"{env.sender_id} -> {rec.state.value}")
elif env.type == "result":
if env.payload.get("found"):
is_first = await self._settle_if_first(
env.sender_id, env.payload["answer"]
)
if is_first:
last_terminate_time = _now()
elif env.type == "ack":
self._acks.add(env.sender_id)
print(f" [ack] {env.sender_id} 已确认终止({len(self._acks)} 个已 ack)")
# 等待所有 Worker 协程收尾
await asyncio.gather(*worker_tasks, return_exceptions=True)
elapsed = time.monotonic() - t0 # 并行执行的墙钟耗时(含收敛静默期)
self._print_table("最终状态")
return {
"winner": self.winner,
"answer": self.answer,
"duplicate_hits": self.duplicate_hits,
"acks": sorted(self._acks),
"settled_once": self._settled,
"terminate_broadcasts": sum(
1 for e in self.bus.history if e.type == "terminate"
),
"parallel_seconds": elapsed,
}
# ————————————————————————— 串行基线(性能对比) —————————————————————————
async def run_sequential(sources: List[Source], question: str) -> dict:
"""
串行基线:逐个来源"抓取 + 判断",命中即止。
用于对照并行执行的墙钟收益(对应书中实验要求"记录并对比并行/串行时间差异")。
这里真实地一个接一个 ``await source.fetch()``,因此耗时是**实测**而非估算——
绝不伪造数据;命中判断复用与并行版本相同的 :func:`judge_answer`。
"""
t0 = time.monotonic()
winner: Optional[str] = None
answer: Optional[str] = None
fetched = 0
for src in sources:
fetched += 1
text = await src.fetch()
ans = await judge_answer(question, text)
if ans:
winner, answer = src.name, ans
break
return {
"seconds": time.monotonic() - t0,
"winner": winner,
"answer": answer,
"fetched": fetched,
"total": len(sources),
}
demo.py¶
"""
实验 10-6 演示入口:同时从多个网站搜集信息的 Agent
=================================================
一条命令即可运行:
python demo.py
演示内容(对应书中强调的机制):
(a) 消息总线的发布/订阅:日志里可见带信封的消息流(BUS 前缀);
(b) N 个子 Agent 并行执行,主协调器实时刷新任务状态表;
(c) 某子 Agent 命中后触发级联终止,其余 Agent 收到 terminate 并优雅退出(ack);
(d) 多个子 Agent 几乎同时命中时,只结算一次、只广播一轮终止(幂等 + 加锁);
(e) (--compare)并行 vs 串行的墙钟耗时实测对比,验证并行化的性能收益。
默认使用离线的关键词判断,保证结果可复现;
若配置了 OPENAI_API_KEY 且未设 USE_LLM=0,子 Agent 会改用真实 LLM 做判断。
命令行参数详见 `python demo.py --help`;不传任何参数即为原有默认行为
(10 个子 Agent、内置问题、离线关键词判断、详细 BUS 日志)。
"""
from __future__ import annotations
import argparse
import asyncio
import json
import os
import time
try:
from dotenv import load_dotenv
load_dotenv()
except Exception: # noqa: BLE001 —— 没装 python-dotenv 也能跑
pass
from agents import Coordinator, WorkerAgent, run_sequential
from llm import llm_available
from message_bus import MessageBus
from sources import DEMO_SOURCES, QUESTION, build_sources
def _parse_args() -> argparse.Namespace:
"""解析命令行参数;不传任何参数时行为与之前完全一致(10 Agent、离线、详细日志)。"""
parser = argparse.ArgumentParser(
prog="demo.py",
formatter_class=argparse.RawDescriptionHelpFormatter,
description=(
"实验 10-6:多个同构子 Agent 并行搜索 + 中心协调的演示。\n"
"展示消息总线发布/订阅、并行派发、实时状态监控、级联终止与竞态处理;\n"
"默认离线关键词判断(结果可复现),不传参数即为原有默认行为。"
),
epilog=(
"示例:\n"
" python demo.py # 默认:10 个 Agent、内置问题、离线可复现\n"
" python demo.py --agents 6 # 改为 6 个并行 Agent\n"
" python demo.py --compare # 额外实测并行 vs 串行的墙钟耗时\n"
" python demo.py --output result.json # 把结论写入 JSON 文件\n"
" python demo.py --use-llm --model gpt-5.6-luna # 用真实 LLM 判断(需配 key)"
),
)
parser.add_argument(
"-q",
"--query",
default=QUESTION,
metavar="问题",
help="研究问题(默认使用内置问题)。注意:离线关键词判断是针对内置来源调校的,"
"自定义问题通常应搭配 --use-llm 的真实 LLM 判断才有意义。",
)
parser.add_argument(
"-n",
"--agents",
type=int,
default=len(DEMO_SOURCES),
metavar="N",
help=f"并行子 Agent 的数量(默认 {len(DEMO_SOURCES)},对应书中约 10 个并行 Agent)。"
"N>=2 时始终包含两个含答案的源以稳定演示竞态与级联终止。",
)
parser.add_argument(
"--model",
default=None,
metavar="MODEL",
help="LLM 模型名(等价于设置环境变量 OPENAI_MODEL;仅在 --use-llm 且配置了 "
"OPENAI_API_KEY 或 OPENROUTER_API_KEY 时生效)。默认沿用环境变量或 gpt-5.6-luna。",
)
parser.add_argument(
"-o",
"--output",
default=None,
metavar="PATH",
help="把最终结论(含并行/串行耗时、winner、竞态统计)以 JSON 写入该路径。默认不写文件。",
)
parser.add_argument(
"--compare",
action="store_true",
help="并行运行结束后,再实测一遍串行基线并打印墙钟耗时对比(验证并行化收益)。默认不运行串行基线。",
)
parser.add_argument(
"--use-llm",
action="store_true",
help="强制启用真实 LLM 判断(等价于环境变量 USE_LLM=1;仍需配置 "
"OPENAI_API_KEY 才会真正生效,否则自动回退离线关键词判断)。默认不启用。",
)
parser.add_argument(
"--quiet",
action="store_true",
help="减少消息总线的逐条 BUS 日志打印(任务状态表/结论/自检不受影响)。默认打印全部日志。",
)
return parser.parse_args()
async def main(args: argparse.Namespace):
if args.agents < 1:
raise SystemExit("错误:--agents 至少为 1")
if args.use_llm:
# 仅设置意图开关;是否真正调用 LLM 仍取决于 llm.llm_available()
# (还需配置 OPENAI_API_KEY),未配置时会自动回退离线关键词判断。
os.environ["USE_LLM"] = "1"
if args.model:
os.environ["OPENAI_MODEL"] = args.model
sources = build_sources(args.agents)
print("=" * 78)
print("实验 10-6 · 同时从多个网站搜集信息的 Agent(并行搜索 + 中心协调)")
print("=" * 78)
print(f"任务问题:{args.query}")
print(f"并行来源数:{len(sources)} 个模拟'网站'(子 Agent 数 = 来源数)")
answer_srcs = [s.name for s in sources if s.holds_answer]
print(f"含答案的源:{answer_srcs}(其中前两个延迟相同,用于演示竞态)")
print("-" * 78)
bus = MessageBus(verbose=not args.quiet)
coordinator = Coordinator(bus, args.query)
# 并行装配 N 个同构子 Agent,每个绑定一个来源
for i, src in enumerate(sources):
w = WorkerAgent(f"worker-{i:02d}", src, bus, args.query)
coordinator.add_worker(w)
result = await coordinator.run()
print("=" * 78)
print("演示结论(自动校验)")
print("=" * 78)
total_msgs = len(bus.history)
print(f"1) 消息总线共传递 {total_msgs} 条带信封消息(发布/订阅正常工作)。")
print(f"2) {len(coordinator.workers)} 个子 Agent 并行执行,状态表全程实时刷新。")
print(f"3) 首个命中并结算的 Worker:{result['winner']}")
print(f" 答案:{result['answer']}")
print(f" 收到 terminate 并 ack 的 Worker:{result['acks']}")
print(f"4) terminate 广播轮数:{result['terminate_broadcasts']}(应为 1,证明只广播一轮)")
print(f" 迟到/并发的重复命中被忽略:{result['duplicate_hits'] or '无(本次无并发迟到命中)'}")
print(f" 是否只结算一次:{result['settled_once']}")
print(f"5) 并行执行墙钟耗时:{result['parallel_seconds']:.2f}s(含收敛静默期)")
# —— (e) 并行 vs 串行的墙钟对比:实测,绝不伪造 ——
seq = None
if args.compare:
print("-" * 78)
print("并行 vs 串行 墙钟对比(--compare,串行基线为实测)")
print("-" * 78)
seq = await run_sequential(sources, args.query)
print(
f" 串行:命中前逐个抓取了 {seq['fetched']}/{seq['total']} 个源,"
f"墙钟耗时 {seq['seconds']:.2f}s,winner={seq['winner']}"
)
print(f" 并行:墙钟耗时 {result['parallel_seconds']:.2f}s,winner={result['winner']}")
if seq["seconds"] > 0 and result["parallel_seconds"] > 0:
speedup = seq["seconds"] / result["parallel_seconds"]
saved = seq["seconds"] - result["parallel_seconds"]
print(f" 加速比 ≈ {speedup:.2f}×,节省约 {saved:.2f}s(并行让最快的源立即结束全局搜索)。")
# —— 结论落盘(可选)——
if args.output:
summary = {
"question": args.query,
"num_agents": len(coordinator.workers),
"sources": [s.name for s in sources],
"judge_mode": "llm" if llm_available() else "keyword_offline",
"total_bus_messages": total_msgs,
"winner": result["winner"],
"answer": result["answer"],
"duplicate_hits": result["duplicate_hits"],
"acks": result["acks"],
"settled_once": result["settled_once"],
"terminate_broadcasts": result["terminate_broadcasts"],
"parallel_seconds": round(result["parallel_seconds"], 3),
"sequential_baseline": (
{
"seconds": round(seq["seconds"], 3),
"fetched": seq["fetched"],
"winner": seq["winner"],
}
if seq
else None
),
"timestamp": time.strftime("%Y-%m-%d %H:%M:%S"),
}
with open(args.output, "w", encoding="utf-8") as f:
json.dump(summary, f, ensure_ascii=False, indent=2)
print(f"\n[已写入] 结论 JSON -> {args.output}")
# —— 断言式自检:仅在离线可复现模式下强断言(LLM 模式命中与否取决于模型)——
if not llm_available():
assert result["winner"] is not None, "应至少有一个 Worker 命中"
assert result["terminate_broadcasts"] == 1, "级联终止只能广播一轮"
assert result["settled_once"] is True, "必须完成且只结算一次"
print("\n[自检通过] 单次结算 + 单轮终止广播 + 级联 ack 均符合预期。")
else:
print("\n[提示] 当前为真实 LLM 判断模式,是否命中取决于模型,不做强断言。")
if __name__ == "__main__":
asyncio.run(main(_parse_args()))
llm.py¶
"""
可选的 LLM 判断层
=================
子 Agent 抓到网页文本后,需要判断"这段内容是否回答了目标问题"。
- 默认使用 ``sources.keyword_judge`` 的可控关键词判断(无需联网、结果可复现);
- 若设置了 ``OPENAI_API_KEY`` 且未强制离线,则调用真实 LLM 做判断,
展示"子 Agent 用大模型做真实决策"这条路径。
协调 / 总线 / 终止 / 竞态这些机制与 LLM 无关,始终是真实实现;
LLM 只影响"单个源里是否命中答案"的判断,属于可插拔部分。
"""
from __future__ import annotations
import json
import os
from typing import Optional
from sources import keyword_judge
def _to_openrouter_model(model: str) -> str:
"""把模型名映射到 OpenRouter 命名空间(用于无 OPENAI_API_KEY 的回退路径)。"""
if "/" in model:
return model # 已是 OpenRouter 命名空间,原样使用
if model.startswith("gpt-"):
return "openai/" + model # gpt-* -> openai/gpt-*
if model.startswith("claude-"):
return "anthropic/claude-opus-4.8"
return "openai/gpt-5.6-luna" # 兜底:当前便宜旗舰
def _loads_lenient(content: str):
"""容错解析 JSON:兼容个别模型把 JSON 包在 ```json ... ``` 代码围栏里的情况。"""
s = (content or "").strip()
if s.startswith("```"):
s = s.split("\n", 1)[-1] if "\n" in s else s
s = s.rsplit("```", 1)[0].strip()
if s.lower().startswith("json"):
s = s[4:].strip()
return json.loads(s)
def llm_available() -> bool:
"""是否具备调用真实 LLM 的条件。
注意:本实验的重点是"并行协调 / 消息总线 / 级联终止 / 竞态结算"这些机制,
而不是 LLM 的检索质量。为保证演示**可复现**(只有真正包含答案的源才命中,
从而稳定触发竞态与级联终止),默认走确定性的关键词判断。
只有显式设置 USE_LLM=1 时才启用真实 LLM 判断(可能对 mock 源产生幻觉,仅供体验)。
Key 解析:优先 OPENAI_API_KEY;没有则回退 OPENROUTER_API_KEY(走 OpenRouter)。
"""
if os.getenv("USE_LLM", "").lower() not in ("1", "true", "yes"):
return False
return bool(os.getenv("OPENAI_API_KEY") or os.getenv("OPENROUTER_API_KEY"))
async def judge_answer(question: str, text: str) -> Optional[str]:
"""
判断 text 是否回答了 question。命中则返回答案字符串,否则返回 None。
优先用 LLM;不可用或出错时回退到关键词判断,保证 demo 始终可跑通。
"""
if not llm_available():
return keyword_judge(text)
try:
# 延迟导入,避免没装 openai 时影响关键词路径
from openai import AsyncOpenAI
model = os.getenv("OPENAI_MODEL", "gpt-5.6-luna")
# 通用回退:优先直连 OPENAI_API_KEY;否则用 OPENROUTER_API_KEY 走 OpenRouter。
if os.getenv("OPENAI_API_KEY"):
client = AsyncOpenAI(
api_key=os.getenv("OPENAI_API_KEY"),
base_url=os.getenv("OPENAI_BASE_URL") or None,
)
else:
client = AsyncOpenAI(
api_key=os.getenv("OPENROUTER_API_KEY"),
base_url="https://openrouter.ai/api/v1",
)
model = _to_openrouter_model(model)
prompt = (
f"问题:{question}\n\n"
f"网页内容:{text}\n\n"
"严格只依据上面的『网页内容』判断,**不得使用你自己的知识**。"
"只有当网页内容里**确实出现了**问题的具体答案时,才算命中。"
"命中则只输出 JSON:{\"found\": true, \"answer\": \"<从网页内容中摘出的答案>\"};"
"若网页内容没有直接给出答案(哪怕你知道答案),一律输出 {\"found\": false}。"
"只输出 JSON,不要其它文字。"
)
resp = await client.chat.completions.create(
model=model,
messages=[{"role": "user", "content": prompt}],
temperature=0,
response_format={"type": "json_object"},
)
data = _loads_lenient(resp.choices[0].message.content)
if data.get("found"):
return data.get("answer") or text
return None
except Exception as exc: # noqa: BLE001 —— 任何异常都回退到离线判断
print(f" [llm] 调用失败,回退关键词判断:{exc}")
return keyword_judge(text)
message_bus.py¶
"""
进程内异步消息总线(Message Bus)
================================
模仿 Redis Pub/Sub 的语义,但完全跑在单进程的 asyncio 事件循环里,
无需真正部署 Redis。它承担实验 10-6 里"中心协调"的通信底座:
- 每条消息都封装在 ``Envelope`` 信封里,带上 sender_id / target / type / payload;
- Agent 通过 ``subscribe()`` 拿到一个订阅句柄,按消息类型接收;
- Agent 通过 ``publish()`` 把消息投递给指定目标或广播给所有人;
- 总线本身不做任何业务判断,只负责"可靠地把信封送达订阅者"。
设计要点:
- 使用 ``asyncio.Queue`` 做每个订阅者的收件箱,天然线程/协程安全;
- ``target`` 为 ``BROADCAST`` 时投递给所有订阅了该类型的人(发送者除外);
- 打印带时间戳的事件日志,方便在演示里"看见"发布/订阅的消息流。
"""
from __future__ import annotations
import asyncio
import itertools
import json
import time
from dataclasses import dataclass, field
from typing import Any, Dict, List, Optional
# 广播目标常量:发给所有订阅者
BROADCAST = "*"
# 全局单调递增的消息序号,方便在日志里追踪顺序
_seq_counter = itertools.count(1)
# 演示启动时间,用于打印相对时间戳(更易读)
_START_TIME = time.monotonic()
def _now() -> float:
"""返回自演示启动以来的秒数(相对时间戳)。"""
return time.monotonic() - _START_TIME
@dataclass
class Envelope:
"""消息信封:总线里流动的最小单元。"""
sender_id: str # 发送者 ID
target: str # 目标 Agent ID,或 BROADCAST 表示广播
type: str # 消息类型:task_assigned / status_update / result / terminate / ack ...
payload: Dict[str, Any] = field(default_factory=dict) # JSON 负载
seq: int = field(default_factory=lambda: next(_seq_counter)) # 全局序号
ts: float = field(default_factory=_now) # 相对时间戳
def short(self) -> str:
"""给日志用的紧凑单行表示。"""
tgt = "ALL" if self.target == BROADCAST else self.target
body = json.dumps(self.payload, ensure_ascii=False)
if len(body) > 80:
body = body[:77] + "..."
return (
f"[t={self.ts:6.2f}s #{self.seq:<3}] "
f"{self.sender_id:>11} -> {tgt:<11} | {self.type:<14} | {body}"
)
class Subscription:
"""订阅句柄:内部就是一个收件箱队列 + 关心的消息类型集合。"""
def __init__(self, owner_id: str, types: Optional[List[str]]):
self.owner_id = owner_id
# types 为 None 表示订阅所有类型
self.types = set(types) if types else None
self.inbox: "asyncio.Queue[Envelope]" = asyncio.Queue()
def accepts(self, env: Envelope) -> bool:
return self.types is None or env.type in self.types
async def get(self) -> Envelope:
return await self.inbox.get()
async def get_nowait_or_wait(self, timeout: float) -> Optional[Envelope]:
"""带超时地取一条消息;超时返回 None(便于子 Agent 在循环里轮询终止信号)。"""
try:
return await asyncio.wait_for(self.inbox.get(), timeout=timeout)
except asyncio.TimeoutError:
return None
class MessageBus:
"""异步消息总线:注册订阅者、投递信封、打印消息流日志。"""
def __init__(self, verbose: bool = True):
# owner_id -> 该 owner 的订阅列表
self._subs: Dict[str, List[Subscription]] = {}
self.verbose = verbose
# 记录全部流过总线的信封,便于事后统计/断言
self.history: List[Envelope] = []
def subscribe(self, owner_id: str, types: Optional[List[str]] = None) -> Subscription:
"""注册一个订阅者,返回订阅句柄。types=None 表示接收所有类型。"""
sub = Subscription(owner_id, types)
self._subs.setdefault(owner_id, []).append(sub)
return sub
async def publish(self, env: Envelope) -> None:
"""把信封投递到总线:广播或点对点。"""
self.history.append(env)
if self.verbose:
print(" BUS " + env.short())
delivered = 0
for owner_id, sub_list in self._subs.items():
# 点对点:只投递给指定目标
if env.target != BROADCAST and owner_id != env.target:
continue
# 广播时不回投给发送者自己
if env.target == BROADCAST and owner_id == env.sender_id:
continue
for sub in sub_list:
if sub.accepts(env):
await sub.inbox.put(env)
delivered += 1
# 让出事件循环,保证消息尽快被对端取走(更接近真实推送时序)
await asyncio.sleep(0)
# —— 便捷构造并发布 ——
async def send(
self,
sender_id: str,
target: str,
type: str,
payload: Optional[Dict[str, Any]] = None,
) -> Envelope:
env = Envelope(sender_id=sender_id, target=target, type=type, payload=payload or {})
await self.publish(env)
return env
sources.py¶
"""
模拟的"多个网站/信息源"
=======================
实验 10-6 强调的是"多个同构 Agent 并行搜索 + 中心协调",浏览器本身不是重点,
因此这里用一批**可控的模拟数据源**代替真实浏览器:
- 每个 Source 对应一个"网站",有名字、内容、以及一个可控的访问延迟;
- ``fetch()`` 模拟"打开网站并抓取",用 ``asyncio.sleep`` 制造不同的耗时,
让"谁先命中"变得可复现(也让级联终止/竞态现象稳定出现);
- 是否包含目标答案由 ``holds_answer`` 决定。
延迟是刻意设计的:让两个源几乎同时命中,以此稳定地触发"竞态"演示。
"""
from __future__ import annotations
import asyncio
from dataclasses import dataclass
from typing import List, Optional
@dataclass
class Source:
name: str # "网站"名字
latency: float # 单步抓取延迟(秒),模拟网络/渲染耗时
content: str # 该网站的可搜索文本
holds_answer: bool # 该网站是否真的包含目标答案
async def fetch(self) -> str:
"""模拟一次抓取:等待 latency 后返回内容。"""
await asyncio.sleep(self.latency)
return self.content
# —— 演示任务:并行查找"世界上最高的山峰是哪座、海拔多少" ——
# 只有部分源包含准确答案;两个正确源的延迟被设计得很接近,用来制造竞态。
QUESTION = "世界上海拔最高的山峰是哪座?海拔约多少米?"
ANSWER_KEYWORDS = ["珠穆朗玛", "珠峰", "Everest", "8848"]
DEMO_SOURCES: List[Source] = [
Source("baike-wiki", 0.9, "百科条目:喜马拉雅山脉横亘于青藏高原南缘。", False),
Source("news-portal", 1.1, "新闻门户:近期多支登山队计划攀登高海拔山峰。", False),
Source("geo-journal", 0.6, "地理期刊:世界最高峰为珠穆朗玛峰,海拔约 8848 米,位于中尼边境。", True),
Source("travel-blog", 1.4, "旅行博客:作者分享了在尼泊尔徒步大本营的见闻。", False),
# 与 geo-journal 完全相同的延迟:两个正确源会几乎同一时刻命中,稳定触发竞态。
Source("forum-qa", 0.6, "问答社区:网友讨论最高峰其实是珠穆朗玛峰,官方海拔 8848.86 米。", True),
Source("edu-site", 1.2, "教育网站:介绍板块运动如何抬升出高大的山脉。", False),
Source("gov-data", 1.6, "政府数据:发布了若干山峰的测绘参数与坐标信息。", False),
Source("random-blog", 0.8, "个人博客:随笔记录了一次雪山摄影旅行。", False),
Source("science-mag", 1.3, "科普杂志:解释高海拔缺氧对人体的影响。", False),
Source("map-service", 1.0, "地图服务:可查询各山峰的等高线与地形剖面。", False),
]
def keyword_judge(text: str) -> Optional[str]:
"""
可控的"判断是否命中答案"逻辑(不依赖 LLM):
命中任一关键词即认为找到了答案,返回抽取到的答案句子;否则返回 None。
"""
for kw in ANSWER_KEYWORDS:
if kw in text:
# 简单地把包含关键词的那句话当作答案返回
for sentence in text.replace("。", "。\n").splitlines():
if kw in sentence:
return sentence.strip(":: ")
return text
return None
def build_sources(n: int) -> List[Source]:
"""
构造 n 个并行来源,供命令行 ``--agents N`` 动态调整并行 Agent 数量。
设计约束(保证演示始终可复现):
- n == 10 时**原样返回** ``DEMO_SOURCES``,默认行为与之前完全一致;
- n >= 2 时始终包含两个"含答案且延迟相同"的源(geo-journal / forum-qa),
从而稳定触发命中、竞态与级联终止;
- 其余用无答案的填充源补齐;n 超过内置源数量时循环生成并加序号后缀。
"""
if n < 1:
raise ValueError("并行 Agent 数量至少为 1")
if n == len(DEMO_SOURCES):
return list(DEMO_SOURCES)
answer_pool = [s for s in DEMO_SOURCES if s.holds_answer]
filler_pool = [s for s in DEMO_SOURCES if not s.holds_answer]
k_answer = min(len(answer_pool), 2 if n >= 2 else 1)
n_filler = n - k_answer
fillers: List[Source] = []
seen: dict = {}
for i in range(n_filler):
base = filler_pool[i % len(filler_pool)]
seen[base.name] = seen.get(base.name, 0) + 1
name = base.name if seen[base.name] == 1 else f"{base.name}-{seen[base.name]}"
fillers.append(Source(name, base.latency, base.content, False))
sources = fillers
# 把含答案的源插入到确定的位置(沿用默认布局:靠前但不在最前,便于观察并行推进)
for j, src in enumerate(answer_pool[:k_answer]):
pos = min(2 + 2 * j, len(sources))
sources.insert(pos, Source(src.name, src.latency, src.content, src.holds_answer))
return sources