跳转至

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│  (同构·并行)
                    └─────┘  └─────┘  └─────┘  └─────┘         └─────┘

对应书中强调的五个机制:

  1. 消息总线(Message Bus):所有通信都是发布到总线的带信封消息,按 type 订阅。日志中每条 BUS ... 行就是一次发布/投递(Redis Pub/Sub 语义的进程内实现)。
  2. 并行派发Coordinator 同时给 10 个子 Agent 发 task_assignedasyncio.create_task 并发执行。
  3. 实时监控(push 范式):子 Agent 执行中主动 status_update 上报进度, 主 Agent 维护任务状态表并在状态机跳变时实时刷新打印。 状态机:已提交 → 执行中 →(需要输入)→ 已完成 / 失败 / 已终止
  4. 级联终止:某子 Agent 命中后,主 Agent 广播 terminate;其余子 Agent 在 循环的安全检查点发现信号后回 ack 并优雅退出(状态置为"已终止")。
  5. 竞态处理:多个子 Agent 可能几乎同时命中,主 Agent 用 asyncio.Lock + 幂等标志 _settled 保证只结算一次、只广播一轮终止;迟到的命中被记录并忽略。

为让"竞态""级联终止"可复现,各来源被赋予不同的模拟延迟,其中 geo-journalforum-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_URLOPENAI_MODEL)。

通用回退:若未设置 OPENAI_API_KEY 但设了 OPENROUTER_API_KEY,则真实 LLM 判断 自动改走 OpenRouter,并把模型名映射到其命名空间(gpt-5.6-lunaopenai/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_KEYOPENROUTER_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只结算一次 == Truewinner 非空。跑通即证明机制正确:

4) terminate 广播轮数:1(应为 1,证明只广播一轮)
   迟到/并发的重复命中被忽略:['worker-04']
[自检通过] 单次结算 + 单轮终止广播 + 级联 ack 均符合预期。

局限与注意事项

本实验把重点放在协调机制(消息总线/并行派发/级联终止/竞态处理)上,这些均为真实 实现;但为了可离线运行与自动验证,以下三处做了简化,是已知局限:

  • 局限·模拟源非真实浏览器:不启动真实浏览器,来源是可控的模拟数据 + 延迟。若要接 真实 Computer Use,只需把 WorkerAgent.run() 里的"抓取一步 + 判断"换成真实浏览器 操作,协调层无需改动。
  • 局限·竞态靠相同延迟稳定复现geo-journalforum-qa 两个正确源被人为设成 相同延迟,才能稳定触发"同时命中";真实环境里竞态是偶发的,但加锁 + 幂等的结算逻辑 对偶发竞态同样成立,不依赖这个人为设定。
  • 局限·进程内总线非真实 RedisMessageBus进程内 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