跳转至

async-agent

第4章 · 工具 · 配套项目 chapter4/async-agent

项目说明

实验 4-5:带并行执行和打断能力的异步 Agent(★★★)

本目录是《深入理解 AI Agent》实验 4-5 的配套可运行代码,实现了设计文档 agent_framework_design.md 中描述的事件驱动异步 Agent 框架(Flux)的核心部分。

在 4-4 的简单事件队列之上,本实验进入异步 Agent 的深水区,聚焦四件事: 异步工具执行、事件队列与批量处理、打断机制、并行工具的取消与状态查询。 Agent 需要同时管理多个并发任务,处理打断与恢复,并根据实时状态动态决策。

本目录提供两条使用路径:

  • 离线演示(推荐先跑,零依赖、无需 API key):把三项核心异步能力单独拎出来、 用可测量的方式演示——并行 vs 串行的墙钟时间对比、打断/取消后恢复、状态检查点持久化与恢复。 这条路径不联网、不调用 LLM,甚至不需要安装 openaipython demo.py 即可直接运行。
  • LLM 场景(还原书中四个验证场景):Agent 的决策由真实 LLM(默认 OpenAI gpt-5.6-luna, function calling)完成,需要配置 API key。

两条路径共用同一套异步运行时;长任务都用模拟的异步"终端命令"(带进度输出)实现,绝不真跑危险命令。


一、架构

对应设计文档第 5 节的事件处理循环,全部基于 asyncio 单线程实现:

                       ┌──────────────┐
   用户消息 / 打断  ──▶ │    inbox     │  所有进来的原始事件
   异步任务完成通知 ──▶ │  (asyncio.Q) │
                       └──────┬───────┘
                   ┌──────────▼───────────┐   判定紧急度 classify_urgency()
                   │     _dispatcher      │──▶ 打断 / 立即处理 / 排队
                   └──────────┬───────────┘
             ┌────────────────┼───────────────────┐
   INTERRUPT │        IMMEDIATE│           DEFERRED│
   取消当前turn+异步工具    直接入 work        进 pending 缓冲,
   并留痕                                     异步结果到达时批量追加
                   ┌──────────▼───────────┐
                   │        work          │  待处理的事件批次
                   └──────────┬───────────┘
                   ┌──────────▼───────────┐
                   │       _worker        │  逐批:追加到轨迹 -> run_llm_turn()
                   │   turn_task 可被取消  │  (打断时 cancel 掉这个子任务)
                   └──────────────────────┘

  TaskManager:管理模拟异步终端任务(start / query / cancel / cancel_all)
               任务自然完成 -> 以"新事件"(async.result) 注入 inbox

代码文件:

文件 作用
events.py 事件模型 Event(含检查点序列化 to_dict/from_dict)、事件类型、紧急度判定 classify_urgency()
tasks.py 模拟异步"终端命令"与 TaskManager(进度推进、按 ID 取消/查询、状态 snapshot/restore
runtime.py AgentRuntime:事件循环、两种处理机制、LLM function calling、工具执行、检查点 save_checkpoint/load_checkpoint
async_demos.py 三个离线演示(无需 API key):并行墙钟对比、打断/恢复、状态检查点
demo.py 统一命令行入口(argparse 子命令):离线演示 + 四个 LLM 验证场景

两种事件处理机制(设计文档 5.1)

  • 取消式处理(Cancellation-Based):紧急事件(用户"取消/停止")到达时, 立即取消正在进行的 LLM turn,并取消所有后台异步工具,把打断事件与取消回执写入轨迹。
  • 排队处理(Queued):非紧急事件(补充性指令)先进入 pending 缓冲,不打断正在进行的工作; 当某个异步工具完成、产生 async.result 事件时,一次性把 pending 里的事件批量追加到轨迹,再触发一次 LLM。

紧急度判定规则(简单可解释):

  1. 含打断关键词(取消/停止/stop…)→ INTERRUPT(取消式处理)
  2. 是一个提问(带问号或疑问词,如"现在几点了?")→ IMMEDIATE(立即回应,但打断后台任务)
  3. 其它补充性指令(如"用日语回复")→ DEFERRED(排队,批量处理)

异步工具

run_terminal_command异步工具:调用后立刻返回 task_id 占位符(不阻塞), 命令在后台按固定速度推进进度;真正完成后,其结果作为一条新事件async.result)注入对话。 另有 query_task / cancel_task 按 ID 查询进度与取消,get_current_time 用于即时提问。

时间轴加速:为便于复现,1 个"模拟秒"默认映射为 0.4 真实秒(FLUX_TICK_REAL 可调)。 速度差 3% / 2% / 1% 每(模拟)秒是否过 50% 的判定逻辑完全保留。


二、运行

命令行入口是 demo.py,用 argparse 子命令组织,python demo.py --help 查看全部用法。

离线演示(无需 API key,开箱即用)

cd chapter4/async-agent

python demo.py              # 默认:依次运行下面三个离线演示
python demo.py offline      # 同上:显式地依次运行三个离线演示
python demo.py parallel     # 能力一:并行 vs 串行工具调用的墙钟时间对比(打印加速比)
python demo.py interrupt    # 能力二:长任务运行中被打断/取消,随后系统恢复
python demo.py state        # 能力三:状态检查点持久化 + 跨会话恢复并校验

这三个演示不联网、不调用 LLM,连 openai 都无需安装——用纯 asyncio 直接测量并行加速、 打断后的状态冻结、以及检查点的落盘与还原。

LLM 验证场景(还原书中四个场景,需要 API key)

pip install -r requirements.txt
cp env.example .env         # 填入 OPENAI_API_KEY

python demo.py scenarios                # 依次运行全部四个场景
python demo.py scenarios --scenario 1   # 只跑场景 1(异步执行 + 即时提问)
python demo.py scenarios --scenario 3   # 只跑场景 3(打断机制)

默认用 OpenAI gpt-5.6-luna。也可切换服务商(OpenAI 兼容接口):

# Moonshot(默认模型为当前的推理模型 kimi-k3)
LLM_PROVIDER=moonshot python demo.py scenarios --scenario 1
# 火山方舟 ARK(LLM_MODEL 填推理接入点 ID)
LLM_PROVIDER=ark LLM_MODEL=ep-xxxx python demo.py scenarios --scenario 1

OpenRouter 通用兜底:未配置 OPENAI_API_KEY(且未用 moonshot/ark provider)时, 只要设置了 OPENROUTER_API_KEYdemo.py 会自动改走 OpenRouter,并把模型名映射为 provider/model 形式(gpt-*openai/…claude-*anthropic/claude-opus-4.8、 含 / 的原样透传)。也可显式 LLM_PROVIDER=openrouter。例如: OPENROUTER_API_KEY=sk-or-xxx LLM_MODEL=openai/gpt-5.6-luna python demo.py scenarios --scenario 1

Moonshot 默认走推理模型 kimi-k3(旧的 kimi-k2-*-previewmoonshot-v1-* 已过时/停用)。 推理模型要求 temperature=1max_tokens>=2048demo.py 会按模型自动套用这套采样参数,无需手动配置。

兼容旧用法:python demo.py --scenario N 会自动等价为 scenarios --scenario N

日志中不同来源用颜色区分:USER(用户)、AGENT(Agent 回复)、TOOL(工具调用)、 TASK(后台异步任务)、TRAJ(轨迹留痕)、STATE(状态检查点)、SYSTEM(框架事件)。


三、离线演示的三项能力(真实测量输出)

以下三段均为真实运行输出节选(无需 API key),演示异步到底带来了什么。

能力一:并行 vs 串行工具调用(python demo.py parallel

四个相互独立的只读感知工具(读文件 / 搜索 / 查库 / 向量检索),串行逐个 await 与并行 asyncio.gather 的墙钟时间对比:

  ── 结果对比 ─────────────────────────────────────────────
  串行总耗时(Σ 各工具)                  4.51s
  并行总耗时(gather)                 1.50s
  并行理论下界(最慢单个)                  1.50s
  加速比 = 串行 / 并行                 3.00x
  ─────────────────────────────────────────────────────────

墙钟时间由「各工具求和」降到「取最大单个」——这正是书中「只读感知工具天然适合并行」的量化落点。

能力二:打断 / 取消 / 恢复(python demo.py interrupt

三个并行后台任务运行中,用户先即时提问(不阻塞任务),随后发出「取消」打断:

[  1.00s] USER   | (即时提问)现在几点了?
[  1.00s] AGENT  | 现在 00:14:26。三个后台任务仍在并行推进,未被这次提问阻塞。
[  2.00s] USER   | (打断)取消
[  2.00s] TASK   | T1 已被取消 🛑(进度停在 39%)
[  2.00s] TASK   | T2 已被取消 🛑(进度停在 26%)
[  2.00s] TASK   | T3 已被取消 🛑(进度停在 13%)

  ── 打断后各任务状态(进度冻结在中途)───────────────────
  task_id 命令                        状态              进度
  T1      python analyze_fast.py    cancelled      39%
  T2      python analyze_mid.py     cancelled      26%
  T3      python analyze_slow.py    cancelled      13%
  ─────────────────────────────────────────────────────────
[  2.05s] SYSTEM | 打断处理完毕,系统恢复空闲,可继续接受新任务……
[  5.52s] TASK   | T4 完成 ✅
[  5.52s] AGENT  | 已从打断中恢复,新任务 T4 正常完成:……

打断只冻结被取消任务的进度,运行时本身无损,随后能立即接受并跑完新任务。

能力三:状态检查点持久化与恢复(python demo.py state

会话 A 产生一段轨迹 + 两个运行中的后台任务,落盘为 checkpoints/agent_state.json; 会话 B 用全新运行时从磁盘恢复并校验:

  ── 恢复校验 ─────────────────────────────────────────────
  轨迹事件数     保存前 3  ->  恢复后 3  [一致 ✓]
  可重建 LLM 上下文消息 4 条(system + 轨迹回放)
  task_id 命令                             保存前进度  恢复后状态           进度
  T1      python analyze_fast.py           21%  suspended      21%
  T2      python analyze_slow.py            7%  suspended       7%
  ─────────────────────────────────────────────────────────

轨迹与任务进度完整落盘并跨会话还原;运行中的任务恢复后标记为 suspended,保留最后已知进度, 供上层决定「重跑」还是「按进度续跑」——这就是异步任务的状态管理。


四、四个 LLM 验证场景

场景 1:异步工具执行

Agent 执行一个长终端命令,期间用户插入提问"现在几点了?"。 因为长命令是异步的、不阻塞,Agent 立即用 get_current_time 回应时间, 等后台任务完成后再把分析结论呈现出来。

场景 2:事件队列与批量处理

Agent 执行长任务期间,用户连续发"记得用日语回复""整理成网页"。 这两条是非紧急指令,先进入排队缓冲;任务完成时,框架把它们一次性批量追加到轨迹, Agent 再综合所有指令,输出日语的 HTML 结果。

场景 3:打断机制

Agent 执行长任务,用户发"取消"。框架立即取消当前执行流并取消后台异步工具, 在轨迹中记录打断事件(user.interrupt)和取消回执(system.note,含被取消的 task_id)。

场景 4:并行工具的取消与状态查询

用户要求"同时运行这三个脚本,哪个先完成就查其余进度,未过 50% 就取消"。 三个脚本速度分别为 3% / 2% / 1% 每秒。Agent 同时启动三个异步任务; 最快的先完成后,Agent 查询另外两个(约 66% 与 33%),取消未过 50% 的那个, 其余完成后整合出报告。


五、LLM 场景真实运行输出(关键片段)

以下均为真实调用 gpt-5.6-luna(OpenAI 兼容接口)的输出节选(时间戳为真实秒,需配置 API key 复现)。

场景 1(异步执行 + 即时提问)

[ 3.97s] AGENT | 任务已在后台启动(task_id:T1)。完成后我会根据日志分析结果给出结论。
[ 4.96s] TASK  | T1 `python analyze_logs.py` 进度 22%      ← 任务仍在后台跑
[ 5.19s] TOOL  | get_current_time -> 2026-07-18 13:43:30  ← 即时提问先回应
[ 6.91s] AGENT | 现在是 2026 年 7 月 18 日 13:43:30。
[12.19s] TRAJ  | + async.result  异步完成 T1               ← 真实结果作为新事件注入
[16.67s] AGENT | 日志分析已完成,结论如下:共扫描 12,840 条记录…  ← 再呈现分析

场景 2(批量处理)

[ 1.50s] SYSTEM | 事件进入排队缓冲(当前积压 1 条)
[ 1.90s] SYSTEM | 事件进入排队缓冲(当前积压 2 条)
[12.05s] TASK   | T1 完成 ✅
[12.05s] SYSTEM | 异步结果到达,批量处理 2 条积压的非紧急事件
[12.06s] TRAJ   | + async.result   异步完成 T1
[12.06s] TRAJ   | + user.input     记得最后用日语回复
[12.06s] TRAJ   | + user.input     把结果整理成一个网页(HTML)
...
[22.38s] AGENT  | <!DOCTYPE html>…<h2>分析結論</h2>… (批量指令一次性满足:日语 + HTML)

场景 3(打断)

[ 2.40s] TASK   | 启动异步任务 T1: `python analyze_logs.py` (速度 4%/模拟秒)
[ 4.00s] USER   | (interrupt) 取消
[ 4.00s] TASK   | T1 已被取消 🛑(进度停在 14%)
[ 4.00s] TRAJ   | + user.interrupt  用户打断:取消
[ 4.00s] TRAJ   | + system.note     打断回执,取消任务 ['T1']
[ 5.04s] AGENT  | 已停止后台任务 T1。

场景 4(并行 + 状态查询 + 按 50% 阈值取消 + 整合报告)

[ 2.82s] TASK | 启动异步任务 T1: `python analyze_fast.py` (速度 3%/模拟秒)
[ 2.82s] TASK | 启动异步任务 T2: `python analyze_mid.py`  (速度 2%/模拟秒)
[ 2.82s] TASK | 启动异步任务 T3: `python analyze_slow.py` (速度 1%/模拟秒)
[16.47s] TASK | T1 完成 ✅                               ← 最快脚本先完成
[19.84s] TOOL | query_task(T2) -> running 84%           ← 查询其余两个进度
[19.84s] TOOL | query_task(T3) -> running 42%
[21.93s] TOOL | cancel_task(T3) -> 已取消 (进度 47%)     ← 未过 50%,取消
[22.89s] TASK | T2 完成 ✅
[26.50s] AGENT | ## 分析汇总报告 … analyze_slow.py:已取消(未超 50%)…


六、注意事项

  • 离线演示(parallel/interrupt/state)无需任何 API key、也无需安装 openai,开箱即跑。
  • 只有 scenarios 子命令需要联网并配置有效的 API keyOPENAI_API_KEY,或切换到 MOONSHOT_API_KEY / ARK_API_KEY)。
  • LLM 决策由真实模型产生,输出措辞每次可能略有不同;四个场景的行为逻辑是稳定可复现的。 若遇到 OpenAI 偶发的高延迟,重跑即可。
  • 时间轴已加速;把 FLUX_TICK_REAL 调大可让演示更接近书中"几十秒"的真实节奏, 调小则更快(过小可能让场景 4 的"未过 50% 就取消"来不及判定)。
  • 所有"终端命令"均为模拟,不会在你的机器上真实执行任何命令。

源代码

async_demos.py

"""离线演示:不依赖任何 LLM / API key,直接驱动异步运行时的底层原语。

`demo.py` 里的四个「场景」需要真实 LLM 做决策;本模块则把实验 4-5 的三项核心
异步能力单独拎出来,用可测量、可复现的方式演示,**无需联网、无需 API key**:

  - demo_parallel  :并行 vs 串行工具调用的【墙钟时间】对比(真实测量,打印加速比)。
  - demo_interrupt :长任务运行中被【打断/取消】,随后系统【恢复】并接受新任务。
  - demo_state     :Agent 状态【检查点持久化】到磁盘,再【跨会话恢复】并校验。

这三个演示共同回答「异步到底带来了什么」——用数字和状态变化说话,而不只是措辞。
"""

from __future__ import annotations

import asyncio
import datetime
import os
import time

from runtime import AgentRuntime, format_log
from events import Event, EventType
from tasks import TaskManager
import tasks


class Logger:
    """与 runtime 同款的彩色时间戳日志器(相对本次演示起点计时)。"""

    def __init__(self) -> None:
        self.t0 = time.time()

    def __call__(self, source: str, text: str) -> None:
        print(format_log(self.t0, source, text), flush=True)


def banner(title: str) -> None:
    print("\n" + "=" * 78)
    print(f"  {title}")
    print("=" * 78, flush=True)


# ============================ 1. 并行 vs 串行 ============================

# 一组相互独立的【只读感知工具】(读文件 / 搜索 / 查库 / 向量检索)。
# 只读、无副作用,因此可以安全地并行——这正是书中「感知工具天然适合并行」的落点。
_PERCEIVE_TOOLS = [
    ("read_config.json", 0.8),
    ("web_search(‘异步 Agent’)", 1.2),
    ("db_query(orders)", 1.5),
    ("vector_lookup(memory)", 1.0),
]


async def _perceive(name: str, latency: float, log: Logger) -> tuple[str, float, float]:
    """模拟一次带 I/O 延迟的只读感知调用;返回 (名称, 标称延迟, 实测耗时)。"""
    t0 = time.time()
    log("TOOL", f"→ {name} 启动(模拟 I/O 耗时 {latency:.1f}s)")
    await asyncio.sleep(latency)
    dt = time.time() - t0
    log("TOOL", f"✓ {name} 完成(实测 {dt:.2f}s)")
    return name, latency, dt


async def demo_parallel() -> None:
    banner("能力一|并行工具调用:并行 vs 串行的墙钟时间对比")
    log = Logger()
    log("SYSTEM", "有 4 个相互独立的只读感知工具需要调用(无副作用,可安全并行)。")

    # —— 串行:一个 await 完再 await 下一个 ——
    log("SYSTEM", "\033[0m[串行] 逐个 await(同步 ReAct 的默认做法)……")
    seq_start = time.time()
    for name, lat in _PERCEIVE_TOOLS:
        await _perceive(name, lat, log)
    seq_total = time.time() - seq_start

    # —— 并行:一次性发起,asyncio.gather 并发等待 ——
    log("SYSTEM", "\033[0m[并行] 一次性发起,asyncio.gather 并发等待……")
    par_start = time.time()
    await asyncio.gather(*[_perceive(name, lat, log) for name, lat in _PERCEIVE_TOOLS])
    par_total = time.time() - par_start

    slowest = max(lat for _, lat in _PERCEIVE_TOOLS)
    speedup = seq_total / par_total if par_total else float("inf")

    print("\n  ── 结果对比 ─────────────────────────────────────────────")
    print(f"  {'工具':<26}{'标称延迟':>10}")
    for name, lat in _PERCEIVE_TOOLS:
        print(f"  {name:<26}{lat:>8.1f}s")
    print("  ─────────────────────────────────────────────────────────")
    print(f"  {'串行总耗时(Σ 各工具)':<26}{seq_total:>8.2f}s")
    print(f"  {'并行总耗时(gather)':<26}{par_total:>8.2f}s")
    print(f"  {'并行理论下界(最慢单个)':<26}{slowest:>8.2f}s")
    print(f"  {'加速比 = 串行 / 并行':<26}{speedup:>8.2f}x")
    print("  ─────────────────────────────────────────────────────────")
    print("  结论:独立的只读调用并行化后,墙钟时间由「求和」降到「取最大」。\n")


# ============================ 2. 打断 / 取消 / 恢复 ============================

async def demo_interrupt() -> None:
    banner("能力二|打断与取消:长任务运行中被打断,随后系统恢复")
    tasks.TICK_REAL = 0.15  # 本演示放慢节奏,留出「跑到一半再打断」的时间窗口
    log = Logger()
    completed: list = []

    async def on_complete(state) -> None:
        completed.append(state)

    tm = TaskManager(on_complete=on_complete, log=log)

    # 1) 并行启动三个后台异步任务
    log("SYSTEM", "启动三个并行后台分析任务(fast/mid/slow)……")
    for cmd in ["python analyze_fast.py", "python analyze_mid.py", "python analyze_slow.py"]:
        tm.start(cmd)

    # 2) 运行期间用户即时提问 —— 后台任务不被阻塞
    await asyncio.sleep(1.0)
    now = datetime.datetime.now().strftime("%H:%M:%S")
    log("USER", "(即时提问)现在几点了?")
    log("AGENT", f"现在 {now}。三个后台任务仍在并行推进,未被这次提问阻塞。")

    # 3) 跑到中途,用户发出打断 —— 立即取消所有在跑的任务
    await asyncio.sleep(1.0)
    log("USER", "(打断)取消")
    cancelled = tm.cancel_all()
    await asyncio.sleep(0.05)  # 让 CancelledError 在各协程内落地
    log("SYSTEM", f"已执行打断:取消了 {cancelled}(进度在被取消处冻结)")

    print("\n  ── 打断后各任务状态(进度冻结在中途)───────────────────")
    print(f"  {'task_id':<8}{'命令':<26}{'状态':<12}{'进度':>6}")
    for s in tm.all_states():
        print(f"  {s.task_id:<8}{s.command:<26}{s.status:<12}{s.progress:>5.0f}%")
    print("  ─────────────────────────────────────────────────────────")

    # 4) 恢复:executor 依然健康,接受并跑完一个新任务
    log("SYSTEM", "打断处理完毕,系统恢复空闲,可继续接受新任务……")
    fresh = tm.start("python re_run_summary.py")
    await fresh._task
    log("AGENT", f"已从打断中恢复,新任务 {fresh.task_id} 正常完成:"
                 f"{completed[-1].result[:36]}……")
    print("  结论:打断只冻结被取消的任务,运行时本身无损,可立即继续工作。\n")


# ============================ 3. 状态检查点:持久化 / 恢复 ============================

def _seed_trajectory(rt: AgentRuntime) -> None:
    """给运行时灌入一段「已发生」的对话轨迹,模拟会话进行到一半。"""
    rt._append(Event(EventType.USER_INPUT,
                     message={"role": "user", "content": "分析今天的日志并总结异常"},
                     label="用户消息:分析日志"))
    rt._append(Event(EventType.AGENT_TOOL_CALL,
                     message={"role": "assistant", "content": "好的,我这就在后台启动分析。",
                              "tool_calls": [{"id": "call_1", "type": "function",
                                              "function": {"name": "run_terminal_command",
                                                           "arguments": '{"command": "python analyze_fast.py"}'}}]},
                     label="调用工具 run_terminal_command"))
    rt._append(Event(EventType.TOOL_RESULT,
                     message={"role": "tool", "tool_call_id": "call_1",
                              "content": "命令已在后台异步启动。task_id=T1。"},
                     label="工具结果 run_terminal_command"))


async def demo_state() -> None:
    banner("能力三|状态管理:检查点持久化与跨会话恢复")
    tasks.TICK_REAL = 0.15
    ckpt_dir = os.path.join(os.path.dirname(os.path.abspath(__file__)), "checkpoints")
    os.makedirs(ckpt_dir, exist_ok=True)
    path = os.path.join(ckpt_dir, "agent_state.json")

    # —— 会话 A:产生一段轨迹 + 两个仍在运行的后台任务,然后落盘 ——
    log = Logger()
    log("SYSTEM", "会话 A 开始:构造轨迹并启动两个后台任务……")
    rt_a = AgentRuntime(client=None, model="demo-offline")
    rt_a._t0 = log.t0  # 让两个运行时共用同一时间基准,便于观察
    _seed_trajectory(rt_a)
    rt_a.tasks.start("python analyze_fast.py")   # 进行中
    rt_a.tasks.start("python analyze_slow.py")   # 进行中
    await asyncio.sleep(1.2)  # 让进度累积到中途

    before_traj = len(rt_a.trajectory)
    before_tasks = {s.task_id: (s.status, s.progress) for s in rt_a.tasks.all_states()}
    rt_a.save_checkpoint(path)

    # 模拟进程退出:取消掉活着的协程
    rt_a.tasks.cancel_all()
    await asyncio.sleep(0.05)
    log("SYSTEM", "会话 A 结束(进程退出,内存中的运行时已销毁)。")

    # —— 会话 B:全新运行时,从磁盘恢复 ——
    log("SYSTEM", "会话 B 开始:新建空运行时,从检查点恢复……")
    rt_b = AgentRuntime(client=None, model="demo-offline")
    rt_b._t0 = log.t0
    data = rt_b.load_checkpoint(path)

    after_traj = len(rt_b.trajectory)
    msgs = rt_b.build_messages()  # 证明恢复后能重建可喂给 LLM 的上下文

    print("\n  ── 恢复校验 ─────────────────────────────────────────────")
    print(f"  轨迹事件数     保存前 {before_traj}  ->  恢复后 {after_traj}  "
          f"[{'一致 ✓' if before_traj == after_traj else '不一致 ✗'}]")
    print(f"  可重建 LLM 上下文消息 {len(msgs)} 条(system + 轨迹回放)")
    print(f"  {'task_id':<8}{'命令':<26}{'保存前进度':>10}  {'恢复后状态':<12}{'进度':>6}")
    for rec in data["tasks"]:
        tid = rec["task_id"]
        before = before_tasks.get(tid, ("-", 0.0))
        st = rt_b.tasks.query(tid)
        print(f"  {tid:<8}{rec['command']:<26}{before[1]:>9.0f}%  "
              f"{st.status:<12}{st.progress:>5.0f}%")
    print("  ─────────────────────────────────────────────────────────")
    print(f"  检查点文件:{path}")
    print("  结论:轨迹与任务进度完整落盘并跨会话还原;运行中的任务被标记为 suspended,")
    print("        保留了最后已知进度,供上层决定「重跑」还是「按进度续跑」。\n")


# 供 demo.py 复用的离线演示注册表
OFFLINE_DEMOS = {
    "parallel": demo_parallel,
    "interrupt": demo_interrupt,
    "state": demo_state,
}

demo.py

"""实验 4-5 命令行入口:带并行执行、打断/取消与状态管理的异步 Agent。

本脚本提供两类演示,用子命令区分:

  【离线演示】不需要任何 API key,直接测量异步运行时的底层行为——
      python demo.py parallel     并行 vs 串行工具调用的墙钟时间对比(打印加速比)
      python demo.py interrupt    长任务运行中被打断/取消,随后系统恢复
      python demo.py state        Agent 状态检查点持久化 + 跨会话恢复并校验
      python demo.py offline       依次运行上面全部三个离线演示(默认行为)

  【LLM 场景】需要 OPENAI_API_KEY(或 MOONSHOT/ARK),由真实模型做决策——
      python demo.py scenarios              依次运行书中四个验证场景
      python demo.py scenarios --scenario 1  只跑场景 1(异步执行 + 即时提问)
      python demo.py scenarios --scenario 3  只跑场景 3(打断机制)

不带任何子命令时运行【离线演示】,因此开箱即用、无需联网。
为兼容旧用法,`python demo.py --scenario N` 等价于 `scenarios --scenario N`。
"""

from __future__ import annotations

import argparse
import asyncio
import os
import sys
import time

try:
    from dotenv import load_dotenv
    load_dotenv()
except Exception:
    pass

from async_demos import OFFLINE_DEMOS, banner
from runtime import AgentRuntime

# openai 仅在运行 LLM 场景时才惰性导入;离线演示不碰它,保证无 key/无 openai 也能跑。


def _completion_params_for(model: str) -> dict:
    """按模型返回安全的采样参数。

    Moonshot kimi-k3 是【推理模型】:必须 temperature=1 且 max_tokens>=2048,
    否则可能报错或截断。其余模型用 temperature=0.2 保证决策稳定。
    """
    if model.startswith("kimi-k3"):
        return {"temperature": 1, "max_tokens": 4096}
    return {"temperature": 0.2}


def _map_model_for_openrouter(model: str) -> str:
    """把常见模型名映射成 OpenRouter 的 `provider/model` 形式。

    - 已含 "/" 的 id(如 anthropic/claude-opus-4.8、google/gemini-2.5-pro)原样透传。
    - gpt-*/o1-*/o3-*/o4-* -> openai/…
    - claude-* -> anthropic/claude-opus-4.8
    - 其它保持原样(交给 OpenRouter 校验)。
    """
    if "/" in model:
        return model
    m = model.lower()
    if m.startswith(("gpt-", "o1-", "o3-", "o4-")):
        return f"openai/{model}"
    if m.startswith("claude-"):
        return "anthropic/claude-opus-4.8"
    return model


def make_client():
    """按 LLM_PROVIDER 选择可用的模型服务(默认 openai)。

    返回 (client, model, completion_params)。

    通用兜底:当直连 provider 的 key 缺失、但存在 OPENROUTER_API_KEY 时,
    自动改走 OpenRouter(api_key=OPENROUTER_API_KEY,base_url=openrouter.ai/api/v1,
    并把模型名映射成 provider/model 形式),从而"有 OpenRouter key 就能跑"。
    """
    from openai import AsyncOpenAI  # 惰性导入:离线演示无需安装 openai
    provider = os.getenv("LLM_PROVIDER", "openai").lower()
    if provider == "moonshot":
        key = os.environ["MOONSHOT_API_KEY"]
        # 默认用当前的推理模型 kimi-k3(旧的 kimi-k2-*-preview 与 moonshot-v1-* 均已过时/停用)。
        model = os.getenv("LLM_MODEL", "kimi-k3")
        client = AsyncOpenAI(api_key=key, base_url="https://api.moonshot.cn/v1")
        return client, model, _completion_params_for(model)
    if provider == "ark":
        key = os.environ["ARK_API_KEY"]
        model = os.getenv("LLM_MODEL")  # ARK 需要填 endpoint id
        if not model:
            raise SystemExit("使用 ARK 时请设置 LLM_MODEL 为你的推理接入点 ID")
        client = AsyncOpenAI(api_key=key, base_url="https://ark.cn-beijing.volces.com/api/v3")
        return client, model, _completion_params_for(model)
    if provider == "openrouter":
        key = os.environ["OPENROUTER_API_KEY"]
        model = _map_model_for_openrouter(os.getenv("LLM_MODEL", "openai/gpt-5.6-luna"))
        client = AsyncOpenAI(api_key=key, base_url="https://openrouter.ai/api/v1")
        return client, model, _completion_params_for(model)
    key = os.getenv("OPENAI_API_KEY")
    or_key = os.getenv("OPENROUTER_API_KEY")
    model = os.getenv("LLM_MODEL", "gpt-5.6-luna")
    # gpt-5.x(含 gpt-5.6*)直连 OpenAI 需要组织验证;只要有 OPENROUTER_API_KEY,
    # 就优先走 OpenRouter;直连 OPENAI_API_KEY 缺失时同样兜底到 OpenRouter。
    if or_key and (not key or model.lower().startswith("gpt-5")):
        mapped = _map_model_for_openrouter(model)
        client = AsyncOpenAI(api_key=or_key, base_url="https://openrouter.ai/api/v1")
        return client, mapped, _completion_params_for(mapped)
    if key:
        base = os.getenv("OPENAI_BASE_URL")
        client = AsyncOpenAI(api_key=key, base_url=base) if base else AsyncOpenAI(api_key=key)
        return client, model, _completion_params_for(model)
    raise SystemExit(
        "未找到可用的 LLM Key。请设置以下任意一项:"
        "OPENAI_API_KEY 或 OPENROUTER_API_KEY(或 LLM_PROVIDER=moonshot 且 MOONSHOT_API_KEY / "
        "LLM_PROVIDER=ark 且 ARK_API_KEY)。"
    )


async def run_runtime(rt: AgentRuntime):
    """在后台跑事件循环。"""
    return asyncio.create_task(rt.serve())


# ------------------------------- 四个场景 -------------------------------

async def scenario_1(client, model, params):
    banner("场景 1|异步工具执行:长任务运行期间即时回应插入的提问")
    rt = AgentRuntime(client, model, completion_params=params)
    serve = await run_runtime(rt)

    # 用户下达一个耗时的日志分析任务
    await rt.submit_user_message(
        "请运行终端命令 `python analyze_logs.py`(这是耗时的日志分析),完成后给我分析结论。",
        urgency="immediate")
    await asyncio.sleep(2.2)  # 任务已在后台跑

    # 期间用户插入一个即时问题
    await rt.submit_user_message("现在几点了?")  # 带问号 -> 立即回应

    await rt.wait_until_idle()
    await rt.stop(); await serve


async def scenario_2(client, model, params):
    banner("场景 2|事件队列与批量处理:非紧急指令累积,任务完成时一次性处理")
    rt = AgentRuntime(client, model, completion_params=params)
    serve = await run_runtime(rt)

    await rt.submit_user_message(
        "请运行终端命令 `python analyze_logs.py`(耗时日志分析),完成后把分析结论告诉我。",
        urgency="immediate")
    await asyncio.sleep(1.5)

    # 连续发两条补充性指令(无问号 -> 非紧急,进入排队缓冲)
    await rt.submit_user_message("记得最后用日语回复")
    await asyncio.sleep(0.4)
    await rt.submit_user_message("把结果整理成一个网页(HTML)")

    await rt.wait_until_idle()
    await rt.stop(); await serve


async def scenario_3(client, model, params):
    banner("场景 3|打断机制:用户'取消'立即终止执行流并取消异步工具")
    rt = AgentRuntime(client, model, completion_params=params)
    serve = await run_runtime(rt)

    await rt.submit_user_message(
        "请运行终端命令 `python analyze_logs.py`(耗时日志分析),完成后给我结论。",
        urgency="immediate")
    await asyncio.sleep(4.0)  # 等后台任务确实跑起来(跑到一半左右)

    await rt.submit_user_message("取消")  # 打断关键词 -> 立即取消

    await rt.wait_until_idle(stable=1.0)
    await rt.stop(); await serve


async def scenario_4(client, model, params):
    banner("场景 4|并行工具的取消与状态查询:三脚本竞速 + 按 50% 阈值取消 + 整合报告")
    rt = AgentRuntime(client, model, completion_params=params)
    serve = await run_runtime(rt)

    await rt.submit_user_message(
        "同时运行这三个分析脚本:`python analyze_fast.py`、`python analyze_mid.py`、`python analyze_slow.py`。"
        "哪个脚本先完成,你就查询另外两个脚本的进度;如果某个脚本进度还没超过 50%,就取消它;"
        "其余脚本完成后,把所有已完成脚本的结果整合成一份报告给我。",
        urgency="immediate")

    await rt.wait_until_idle(stable=1.5, timeout=60)
    await rt.stop(); await serve


SCENARIOS = {1: scenario_1, 2: scenario_2, 3: scenario_3, 4: scenario_4}


# ------------------------------- 子命令实现 -------------------------------

async def run_offline(names: list[str]) -> None:
    """运行离线演示(无需 API key)。"""
    for name in names:
        await OFFLINE_DEMOS[name]()


async def run_scenarios(which: int | None) -> None:
    """运行 LLM 驱动的验证场景(需要 API key)。"""
    client, model, params = make_client()
    print(f"使用模型:{model}")
    todo = [which] if which else [1, 2, 3, 4]
    for i in todo:
        await SCENARIOS[i](client, model, params)
        await asyncio.sleep(0.5)


def build_parser() -> argparse.ArgumentParser:
    parser = argparse.ArgumentParser(
        prog="demo.py",
        formatter_class=argparse.RawDescriptionHelpFormatter,
        description="实验 4-5:带并行执行、打断/取消与状态管理的异步 Agent 演示。",
        epilog=(
            "示例:\n"
            "  python demo.py                     # 默认:依次运行三个离线演示(无需 API key)\n"
            "  python demo.py parallel            # 并行 vs 串行的墙钟时间对比(打印加速比)\n"
            "  python demo.py interrupt           # 长任务运行中被打断/取消,随后恢复\n"
            "  python demo.py state               # 状态检查点持久化 + 跨会话恢复并校验\n"
            "  python demo.py scenarios --scenario 3   # LLM 场景 3:打断机制(需 API key)\n"
            "\n离线演示不联网、不需要任何 key;scenarios 子命令需要 OPENAI_API_KEY(或 MOONSHOT/ARK)。"
        ),
    )
    sub = parser.add_subparsers(dest="command", metavar="<子命令>")

    sub.add_parser("parallel", help="并行 vs 串行工具调用的墙钟时间对比(离线,无需 key)")
    sub.add_parser("interrupt", help="长任务运行中被打断/取消,随后系统恢复(离线,无需 key)")
    sub.add_parser("state", help="Agent 状态检查点持久化与跨会话恢复(离线,无需 key)")
    sub.add_parser("offline", help="依次运行上面三个离线演示(默认行为)")

    ps = sub.add_parser("scenarios", help="书中四个 LLM 验证场景(需要 API key)")
    ps.add_argument("--scenario", type=int, choices=[1, 2, 3, 4],
                    help="只运行指定场景(1 异步执行 / 2 批量处理 / 3 打断 / 4 并行取消);不填则全部")
    return parser


async def main() -> None:
    # 兼容旧用法:`python demo.py --scenario N` 等价于 `scenarios --scenario N`
    argv = sys.argv[1:]
    if argv and argv[0].startswith("-") and argv[0] not in ("-h", "--help"):
        argv = ["scenarios"] + argv

    args = build_parser().parse_args(argv)
    cmd = args.command or "offline"

    if cmd == "scenarios":
        await run_scenarios(args.scenario)
    elif cmd == "offline":
        await run_offline(["parallel", "interrupt", "state"])
    else:  # parallel / interrupt / state
        await run_offline([cmd])

    print("\n演示结束。")


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

events.py

"""事件模型(对应设计文档中的 Event / Trajectory 概念)。

Flux 把 Agent 的一切经历都抽象成"事件",按时间顺序追加到轨迹(trajectory)里。
本文件定义事件类型、事件对象,以及"事件紧急度"的判定逻辑——这是实验 4-5 里
"批量处理 vs 立即打断"两种处理机制的分类依据。
"""

from __future__ import annotations

import time
from dataclasses import dataclass, field
from typing import Optional


class EventType:
    """事件类型常量(对应设计文档第 2 节 Inputs / Interrupts / Thinking / Actions)。"""

    USER_INPUT = "user.input"          # 用户输入(非紧急,走"排队处理")
    USER_INTERRUPT = "user.interrupt"  # 用户打断(紧急,走"取消式处理")
    AGENT_OUTPUT = "agent.output"      # Agent 面向用户的最终回复
    AGENT_TOOL_CALL = "agent.tool_call"  # Agent 发起的工具调用(Action)
    TOOL_RESULT = "tool.result"        # 工具返回结果(同步工具 / 异步占位符)
    ASYNC_RESULT = "async.result"      # 异步工具真正完成后注入的新事件
    SYSTEM_NOTE = "system.note"        # 框架注入的系统提示(如取消回执)


class Urgency:
    """事件紧急度:决定采用哪种事件处理机制。"""

    INTERRUPT = "interrupt"  # 取消式处理:立刻打断当前执行并取消异步工具
    IMMEDIATE = "immediate"  # 立即处理:不打断后台异步任务,但马上回应(如用户提问)
    DEFERRED = "deferred"    # 排队处理:累积到 pending 队列,任务完成时批量追加


# 打断类关键词:命中即视为紧急打断
_INTERRUPT_KEYWORDS = ["取消", "停止", "中止", "打住", "别做了", "stop", "cancel", "abort"]

# 疑问类信号:命中即视为需要"立即回应"(而不是排队)
_QUESTION_MARKS = ("?", "?")
_QUESTION_KEYWORDS = ["几点", "多少", "怎么", "如何", "为什么", "是不是", "有没有",
                      "吗", "呢", "what", "when", "how", "why", "which"]


def classify_urgency(text: str) -> str:
    """根据用户消息内容判定紧急度。

    规则(简单、可解释,便于书中讲清楚):
      1. 含打断关键词(取消/停止/stop...) -> INTERRUPT(紧急,取消式处理)
      2. 是一个提问(带问号或疑问词)      -> IMMEDIATE(立即回应,但不打断后台任务)
      3. 其它(补充性指令,如"用日语回复")-> DEFERRED(排队,批量处理)
    """
    low = text.lower()
    if any(kw in text or kw in low for kw in _INTERRUPT_KEYWORDS):
        return Urgency.INTERRUPT
    if text.strip().endswith(_QUESTION_MARKS) or any(kw in text or kw in low for kw in _QUESTION_KEYWORDS):
        return Urgency.IMMEDIATE
    return Urgency.DEFERRED


@dataclass
class Event:
    """一条轨迹事件。

    message 字段保存"可直接喂给 LLM 的 OpenAI 消息字典"(保证上下文的高保真回放);
    没有 message 的事件(若有)只用于日志。
    """

    type: str
    message: Optional[dict] = None        # OpenAI chat 格式消息,供构建 LLM 上下文
    label: str = ""                       # 人类可读的日志标签
    task_id: Optional[str] = None         # 关联的异步任务 ID(若有)
    urgency: Optional[str] = None         # 仅用户输入事件会带
    ts: float = field(default_factory=time.time)

    def to_dict(self) -> dict:
        """序列化为纯 JSON 可写的字典(用于状态检查点持久化)。"""
        return {
            "type": self.type, "message": self.message, "label": self.label,
            "task_id": self.task_id, "urgency": self.urgency, "ts": self.ts,
        }

    @classmethod
    def from_dict(cls, d: dict) -> "Event":
        """从检查点字典还原事件对象。"""
        return cls(
            type=d["type"], message=d.get("message"), label=d.get("label", ""),
            task_id=d.get("task_id"), urgency=d.get("urgency"),
            ts=d.get("ts", time.time()),
        )

runtime.py

"""Flux 异步 Agent 运行时(实验 4-5 核心)。

实现设计文档第 5 节的事件处理循环,重点覆盖实验 4-5 的四个能力:
  1. 异步工具执行:run_terminal_command 立即返回占位符,任务在后台跑。
  2. 事件队列与批量处理:非紧急事件进 pending,异步结果到达时一次性批量追加。
  3. 打断机制:用户"取消/停止"立即取消当前 turn + 所有异步工具,并留痕。
  4. 并行工具的取消与状态查询:query_task / cancel_task 按 ID 操作;
     异步完成后以"新事件"把真实结果注入对话。

架构(三个协程协作,全部基于 asyncio 单线程):
  - inbox 队列:所有进来的事件(用户输入、打断、异步完成通知)先入 inbox。
  - _dispatcher:从 inbox 取事件 -> 判定紧急度 -> 分流(立即处理 / 排队 / 打断)。
  - _worker   :从 work 队列取"事件批次" -> 追加到轨迹 -> 跑一轮 LLM(run_llm_turn)。
    每一轮 LLM 作为可取消的子任务(turn_task),打断时直接 cancel 它。
"""

from __future__ import annotations

import asyncio
import datetime
import json
import time
from typing import Optional

from events import Event, EventType, Urgency, classify_urgency
from tasks import TaskManager, TaskState

# ------------------------- LLM 工具定义(function calling) -------------------------

TOOL_SCHEMAS = [
    {
        "type": "function",
        "function": {
            "name": "run_terminal_command",
            "description": ("异步执行一个(模拟的)耗时终端命令,例如日志分析脚本。"
                            "调用后命令在后台运行,本工具立即返回一个 task_id 占位符,"
                            "不会阻塞。任务真正完成后,其结果会作为一条新的系统事件出现在对话中。"),
            "parameters": {
                "type": "object",
                "properties": {
                    "command": {"type": "string", "description": "要执行的终端命令,如 `python analyze_logs.py`"},
                },
                "required": ["command"],
            },
        },
    },
    {
        "type": "function",
        "function": {
            "name": "get_current_time",
            "description": "立即返回当前时间。用于回答用户'现在几点了'之类的即时问题。",
            "parameters": {"type": "object", "properties": {}},
        },
    },
    {
        "type": "function",
        "function": {
            "name": "query_task",
            "description": "查询某个后台异步任务的当前进度与状态。",
            "parameters": {
                "type": "object",
                "properties": {"task_id": {"type": "string", "description": "任务 ID,如 T1"}},
                "required": ["task_id"],
            },
        },
    },
    {
        "type": "function",
        "function": {
            "name": "cancel_task",
            "description": "按 task_id 取消一个正在运行的后台异步任务。",
            "parameters": {
                "type": "object",
                "properties": {"task_id": {"type": "string", "description": "任务 ID,如 T1"}},
                "required": ["task_id"],
            },
        },
    },
]

SYSTEM_PROMPT = """你是一个异步 Agent(基于 Flux 框架)。你可以调用工具来完成任务。

关键行为准则:
1. run_terminal_command 是【异步】的:调用后命令在后台运行并立即返回 task_id。
   你应当简要告知用户"任务已在后台启动",然后【结束本轮回复,不要空等结果】。
2. 当你看到形如 "[系统事件|异步任务完成] task_id=... 结果:..." 的消息时,
   说明后台任务真的完成了,这时再基于结果给出分析/整合结论。
3. 如果用户在后台任务运行期间提出简短问题(例如"现在几点了?"),
   立即用对应工具(如 get_current_time)回答,【不要等待】后台任务。
4. 你可以用 query_task 查询任意后台任务进度,用 cancel_task 按 ID 取消任务。
5. 收到 "[用户打断]" 时,立即停止当前工作并简短确认已停止。
6. 严格按用户给出的计划执行(例如"谁先完成就查其余进度,未过 50% 就取消")。
   注意:只取消【进度未超过 50%】的任务;进度已超过 50% 的任务应【保留并等待其完成】,不要取消它。
   每个还在运行的任务只需查询一次进度即可做出取消/保留决定,不要反复查询。
7. 回答简洁、用中文,除非用户明确要求其它语言或格式。
"""

MAX_STEPS = 8  # 单轮内最多的工具调用往返次数(防止死循环)

# 日志配色(各来源一种颜色),供 runtime 与离线演示脚本共用。
_LOG_COLORS = {
    "USER": "\033[96m", "AGENT": "\033[92m", "TOOL": "\033[93m",
    "TASK": "\033[95m", "SYSTEM": "\033[90m", "TRAJ": "\033[94m",
    "STATE": "\033[95m",
}


def format_log(t0: float, source: str, text: str) -> str:
    """把一条日志渲染成「[相对秒] 来源 | 文本」的彩色字符串。"""
    color = _LOG_COLORS.get(source, "")
    reset = "\033[0m" if color else ""
    return f"[{time.time() - t0:6.2f}s] {color}{source:6}{reset} | {text}"


class AgentRuntime:
    def __init__(self, client, model: str, start_time: Optional[float] = None,
                 completion_params: Optional[dict] = None):
        self.client = client
        self.model = model
        # 传给 chat.completions.create 的采样参数。默认 temperature=0.2 适合 gpt-5.6-luna;
        # 推理模型(如 Moonshot kimi-k3)需要 temperature=1 且 max_tokens>=2048,由 make_client 传入。
        self.completion_params = completion_params or {"temperature": 0.2}
        self._t0 = start_time or time.time()

        self.trajectory: list[Event] = []          # 轨迹(工作记忆)
        self.inbox: asyncio.Queue = asyncio.Queue()  # 所有进来的原始事件
        self.work: asyncio.Queue = asyncio.Queue()   # 待处理的事件批次
        self.pending: list[Event] = []               # 非紧急事件的排队缓冲

        self.tasks = TaskManager(on_complete=self._on_task_complete, log=self.log)
        self.turn_task: Optional[asyncio.Task] = None
        self.running = True
        self._STOP = object()

    # ------------------------------- 日志 -------------------------------

    def log(self, source: str, text: str) -> None:
        print(format_log(self._t0, source, text), flush=True)

    def _append(self, event: Event) -> None:
        """把事件追加到轨迹,并打印轨迹留痕。"""
        self.trajectory.append(event)
        self.log("TRAJ", f"+ {event.type:18} {event.label}")

    def build_messages(self) -> list[dict]:
        """把轨迹渲染成 OpenAI chat 消息列表。"""
        msgs = [{"role": "system", "content": SYSTEM_PROMPT}]
        for e in self.trajectory:
            if e.message:
                msgs.append(e.message)
        return msgs

    # ------------------------- 对外接口:提交事件 -------------------------

    async def submit_user_message(self, text: str, urgency: Optional[str] = None) -> None:
        """提交一条用户消息(demo 用它模拟用户输入)。"""
        u = urgency or classify_urgency(text)
        if u == Urgency.INTERRUPT:
            ev = Event(EventType.USER_INTERRUPT, urgency=u,
                       message={"role": "user", "content": f"[用户打断] {text}"},
                       label=f"用户打断:{text}")
        else:
            ev = Event(EventType.USER_INPUT, urgency=u,
                       message={"role": "user", "content": text},
                       label=f"用户消息({u}):{text}")
        self.log("USER", f"({u}) {text}")
        await self.inbox.put(ev)

    async def _on_task_complete(self, state: TaskState) -> None:
        """异步任务自然完成 -> 把真实结果作为【新事件】注入 inbox。"""
        ev = Event(
            EventType.ASYNC_RESULT, task_id=state.task_id,
            message={"role": "user",
                     "content": (f"[系统事件|异步任务完成] task_id={state.task_id} "
                                 f"命令=`{state.command}` 结果:{state.result}")},
            label=f"异步完成 {state.task_id}",
        )
        await self.inbox.put(ev)

    # ------------------------------- 主循环 -------------------------------

    async def serve(self) -> None:
        dispatcher = asyncio.create_task(self._dispatcher())
        worker = asyncio.create_task(self._worker())
        await asyncio.gather(dispatcher, worker)

    def _is_idle(self) -> bool:
        return (not self.tasks.any_running()
                and self.work.empty()
                and self.inbox.empty()
                and (self.turn_task is None or self.turn_task.done()))

    def _drain_pending(self) -> list[Event]:
        drained, self.pending = self.pending, []
        return drained

    async def _dispatcher(self) -> None:
        """事件分流:实现设计文档 5.1 的两种处理机制。"""
        while self.running:
            ev = await self.inbox.get()
            if ev is self._STOP:
                await self.work.put(self._STOP)
                break

            if ev.type == EventType.USER_INTERRUPT:
                # —— 取消式处理:立刻打断当前 turn + 取消所有异步工具 ——
                await self._handle_interrupt(ev)

            elif ev.type == EventType.ASYNC_RESULT:
                # —— 异步结果到达:批量把 pending 一并追加,再触发 LLM ——
                batch = [ev] + self._drain_pending()
                if len(batch) > 1:
                    self.log("SYSTEM", f"异步结果到达,批量处理 {len(batch)-1} 条积压的非紧急事件")
                await self.work.put(batch)

            elif ev.type == EventType.USER_INPUT:
                if ev.urgency == Urgency.IMMEDIATE:
                    # 立即处理(如用户提问),不打断后台异步任务
                    await self.work.put([ev])
                elif self._is_idle():
                    # 空闲时,普通指令也直接处理(例如一开始下达的任务)
                    await self.work.put([ev])
                else:
                    # 排队处理:累积到 pending,等下一次异步结果时批量追加
                    self.pending.append(ev)
                    self.log("SYSTEM", f"事件进入排队缓冲(当前积压 {len(self.pending)} 条)")

    async def _handle_interrupt(self, ev: Event) -> None:
        # 1) 取消正在进行的 LLM turn
        if self.turn_task and not self.turn_task.done():
            self.turn_task.cancel()
        # 2) 取消所有后台异步工具
        cancelled = self.tasks.cancel_all()
        # 3) 组装打断批次:打断事件 + 系统回执 + 被丢弃的积压事件(留痕)
        note = Event(
            EventType.SYSTEM_NOTE,
            message={"role": "user",
                     "content": (f"[系统] 已执行打断:取消了后台任务 {cancelled or '(无)'}。"
                                 f"请向用户简短确认已停止。")},
            label=f"打断回执,取消任务 {cancelled or '(无)'}",
        )
        batch = [ev, note] + self._drain_pending()
        await self.work.put(batch)

    async def _worker(self) -> None:
        """逐批处理事件:追加到轨迹后跑一轮可被取消的 LLM。"""
        while self.running:
            batch = await self.work.get()
            if batch is self._STOP:
                break
            self.turn_task = asyncio.create_task(self._process_batch(batch))
            try:
                await self.turn_task
            except asyncio.CancelledError:
                self.log("SYSTEM", "当前 LLM turn 已被打断取消")

    async def _process_batch(self, batch: list[Event]) -> None:
        for e in batch:
            self._append(e)
        await self.run_llm_turn()

    # ------------------------------- LLM turn -------------------------------

    async def run_llm_turn(self) -> None:
        """调用 LLM 做决策;同步工具就地执行并回填,异步工具启动后回占位符。"""
        for _ in range(MAX_STEPS):
            messages = self.build_messages()
            _t = time.time()
            resp = await self.client.chat.completions.create(
                model=self.model, messages=messages,
                tools=TOOL_SCHEMAS, tool_choice="auto", **self.completion_params,
            )
            self.log("SYSTEM", f"LLM 调用耗时 {time.time()-_t:.2f}s({len(messages)} 条消息)")
            msg = resp.choices[0].message

            assistant_msg: dict = {"role": "assistant", "content": msg.content or ""}
            if msg.tool_calls:
                assistant_msg["tool_calls"] = [
                    {"id": tc.id, "type": "function",
                     "function": {"name": tc.function.name, "arguments": tc.function.arguments}}
                    for tc in msg.tool_calls
                ]

            self._append(Event(
                EventType.AGENT_TOOL_CALL if msg.tool_calls else EventType.AGENT_OUTPUT,
                message=assistant_msg,
                label=("调用工具 " + ", ".join(tc.function.name for tc in msg.tool_calls)
                       if msg.tool_calls else "回复用户"),
            ))

            if msg.content and msg.content.strip():
                self.log("AGENT", msg.content.strip())

            if not msg.tool_calls:
                return  # 本轮结束:Agent 给出了最终回复

            # 执行每个工具调用
            for tc in msg.tool_calls:
                name = tc.function.name
                try:
                    args = json.loads(tc.function.arguments or "{}")
                except json.JSONDecodeError:
                    args = {}
                result_text = self._exec_tool(name, args)
                self._append(Event(
                    EventType.TOOL_RESULT,
                    message={"role": "tool", "tool_call_id": tc.id, "content": result_text},
                    label=f"工具结果 {name}",
                ))

    def _exec_tool(self, name: str, args: dict) -> str:
        """执行工具,返回给 LLM 的文本结果。"""
        if name == "run_terminal_command":
            command = args.get("command", "")
            state = self.tasks.start(command)
            return (f"命令已在后台【异步】启动。task_id={state.task_id},命令=`{command}`。"
                    f"我不会阻塞等待;任务完成后其结果会以系统事件形式返回。"
                    f"可用 query_task('{state.task_id}') 查询进度或 cancel_task('{state.task_id}') 取消。")

        if name == "get_current_time":
            now = datetime.datetime.now().strftime("%Y-%m-%d %H:%M:%S")
            self.log("TOOL", f"get_current_time -> {now}")
            return f"当前时间是 {now}。"

        if name == "query_task":
            tid = args.get("task_id", "")
            st = self.tasks.query(tid)
            if not st:
                return f"未找到任务 {tid}。"
            self.log("TOOL", f"query_task({tid}) -> {st.status} {st.progress:.0f}%")
            return f"task_id={tid} 命令=`{st.command}` 状态={st.status} 进度={st.progress:.0f}%。"

        if name == "cancel_task":
            tid = args.get("task_id", "")
            st = self.tasks.query(tid)
            progress = f"{st.progress:.0f}%" if st else "未知"
            ok = self.tasks.cancel(tid)
            self.log("TOOL", f"cancel_task({tid}) -> {'已取消' if ok else '无法取消'} (进度 {progress})")
            return (f"任务 {tid} 已取消(取消时进度 {progress})。" if ok
                    else f"任务 {tid} 无法取消(可能已完成或不存在)。")

        return f"未知工具:{name}"

    # ------------------------------- 收尾 -------------------------------

    async def wait_until_idle(self, stable: float = 1.3, timeout: float = 90.0) -> None:
        """阻塞直到系统持续空闲 stable 秒(或超时)。"""
        start = time.time()
        last_busy = time.time()
        while True:
            busy = (self.tasks.any_running() or not self.work.empty()
                    or not self.inbox.empty() or bool(self.pending)
                    or (self.turn_task is not None and not self.turn_task.done()))
            now = time.time()
            if busy:
                last_busy = now
            elif now - last_busy >= stable:
                return
            if now - start >= timeout:
                self.log("SYSTEM", "wait_until_idle 超时返回")
                return
            await asyncio.sleep(0.1)

    async def stop(self) -> None:
        self.running = False
        await self.inbox.put(self._STOP)

    # ------------------------- 状态检查点(持久化 / 恢复) -------------------------

    def snapshot(self) -> dict:
        """把 Agent 的可持久化状态导出为一个 JSON 友好的字典。

        状态 = 轨迹(工作记忆)+ 全部异步任务的最后已知状态。这是「跨会话恢复」
        的基础:进程重启后,能据此还原对话上下文与后台任务的进度。
        """
        return {
            "model": self.model,
            "saved_at": datetime.datetime.now().isoformat(timespec="seconds"),
            "trajectory": [e.to_dict() for e in self.trajectory],
            "tasks": self.tasks.snapshot(),
        }

    def save_checkpoint(self, path: str) -> str:
        """把当前状态写入检查点文件(JSON),返回文件路径。"""
        data = self.snapshot()
        with open(path, "w", encoding="utf-8") as f:
            json.dump(data, f, ensure_ascii=False, indent=2)
        self.log("STATE", f"已保存检查点 -> {path}"
                          f"({len(data['trajectory'])} 条轨迹事件,{len(data['tasks'])} 个任务)")
        return path

    def load_checkpoint(self, path: str) -> dict:
        """从检查点文件恢复轨迹与任务状态(原地覆盖当前状态)。"""
        with open(path, "r", encoding="utf-8") as f:
            data = json.load(f)
        self.trajectory = [Event.from_dict(d) for d in data.get("trajectory", [])]
        self.tasks.restore(data.get("tasks", []))
        self.log("STATE", f"已从检查点恢复 <- {path}"
                          f"({len(self.trajectory)} 条轨迹事件,{len(data.get('tasks', []))} 个任务)")
        return data

tasks.py

"""模拟的异步"终端命令"与任务管理器。

为了安全(绝不真跑危险命令)与可复现,长任务用"带进度输出的模拟脚本"实现:
每个模拟脚本以固定的"每(模拟)秒进度百分比"推进,直到 100% 完成。

时间轴加速:真实世界里每 TICK_REAL 秒代表 1 个"模拟秒"。
默认 TICK_REAL=0.4,即 2.5 倍速——保留"3%/2%/1% 的速度差 + 是否过 50% 的判定"逻辑,
但把几十秒的等待压缩到几秒,方便演示复现。
"""

from __future__ import annotations

import asyncio
import os
from dataclasses import dataclass, field
from typing import Awaitable, Callable, Dict, Optional


def _env_float(name: str, default: float) -> float:
    """读取浮点环境变量;值非法时回退到默认值并打印警告。"""
    raw = os.getenv(name)
    if raw is None:
        return default
    try:
        return float(raw)
    except ValueError:
        print(f"⚠️ 环境变量 {name}={raw!r} 非法(应为数字),使用默认值 {default}")
        return default


# 1 个"模拟秒"对应的真实秒数(可用环境变量覆盖)。
# 默认 0.4(2.5 倍速):既压缩了等待,又给模型的"查询-判定-取消"决策留足时间窗口,
# 保证场景 4 里"慢脚本尚未过 50% 就被取消"能稳定复现。
TICK_REAL = _env_float("FLUX_TICK_REAL", 0.4)

# 不同脚本的"每模拟秒进度%"档位。实验 4-5 场景 4 需要 3% / 2% / 1% 的速度差。
_SCRIPT_RATES = [
    ("fast", 3.0),
    ("mid", 2.0),
    ("slow", 1.0),
]
_DEFAULT_RATE = 4.5  # 场景 1/2/3 的普通长任务(约 22 模拟秒 ≈ 5.5 真实秒完成)


def resolve_rate(command: str) -> float:
    """根据命令字符串推断该模拟脚本的推进速度(%/模拟秒)。"""
    low = command.lower()
    for key, rate in _SCRIPT_RATES:
        if key in low:
            return rate
    return _DEFAULT_RATE


@dataclass
class TaskState:
    """一个异步终端任务的实时状态。"""

    task_id: str
    command: str
    rate: float
    progress: float = 0.0
    status: str = "running"          # running | completed | cancelled
    result: str = ""
    _task: Optional[asyncio.Task] = field(default=None, repr=False)


class TaskManager:
    """管理所有异步终端任务:启动、查询进度、取消。

    on_complete 回调会在任务自然完成时被调用(用于把真实结果作为"新事件"注入对话)。
    """

    def __init__(self, on_complete: Callable[[TaskState], Awaitable[None]],
                 log: Callable[[str, str], None]):
        self._on_complete = on_complete
        self._log = log
        self._tasks: Dict[str, TaskState] = {}
        self._counter = 0

    def start(self, command: str) -> TaskState:
        """启动一个异步终端命令,立即返回其状态(含 task_id 占位符)。"""
        self._counter += 1
        task_id = f"T{self._counter}"
        state = TaskState(task_id=task_id, command=command, rate=resolve_rate(command))
        self._tasks[task_id] = state
        state._task = asyncio.create_task(self._run(state))
        self._log("TASK", f"启动异步任务 {task_id}: `{command}` (速度 {state.rate:.0f}%/模拟秒)")
        return state

    async def _run(self, state: TaskState) -> None:
        """后台推进进度,直到完成或被取消。"""
        next_milestone = 20.0
        try:
            while state.progress < 100.0:
                await asyncio.sleep(TICK_REAL)
                state.progress = min(100.0, state.progress + state.rate)
                if state.progress >= next_milestone:
                    self._log("TASK", f"{state.task_id} `{state.command}` 进度 {state.progress:.0f}%")
                    next_milestone += 20.0
            state.status = "completed"
            state.result = (f"命令 `{state.command}` 执行完毕:共扫描 12,840 条记录,"
                            f"发现 3 个异常峰值、1 处可疑错误码(HTTP 503 突增),"
                            f"平均响应时间 128ms。")
            self._log("TASK", f"{state.task_id} 完成 ✅")
            await self._on_complete(state)
        except asyncio.CancelledError:
            # 被取消:标记状态并静默退出(不再注入完成结果)
            state.status = "cancelled"
            self._log("TASK", f"{state.task_id} 已被取消 🛑(进度停在 {state.progress:.0f}%)")
            raise

    def query(self, task_id: str) -> Optional[TaskState]:
        return self._tasks.get(task_id)

    def cancel(self, task_id: str) -> bool:
        """按 ID 取消单个任务。"""
        state = self._tasks.get(task_id)
        if state and state.status == "running":
            state.status = "cancelled"
            if state._task:
                state._task.cancel()
            return True
        return False

    def cancel_all(self) -> list[str]:
        """取消所有仍在运行的任务,返回被取消的 task_id 列表。"""
        cancelled = []
        for tid, state in self._tasks.items():
            if state.status == "running":
                state.status = "cancelled"
                if state._task:
                    state._task.cancel()
                cancelled.append(tid)
        return cancelled

    def any_running(self) -> bool:
        return any(s.status == "running" for s in self._tasks.values())

    def all_states(self) -> list[TaskState]:
        return list(self._tasks.values())

    # --------------------------- 状态检查点(持久化) ---------------------------

    def snapshot(self) -> list[dict]:
        """导出所有任务的最后已知状态,供检查点持久化。

        注意:正在跑的 asyncio 协程无法序列化,只能记录其最后已知进度;
        重启后据此决定「重跑」还是「按进度续跑」,这正是异步任务状态管理的意义。
        """
        return [
            {"task_id": s.task_id, "command": s.command, "rate": s.rate,
             "progress": s.progress, "status": s.status, "result": s.result}
            for s in self._tasks.values()
        ]

    def restore(self, records: list[dict]) -> None:
        """从检查点还原任务的历史状态(不重启协程)。

        还原时把「运行中」的任务标记为 suspended(挂起)——它没有活着的协程,
        只保留了被打快照那一刻的进度,等待上层逻辑决定如何续跑。
        """
        for r in records:
            status = "suspended" if r["status"] == "running" else r["status"]
            state = TaskState(
                task_id=r["task_id"], command=r["command"], rate=r["rate"],
                progress=r["progress"], status=status, result=r.get("result", ""),
            )
            self._tasks[state.task_id] = state
            try:
                self._counter = max(self._counter, int(state.task_id.lstrip("T") or 0))
            except ValueError:
                pass

test_tasks_env.py

"""回归测试:FLUX_TICK_REAL 等浮点环境变量非法时不得让模块导入崩溃。

tasks.py 原来在模块导入时用裸 float() 解析 FLUX_TICK_REAL,
FLUX_TICK_REAL=abc 会让整个演示脚本以 ValueError 崩溃;现在回退到默认值并打印警告。
"""
import importlib
import os
import sys

sys.path.insert(0, os.path.dirname(__file__))

import tasks


def test_env_float_falls_back_on_malformed(monkeypatch, capsys):
    monkeypatch.setenv("FLUX_TICK_REAL", "abc")
    assert tasks._env_float("FLUX_TICK_REAL", 0.4) == 0.4
    assert "FLUX_TICK_REAL" in capsys.readouterr().out


def test_env_float_parses_valid_value(monkeypatch):
    monkeypatch.setenv("FLUX_TICK_REAL", "0.1")
    assert tasks._env_float("FLUX_TICK_REAL", 0.4) == 0.1


def test_env_float_default_when_unset(monkeypatch):
    monkeypatch.delenv("FLUX_TICK_REAL", raising=False)
    assert tasks._env_float("FLUX_TICK_REAL", 0.4) == 0.4


def test_module_reload_survives_malformed_env(monkeypatch):
    """模块级 TICK_REAL 在环境变量非法时不得抛出 ValueError。"""
    monkeypatch.setenv("FLUX_TICK_REAL", "fast")
    importlib.reload(tasks)
    assert tasks.TICK_REAL == 0.4