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,甚至不需要安装
openai,python 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。
紧急度判定规则(简单可解释):
- 含打断关键词(取消/停止/stop…)→
INTERRUPT(取消式处理) - 是一个提问(带问号或疑问词,如"现在几点了?")→
IMMEDIATE(立即回应,但不打断后台任务) - 其它补充性指令(如"用日语回复")→
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_KEY,demo.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 1Moonshot 默认走推理模型
kimi-k3(旧的kimi-k2-*-preview与moonshot-v1-*已过时/停用)。 推理模型要求temperature=1且max_tokens>=2048,demo.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 key(OPENAI_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