第 23 章
第 5 周:长任务跑到一半,它现在到底是什么状态?
给测试 Agent 的长任务写一个显式状态机:8 个状态、7 种事件、12 条合法转换,56 格逐格测非法转换;状态和事件日志在同一个 SQLite 事务里持久化,用版本号挡并发领取,用租约挡僵尸 worker。配套 77 个本地跑通的测试和 5 个变异,一段讲解视频,三个互动演示。
BugHunt-Bench 里,让测试 Agent 探索一遍 Conduit 并提交 Bug 报告,一个任务要跑十几分钟。这十几分钟里,用户可能点了取消,模型 API 可能回了 429 要重试,跑任务的 worker 进程可能被 OOM 杀掉。任务列表上那一行该显示什么?“运行中”、“已取消”,还是"已取消但其实还在跑"?这一章把任务的生命周期写成一个显式的状态机,再从测开的角度回答三个问题:非法转换怎么测全,状态怎么持久化才不会半写,并发和进程崩溃时状态怎么保持正确。
讲解视频
互动演示
三个演示。第一个可以对一个任务逐个触发事件,看状态图上走哪条边、数据库那一行和事件日志怎么变,也可以模拟 worker 被 kill -9;切到"朴素写法"能看到非法事件被直接接受。第二个是 56 格转换矩阵,点一格就触发对应的(状态, 事件),还能把 56 个用例分别跑在两种实现上。第三个单步演示三种并发交错:两个 worker 抢同一个任务、排队时取消后旧读数的 worker 开跑、僵尸 worker 醒来上报。页面底部有自动判分的练习。
为什么不用几个布尔字段
第一版代码常常是 is_running、is_done、is_cancelled 三个布尔字段,各处代码改自己关心的那个。三个布尔有 2³ = 8 种组合,is_running=True, is_cancelled=True 算什么?是"已请求取消、worker 还没停",还是"取消失败了"?没人写下来,每个调用方各自理解,bug 就出在理解不一致的地方。
状态机做的事情很简单:把"当前在哪个状态、收到什么事件、能去哪"写成一张表。表里有的才允许,其他组合都是非法转换,直接拒绝,而且不落库。这张表同时就是规格,测试可以逐格对照。
8 个状态、7 种事件、12 条合法转换
| 当前状态 | 事件 | 目标状态 | 说明 |
|---|---|---|---|
submitted | enqueue | queued | 进队列,等 worker 领取 |
submitted / queued / retry_wait | cancel | cancelled | 还没在跑,直接取消 |
queued | start | running | worker 领取,attempts + 1,拿到租约 |
running | succeed | succeeded | 终态 |
running | fail | retry_wait 或 failed | 可重试且 attempts < max_attempts 才进 retry_wait,否则终态失败 |
running | cancel | cancelling | 正在跑的任务不能瞬间停下,先标记"已请求取消" |
retry_wait | retry_due | queued | 退避时间到,回到队列 |
cancelling | cancel_ack / fail | cancelled | worker 确认停下;停下前失败了也不再重试 |
cancelling | succeed | succeeded | 取消信号到达前已经跑完 |
按单元格数:8 × 7 = 56 个(状态, 事件)组合,其中 12 个合法(running + fail 算一格,目标由守卫条件决定),2 个幂等 no-op(对 cancelling 和 cancelled 再发一次 cancel:不报错、不写库、不记事件),剩下 42 个非法。
状态机本身是一个纯函数,输入旧任务和事件,返回一个新的 Task,不修改原对象(frozen=True 的 dataclass),也不碰数据库:
TRANSITIONS = {
(State.SUBMITTED, Event.ENQUEUE): State.QUEUED,
(State.SUBMITTED, Event.CANCEL): State.CANCELLED,
(State.QUEUED, Event.START): State.RUNNING,
(State.QUEUED, Event.CANCEL): State.CANCELLED,
(State.RUNNING, Event.SUCCEED): State.SUCCEEDED,
(State.RUNNING, Event.FAIL): None,
(State.RUNNING, Event.CANCEL): State.CANCELLING,
(State.RETRY_WAIT, Event.RETRY_DUE): State.QUEUED,
(State.RETRY_WAIT, Event.CANCEL): State.CANCELLED,
(State.CANCELLING, Event.CANCEL_ACK): State.CANCELLED,
(State.CANCELLING, Event.SUCCEED): State.SUCCEEDED, # 取消信号到达前已经跑完:如实记成功
(State.CANCELLING, Event.FAIL): None,
}
# 重复取消是幂等的:不报错、不写库、不记事件
NO_OPS = frozenset({(State.CANCELLING, Event.CANCEL), (State.CANCELLED, Event.CANCEL)})
def apply(task: Task, event: Event, *, retryable: bool = True) -> Task:
key = (task.state, event)
if key in NO_OPS:
return task
if key not in TRANSITIONS:
raise IllegalTransition(task.state, event)
target = TRANSITIONS[key] or _after_fail(task, retryable)
attempts = task.attempts + 1 if event is Event.START else task.attempts
return replace(task, state=target, attempts=attempts, version=task.version + 1)
表里有三个值得拿出来讨论的设计决定,都是取舍,不是唯一答案:
- 取消是请求,不是命令。
running → cancelling → cancelled,中间要等 worker 在检查点看到取消请求、自己停下来。直接把状态改成 cancelled,worker 并不知道,还会继续产生副作用。 cancelling时收到succeed,记成功。worker 在看到取消请求之前已经把 Bug 报告提交出去了,副作用真实发生,状态要和现实一致。如果业务要求取消后的结果作废,那是另一套补偿逻辑(撤回报告),不该靠改状态名掩盖。- 重复取消是 no-op,取消已完成的任务是非法。客户端超时后重发取消请求很常见,不该报错;但"取消一个已经 succeeded 的任务"应该让调用方知道。
非法转换怎么测
配套代码在 week05_长任务可靠性/code/task_state_machine/,测试分两层:test_machine.py 测纯函数,test_store.py 测持久化。
56 格全覆盖。用 pytest.mark.parametrize 把每个组合都跑一遍:
# 独立写一份"规格",不从实现里的 TRANSITIONS 抄(否则测试和实现错得一模一样)
LEGAL = {
(S.SUBMITTED, E.ENQUEUE): S.QUEUED,
# ……共 12 条,running + fail 按 attempts=1、max=3 写成 retry_wait
}
NO_OP = {(S.CANCELLING, E.CANCEL), (S.CANCELLED, E.CANCEL)}
ALL_PAIRS = list(itertools.product(State, Event))
@pytest.mark.parametrize("state,event", ALL_PAIRS, ids=lambda x: x.value)
def test_every_state_event_pair(state, event):
task = at(state)
if (state, event) in LEGAL:
new = apply(task, event)
assert new.state is LEGAL[(state, event)]
assert new.version == task.version + 1
assert task.state is state, "原对象不能被修改"
elif (state, event) in NO_OP:
assert apply(task, event) is task
else:
with pytest.raises(IllegalTransition):
apply(task, event)
这里最关键的一点是期望表独立手写。如果测试直接 from machine import TRANSITIONS 再遍历它生成期望,实现里写错一格,测试会跟着一起错,永远是绿的。快照测试也有同样的问题:快照是从实现里录的,第一次录下来的错误会被当成"正确"。
其他几类测试:
- 守卫条件测边界:
max_attempts = 3时,attempts 为 1、2 的失败进 retry_wait,第 3 次失败进 failed;不可重试的错误第 1 次就进 failed。 - 终态吸收:对三个终态逐个发 7 种事件,要么抛
IllegalTransition,要么状态不变。 - 随机游走查不变量:固定种子
20261001,生成 2,000 条长度 40 的随机事件序列,断言终态不会再变、attempts 不超过 max_attempts、version 等于真实发生的状态变化次数。种子固定,失败时才能复现。 - 持久化层"拒绝即不落库":非法转换之后,tasks 表、租约字段、事件日志三样都要和之前完全一样(测试逐字段比较查询出的值)。
本地运行结果:
$ python3.13 -m pytest -q
........................................................................ [ 93%]
..... [100%]
77 passed in 1.82s
测试全绿不代表测试有用。mutation_check.py 往实现里故意埋 5 个 bug,在临时目录里各跑一遍测试(手工版的变异测试):
$ python3.13 mutation_check.py
抓到 终态不再吸收:cancelled 还能 enqueue <- test_every_state_event_pair, test_illegal_transition_leaves_db_untouched, test_random_walk_invariants, test_terminal_states_are_absorbing
抓到 重试上限差一:attempts <= max_attempts <- test_fail_guard, test_random_walk_invariants, test_reclaim_respects_max_attempts_and_cancel, test_start_counts_attempts
抓到 去掉乐观锁的版本检查 <- test_claim_race_with_threads, test_two_workers_claim_same_task
抓到 去掉租约持有者检查 <- test_zombie_worker_cannot_overwrite
抓到 去掉所有显式事务(BEGIN IMMEDIATE / COMMIT / ROLLBACK) <- test_crash_between_update_and_log_rolls_back
5/5 个变异被测试抓到
互动演示 2 里还可以把同一套 56 格用例跑在"朴素写法"上(每个事件直接改成固定目标,不查表;fail 按次数决定进 retry_wait 还是 failed):只通过 10 格。它连两条合法转换都做错了(running + cancel 直接变 cancelled,cancelling + fail 又进了重试),2 个 no-op 都多写了一次库,42 个非法组合全部被接受。
状态持久化:同一个事务、版本号、租约
状态只放内存里,进程一重启就丢了。这里用 SQLite(Python 标准库 sqlite3,本机 SQLite 3.50.4)。计划里提过 Redis 队列,本章先用 SQLite:单文件、有事务、不用起服务,每个测试用例一个临时库,跑得快也互不干扰。两张表:tasks 存当前状态,task_events 存每一次转换,(task_id, version) 唯一。state 列加了 CHECK 约束,拼错的状态名在数据库层就会被拒绝。
一次状态变更的核心代码:
def _fire_locked(self, task_id, event, retryable, actor, expected_version, check_lease):
current = self.get(task_id)
if expected_version is not None and current.version != expected_version:
raise VersionConflict(f"{task_id}: expected v{expected_version}, found v{current.version}")
new = apply(current, event, retryable=retryable) # 先判合法性,非法直接抛 IllegalTransition
if check_lease and event in WORKER_EVENTS:
owner = self.lease(task_id)[0]
if owner != actor:
raise StaleLease(f"{task_id}: lease owner is {owner!r}, not {actor!r}")
if new is current: # 幂等 no-op:不写库、不记事件
return current
self._write(current, new, event, retryable, actor)
return new
_write 里先 UPDATE tasks ... WHERE id = ? AND version = ?,检查 rowcount 是不是 1,再 INSERT INTO task_events。整个过程包在 BEGIN IMMEDIATE 和 COMMIT / ROLLBACK 之间。逐条说为什么:
1. 状态和事件日志同一个事务写。 连接用 isolation_level=None 打开。Python 文档的说法是:设为 None 时 sqlite3 不会隐式开启事务,交给用户用显式 SQL 自己管理。如果先 UPDATE 成功、写日志时进程崩了,tasks 表说 queued,日志里却没有这条转换,重放就对不上。测试用 monkeypatch 让写日志时抛 RuntimeError("disk full"),断言 tasks 表的 UPDATE 也一起回滚了。
2. 检查和写入必须是一个原子操作。 “先 SELECT 出 queued,在代码里判断,再 UPDATE"是两步,中间别的连接可以插进来。SQLite 文档说明:默认的 DEFERRED 事务先以读事务开始,之后的写语句如果无法升级(别的连接已经在写)就返回 SQLITE_BUSY;IMMEDIATE 则在 BEGIN 时就启动写事务。另一个连接已经在写时,BEGIN IMMEDIATE 会返回 SQLITE_BUSY;Python sqlite3.connect 的 timeout 参数(默认 5 秒)是遇到锁时等待多久才抛 OperationalError。所以在 BEGIN IMMEDIATE 里做"读 → 判断 → 写”,判断用的就是最新数据。
3. 版本号(乐观锁)。 worker 轮询队列时读到的 version 可能已经过期,领取时带上 expected_version,对不上就抛 VersionConflict,什么都不写。UPDATE 带 WHERE version = ? 再看 rowcount(Python 文档:对 UPDATE 返回被修改的行数)是第二道保险。说实话,在单个 SQLite 文件上,BEGIN IMMEDIATE 已经把写串行化,这个 WHERE 条件不会被触发(把它改成恒真,77 个测试照样全过)。但只要读和写之间没有一把写锁,它就是唯一的防线:比如不开显式事务、PostgreSQL / MySQL 在默认隔离级别下先读后写,或者 Redis。Redis 文档说明 WATCH 为事务提供检查再写入(CAS)语义,被 WATCH 的 key 在 EXEC 之前被改过,整个事务就放弃。
4. 租约。 start 时记下 lease_owner 和 lease_until(演示里租约 30 秒),worker 定期 heartbeat() 续约。只有租约持有者能上报 succeed / fail / cancel_ack。本实现的过期是惰性的:heartbeat() 和上报只核对 lease_owner,不比较 lease_until 和当前时间,所以巡检回收之前,原持有者的心跳和上报仍然有效。这是设计选择:租约到期但还没被回收的 worker 其实没有被别人取代,接受它的结果不会造成冲突。
5. 重启后回收。 worker 被 kill -9 后,库里它的任务还是 running,不会再有人改它。巡检 reclaim_expired() 在一个事务里查出租约过期的 running / cancelling 任务,按一次可重试失败走状态机:attempts 还够就进 retry_wait,用完了进 failed,原来在 cancelling 的进 cancelled。“启动时把所有 running 改回 queued"是常见的错误做法:分不清 worker 是死了还是慢,原 worker 还活着时会双跑;不计 attempts,一个必崩的任务会无限重跑。
6. 事件日志可以重放。 从 submit 开始把日志里的事件逐条 apply,结果必须和 tasks 表一致。test_replay_matches_table_after_random_ops 对 200 个随机任务各触发 25 个随机事件后做这个检查;test_state_survives_restart 关掉连接再重新打开,断言状态、attempts、version 都还在。
三个并发交错,确定性地复现
并发 bug 靠多线程碰运气很难稳定复现。race_demo.py 用两个连接手工安排交错顺序,每次运行结果都一样(test_store.py 里另有一个 8 线程 × 50 轮的压力测试,断言每轮恰好一个 worker 领到)。实际输出:
$ python3.13 race_demo.py
场景 1:两个 worker 同时轮询到同一个 queued 任务
朴素写法:2 个 worker 都以为自己领到了,最终 owner = worker-B
状态机 worker-A: 写入成功
状态机 worker-B: 拒绝(VersionConflict)
状态机:1 个 worker 领到
场景 2:用户在排队时取消了任务,一个早先读到 queued 的 worker 还是开跑了
朴素写法:最终状态 = running(取消被覆盖)
状态机 worker-A: 拒绝(VersionConflict)
状态机:最终状态 = cancelled
场景 3:worker-A 卡住超过 30 秒,任务被收回并交给 worker-B,之后 A 醒来上报成功
朴素写法:两次 UPDATE 都写入成功,状态被写了两次 succeeded,最终 owner = worker-B
状态机 + 版本号,但不查租约持有者:
租约过期,回收:['bug-9']
worker-A 上报 succeed: 写入成功
worker-B 上报 succeed: 拒绝(IllegalTransition)
最终:succeeded,记下的是 worker-A 的结果
状态机 + 版本号 + 租约检查:
租约过期,回收:['bug-9']
worker-A 上报 succeed: 拒绝(StaleLease)
worker-B 上报 succeed: 写入成功
最终:succeeded,attempts = 2,version = 6
(0, 'submit', 1, None, 'submitted', 'client')
(1, 'enqueue', 1, 'submitted', 'queued', 'client')
(2, 'start', 1, 'queued', 'running', 'worker-A')
(3, 'fail', 1, 'running', 'retry_wait', 'reaper')
(4, 'retry_due', 1, 'retry_wait', 'queued', 'client')
(5, 'start', 1, 'queued', 'running', 'worker-B')
(6, 'succeed', 1, 'running', 'succeeded', 'worker-B')
场景 1、2 里,即使去掉版本号检查,状态机也会拦住:B 领取时库里已经是 running,running + start 是非法转换;场景 2 的 cancelled + start 同理。这正是"在同一个写事务里检查"的效果。版本号让错误更明确(“你读到的已经过期”),变异测试里去掉它,测试也确实会失败,只是失败原因变成了意外的 IllegalTransition。
场景 3 是长任务特有的僵尸 worker:worker-A 卡了 31 秒(GC 停顿、网络分区、被调试器挂起都可能),租约过期,任务被回收、重新排队、交给 worker-B。A 醒来后上报 succeed。输出里对照了三种实现。朴素写法什么都不查,两次 UPDATE 都写进去,谁的结果算数取决于谁最后写。只去掉租约检查、保留状态机时,A 的迟到上报会把任务记成 succeeded,B 跑完再报反而被当成非法转换(succeeded + succeed)拒绝:最终记下的是过期 worker 的结果。检查租约后,状态层只认 B。
注意租约只保护状态这一行:A 卡住之前如果已经把 Bug 报告提交出去,这份报告照样存在,B 跑完还会再提交一份。要防止外部副作用重复,得让副作用的接收方校验 fencing token 或幂等键,下一章讲。
对照:Codex 和 Celery 的状态
| 系统 | 状态 | 和本章的对应 |
|---|---|---|
| Codex app-server(turn) | Completed / Interrupted / Failed / InProgress | 只有 4 个,枚举里没有排队、重试等待这类中间状态;Interrupted 大致对应我们的 cancelled |
| Codex app-server(命令执行 item) | InProgress / Completed / Failed / Declined | 多一个 Declined:主要是审批被拒(源码注释承认部分运行前的拒绝也会报成 Declined),共同点是命令根本没跑。“没跑"和"跑了失败"是两种状态 |
| Celery | PENDING / STARTED / RETRY / FAILURE / SUCCESS / REVOKED | PENDING 不是记录下来的状态,而是任何未知 task id 的默认状态;STARTED 默认不记录,要开 task_track_started,或给单个任务设 track_started=True;REVOKED 表示任务被撤销 |
Codex 的两个枚举见
v2/turn.rs:33-38
和
v2/item.rs:1077-1082
(commit 7993248);Declined 的范围见
core/src/tools/events.rs:453-456
的注释:ToolError::Rejected 既用于用户拒绝审批,也用于部分运行时拒绝(比如 setup 失败),所以一部分非用户原因的失败也会报成 Declined。turn / item 的层级在「一个生产级 Agent 长什么样?Codex 的仓库地图」那一章讲过。Celery 的说法出自官方文档
Next Steps
、
celery.result
和 Tasks 用户指南的
States
一节。
Celery 的 PENDING 是个经典测试坑:task id 拼错了,查到的也是 PENDING,看上去像"还在排队”。本章的实现查不存在的 id 直接抛 KeyError,“不存在"和"在排队"是两回事。
面试怎么讲
一句话版本:长任务用显式状态机管理,8 个状态、7 种事件、12 条合法转换写成一张表,表外的组合一律拒绝且不落库;状态和事件日志在同一个事务里写,领取任务带版本号,运行中靠租约证明自己还活着,进程崩溃后由巡检把过期租约按可重试失败回收。测试按 56 格逐格对照独立写的规格,再用随机游走查不变量,用变异测试证明这些测试真能抓 bug。
可能的追问:
- 为什么要事件日志,tasks 表不够吗? tasks 表只有"现在”。排查"它是怎么变成 failed 的"要靠日志;日志能重放,“重放结果等于 tasks 表"是一条很好的不变量。
- 租约设多长? 要比心跳间隔长几倍,留出网络抖动和停顿的余量。太短会把慢 worker 误判成死了(回收之后它的迟到上报会被拒绝:任务还在 retry_wait / queued 时是非法转换,已经被别的 worker 领走时是 StaleLease;但这次运行白跑了,外部副作用也可能已经发生),太长则崩溃后要等很久才回收。具体倍数要看心跳的实测延迟,没有通用答案。
- 取消一个正在跑的任务,worker 怎么知道? worker 在检查点(比如每轮 agent loop 开始前)读一下任务状态,看到 cancelling 就停下并上报 cancel_ack。检查点之间的那段时间取消不会生效,所以 cancelling 是一个真实存在、可能持续一段时间的状态。
常见错误说法
- “状态就是几个布尔字段”:组合会爆炸,每种组合的含义没人写下来。
- “非法转换测几个典型的就够了”:8 × 7 的表只有 56 格,全测成本很低,漏掉的往往是没想到的那几格。
- “测试的期望值直接从实现的转换表里读”:实现错了,测试跟着错,永远是绿的。
- “取消就是把状态改成 cancelled”:正在跑的 worker 不知道,要有 cancelling 和 worker 确认。
- “进程重启后把 running 全部改回 queued”:分不清死和慢,会双跑,也不计重试次数。
- “有了状态校验就不需要事务”:校验和写入不在一个原子操作里,校验时读到的可能已经是旧数据。
下一章:重试、指数退避和幂等键。retry_wait 要等多久,重试为什么会导致重复副作用,怎么测。