跳转到内容

状态、事件与复杂度

一个异步任务表最初可能只有三列:idstatusresult。创建时写 pending,Worker 开始时写 running,结束时写 completed。后来加入取消、超时、重试、人工审核和恢复,代码逐渐出现下面的写法:

await db.task.update({ id, status: 'running' });
// 另一个地方
await db.task.update({ id, status: 'cancelled' });
// 失败处理器
await db.task.update({ id, status: 'pending' });
// 管理命令
await db.task.update({ id, status: requestedStatus });

每一行单独看都合理,组合起来却没有一个位置回答:已取消任务还能不能完成?两个 Worker 同时启动时谁获胜?completed 后收到迟到的超时事件怎么办?重试是新任务还是同一任务的新尝试?恢复进程如何区分重复事件与缺失事件?

这不是“状态字段太多”的问题,而是写入没有协议。状态机(State Machine)状态机State Machine用有限状态和明确转换规则描述系统如何变化。打开术语条目 → 把所有合法变化集中到转换函数:当前状态与输入事件决定下一状态,非法组合返回明确拒绝。不变量(Invariant)不变量Invariant在所有合法状态与状态转换中都必须成立、可以被程序检查的条件。打开术语条目 → 则定义任何路径都不能破坏的业务事实。再配合 事件标识(Event ID)事件标识Event ID标识事件在一个有序事件流中的位置,用于断线续传、去重和回放。打开术语条目 →、版本检查、持久化事件和 快照(Snapshot)快照Snapshot某个已知版本上的完整状态副本,用于缩短恢复路径,但不能替代增量事件和校验。打开术语条目 →,任务生命周期才能测试、回放和恢复。

本课构建一个可测试任务状态机。它不把 Event Sourcing 当成所有系统的默认架构,而是借助事件记录解释状态变化,并明确什么时候只保存当前状态更简单。

1. 具体工程问题:迟到的完成覆盖了取消

Section titled “1. 具体工程问题:迟到的完成覆盖了取消”

任务 task-42 的正常路径:

queued → running → completed

用户在运行中取消:

queued → running → cancelling → cancelled

真实系统里,取消不是瞬时的。控制服务先记录取消意图,向 Worker 发送信号;Worker 可能在信号到达前已经完成外部工作。于是两个请求并发发生:

请求 A:CancelTask(expectedVersion=8)
请求 B:CompleteTask(expectedVersion=8, artifactHash=...)

如果两个处理器只是执行 UPDATE tasks SET status = ? WHERE id = ?,后提交者覆盖前提交者。最终状态取决于时序,而不是业务规则。更糟的是,completedcancelled 都可能缺少与之匹配的副作用证据。

我们先写出必须保持的事实:

  1. completedfailedcancelled 是终态,不能被普通命令重新打开;
  2. completed 必须携带经过验证的 Artifact 身份;
  3. cancelled 表示执行已停止或已通过对账确认不会再产生新副作用;只记录取消请求时应处于 cancelling
  4. 每个已提交事件都有唯一 Event ID;同一事件重放不能重复改变状态;
  5. 命令必须声明它读取的 expectedVersion;过期命令不能静默覆盖新状态;
  6. 状态版本每接受一个新事件恰好增加一;拒绝和重复事件不增加版本。

状态列表应从这些不变量推导,而不是先画一堆方框再找含义。

命令、事件、状态与不变量

转换函数是唯一写入口
COMMANDStartTask表达意图;可能被拒绝
GUARDstate = queued版本、权限、前置条件
EVENTTaskStarted已经提交的事实
STATErunning由事件折叠得到
INVARIANT终态不可重新启动`completed + StartTask` 必须返回可分类拒绝,而不是偷偷改回 running。
CONCURRENCYexpectedVersion = 7只有当前版本仍为 7 时才能提交事件;否则重新读取并重算命令。
EVIDENCEeventId + previous + next记录转换前后状态、命令身份、事件身份与拒绝原因。
命令不是事实,事件不是任意消息,状态也不是可随处修改的字段。把写入集中到转换函数,才能测试非法路径和并发冲突。

2. 命令、事件、状态和观察不是同一种数据

Section titled “2. 命令、事件、状态和观察不是同一种数据”

2.1 Command 表达意图,可能被拒绝

Section titled “2.1 Command 表达意图,可能被拒绝”

StartTaskRequestCancellationCompleteTask 是命令。它们来自用户、调度器、Worker 或恢复器,表达“希望系统发生什么”。命令可能因为权限、当前状态、版本冲突、参数无效或预算耗尽而被拒绝。

命令命名用祈使语义:

const command = { type: 'StartTask', taskId: 'task-42', expectedVersion: 3 } as const;

不要把 TaskStarted 作为输入命令名称;这会把“请求开始”和“已经开始”混在一起。

事件在事务提交后才成立:

const event = {
eventId: 'evt-0004',
type: 'TaskStarted',
taskId: 'task-42',
sequence: 4,
attempt: 1,
} as const;

事件使用过去式。它不保证外部世界所有事情都成功,只保证系统定义的事实已经被权威存储接受。例如 CancellationRequested 只证明取消意图已登记,不证明 Worker 已停止;真正停止后再产生 TaskCancelled

2.3 State 是事件按顺序折叠后的当前结论

Section titled “2.3 State 是事件按顺序折叠后的当前结论”

状态可以被定义为:

Sₙ = fold(S₀, [E₁, E₂, …, Eₙ])

其中 S0 是初始状态,E_i 是按任务序号排序的已接受事件。转换函数:

δ(Sᵢ₋₁, Eᵢ) → Sᵢ 或 Reject(reason)

这里的 Reject 不能被忽略。若事件存储里出现无法应用的事件,说明历史损坏、版本不兼容或读取顺序错误;恢复器应停止并报告,而不是跳过后继续生成一个“尽量合理”的状态。

2.4 Observation 是可能过期的读模型

Section titled “2.4 Observation 是可能过期的读模型”

列表页显示的进度、队列长度、预计完成时间通常是派生观察。它们可以来自缓存或异步投影,允许短暂滞后。命令处理器不能拿一个可能过期的列表页状态直接决定写入;它要读取权威版本,或在数据库条件更新中验证版本。

数据时间语义能否被拒绝是否可重建典型持久化
Command希望发生可以不一定请求日志、Inbox
Event已经提交不再撤回,只能追加后续事实是历史本身追加事件表
State截至当前的结论不适用可由事件重建,或直接存储状态行、Snapshot
Observation面向查询的视图不适用可由权威状态重建缓存、投影表

本课任务状态:

export type TaskStatus =
| 'queued'
| 'running'
| 'cancelling'
| 'completed'
| 'failed'
| 'cancelled';

为什么需要 cancelling?因为“收到取消请求”和“执行已停止”之间存在时间窗口。直接从 running 改成 cancelled 会制造错误事实:外部工具可能仍在写文件或发请求。

为什么没有 retrying?重试在本设计中是同一任务的下一次 attempt,它先通过 TaskFailed(retryable=true) 回到 queued,再接受新的 TaskStarted。如果产品需要展示退避等待,可以引入 waiting_retry 并保存 notBefore;状态应服务于可执行规则,而不是为了显示更丰富。

为什么 failed 是终态?本课把“可重试失败”建模为 AttemptFailed 事件,并回到 queuedTaskFailed 表示任务预算耗尽或错误不可恢复。使用不同事件避免“failed 有时终态、有时不是”的歧义。

当前状态输入事件Guard下一状态必须附带的证据
queuedTaskStartedattempt 等于上次 + 1runningworkerId、leaseUntil、attempt
queuedCancellationRequested无执行中副作用cancelledrequester、reason
runningCancellationRequested请求身份有效cancellingrequester、reason、signalId
runningTaskCompletedArtifact 已验证;attempt 匹配completedartifactHash、bytes、attempt
runningAttemptFailedretryable 且剩余预算 > 0queuederrorCode、attempt、nextAttempt
runningTaskFailed不可重试或预算耗尽failederrorCode、attempt
cancellingTaskCancelledWorker 已停止或完成对账cancelledstoppedAtSequence、effectCheck
cancellingTaskCompleted由业务策略决定completed 或拒绝Artifact、取消竞争裁决
任一终态任一普通事件拒绝TERMINAL_STATE

cancelling + TaskCompleted 没有通用答案。如果完成的副作用不可撤回且 Artifact 已经产生,系统可能接受完成并记录取消过晚;如果取消承诺比完成优先,系统则拒绝完成并进入对账。规则必须在状态机中明确,不能留给数据库提交顺序决定。本课选择:只要 TaskCompletedattempt 匹配,并且 Artifact 已验证,接受 completed,同时在 Trace 中记录 completionWonCancellationRace=true

4. 状态转换函数:保持纯函数边界

Section titled “4. 状态转换函数:保持纯函数边界”

核心转换函数不访问数据库、不读取当前时间、不调用工具。输入是完整当前状态和一个候选事件,输出是新状态或结构化拒绝:

export type TransitionResult =
| { accepted: true; next: TaskState; evidence: TransitionEvidence }
| { accepted: false; reason: RejectReason; current: TaskState };

纯函数有三个直接收益:

  • 可以枚举所有状态 × 事件组合;
  • 可以对同一历史稳定回放;
  • 数据库并发、消息去重和业务规则能分开测试。

副作用流程是“先判断可接受,再原子持久化事件和新状态”,而不是在纯函数里直接写库。

初始状态:

status=queued, version=0, attempt=0

按序事件:

E1 TaskStarted(attempt=1)
E2 CancellationRequested(signalId=sig-9)
E3 TaskCompleted(attempt=1, artifactHash=h1)
E4 TaskCancelled(stoppedAtSequence=3)

逐步折叠:

步骤输入状态事件结果
0queued/v0/a0初始状态
1queued/v0/a0TaskStarted(a1)running/v1/a1
2running/v1/a1CancellationRequestedcancelling/v2/a1
3cancelling/v2/a1TaskCompleted(a1,h1)completed/v3/a1,完成赢得竞争
4completed/v3/a1TaskCancelled拒绝:TERMINAL_STATE,版本保持 3

注意 E4 即使有更大的消息到达时间,也不能覆盖 E3。若 E4 已经被错误写入事件存储,重放器应报告非法历史;若它只是队列中的候选消息,Inbox 处理器记录拒绝后确认消息即可。

再看重复事件:E3 因 Worker 没收到确认而重投,Event ID 仍是同一个。处理器先查 Inbox,发现已处理,返回第一次结果的摘要;不再次调用转换函数,不增加版本。若 Worker 使用新 Event ID 重发同一业务事实,系统还需要业务去重键,例如 (taskId, type, attempt) 唯一约束。Event ID 去重只防止同一消息重复,不能防止发送方制造语义重复。

6. 乐观并发:过期写入必须失败

Section titled “6. 乐观并发:过期写入必须失败”

乐观并发控制(Optimistic Concurrency Control)乐观并发控制Optimistic Concurrency Control提交写入时检查读取版本是否仍然有效,冲突时拒绝并要求重新计算。打开术语条目 → 让命令声明它基于哪个版本做决定。数据库更新可以写成:

UPDATE task_state
SET status = $1,
version = version + 1,
state_json = $2
WHERE task_id = $3
AND version = $4;

受影响行数为 0 时,不代表数据库坏了;它表示状态已被其他事务推进。处理器应重新读取,然后重新评估命令。不能直接把 expectedVersion 改成最新值再强行写入,因为新状态下原命令可能已非法。

以取消与完成竞争为例:

A读取 running/v8,准备 CancellationRequested
B读取 running/v8,准备 TaskCompleted
A先提交:cancelling/v9
B条件更新 version=8,影响0行
B重新读取 cancelling/v9
B按明确策略重新计算,接受 TaskCompleted → completed/v10

最终结果由 cancelling + TaskCompleted 的规则决定,而不是“最后提交者覆盖”。

6.1 为什么事务隔离级别仍然重要

Section titled “6.1 为什么事务隔离级别仍然重要”

版本条件更新避免丢失更新,但命令可能还读取其他行,例如剩余配额、租约、Artifact 记录。若业务不变量跨多行,必须让检查与写入处于合适事务边界,并处理序列化冲突或唯一约束冲突。不能因为用了 version 就声称整个业务操作线性一致。

本课来源中的 PostgreSQL 隔离文档用于核对各隔离级别的数据库语义;实际实现要针对所用数据库验证,不应把示例 SQL 直接外推到不同存储。

分布式系统常同时出现两种顺序:

  1. 任务内权威序号:事件被权威存储接受时分配的 sequence;用于重放;
  2. 消息到达顺序:消费者看到消息的先后;可能重复、延迟或乱序。

如果消息 sequence=12 先于 sequence=11 到达,消费者有三种策略:

策略适用条件风险
拒绝并稍后重试上游保证缺失消息最终会到长时间阻塞、重复重试
暂存 Gap Buffer乱序窗口有明确上限Buffer 膨胀,需要超时和补拉
从权威存储按序补拉有按任务查询历史的接口增加读取,需防止权限与版本错误

本课处理器遇到 sequence > current.version + 1 时返回 SEQUENCE_GAP,由恢复器从事件表补拉。遇到 sequence <= current.version 时不能直接判定重复,要先查 Event ID:

  • Event ID 已处理:幂等返回;
  • Event ID 未处理但序号已占用:历史冲突,进入隔离;
  • 同一 Event ID 对应不同 Payload Hash:协议破坏,不能选择其中一个继续。

下面是可运行的核心实现,不含具体数据库适配器。存储接口、内存实现和断言测试都给出,输入输出明确,不执行任意代码。

examples/task-state-machine.ts
import assert from 'node:assert/strict';
import { createHash } from 'node:crypto';
export type TaskStatus =
| 'queued'
| 'running'
| 'cancelling'
| 'completed'
| 'failed'
| 'cancelled';
export interface TaskState {
taskId: string;
status: TaskStatus;
version: number;
attempt: number;
workerId: string | null;
leaseUntilEpochMs: number | null;
cancellationSignalId: string | null;
artifactSha256: string | null;
terminalErrorCode: string | null;
}
interface EventBase {
eventId: string;
taskId: string;
sequence: number;
}
export type TaskEvent =
| (EventBase & {
type: 'TaskStarted';
attempt: number;
workerId: string;
leaseUntilEpochMs: number;
})
| (EventBase & {
type: 'CancellationRequested';
signalId: string;
requester: 'user' | 'system';
reason: string;
})
| (EventBase & {
type: 'TaskCompleted';
attempt: number;
artifactSha256: string;
artifactVerified: true;
})
| (EventBase & {
type: 'AttemptFailed';
attempt: number;
errorCode: string;
retryable: true;
nextAttempt: number;
})
| (EventBase & {
type: 'TaskFailed';
attempt: number;
errorCode: string;
})
| (EventBase & {
type: 'TaskCancelled';
stoppedAtSequence: number;
effectsReconciled: true;
});
export type RejectReason =
| 'TASK_ID_MISMATCH'
| 'DUPLICATE_OR_OLD_SEQUENCE'
| 'SEQUENCE_GAP'
| 'TERMINAL_STATE'
| 'ILLEGAL_TRANSITION'
| 'ATTEMPT_MISMATCH'
| 'INVALID_NEXT_ATTEMPT'
| 'ARTIFACT_NOT_VERIFIED'
| 'EFFECTS_NOT_RECONCILED';
export type TransitionResult =
| {
accepted: true;
next: TaskState;
evidence: {
previousStatus: TaskStatus;
nextStatus: TaskStatus;
eventId: string;
completionWonCancellationRace: boolean;
};
}
| { accepted: false; reason: RejectReason; current: TaskState };
export function initialState(taskId: string): TaskState {
return {
taskId,
status: 'queued',
version: 0,
attempt: 0,
workerId: null,
leaseUntilEpochMs: null,
cancellationSignalId: null,
artifactSha256: null,
terminalErrorCode: null,
};
}
function reject(current: TaskState, reason: RejectReason): TransitionResult {
return { accepted: false, reason, current };
}
export function transition(current: TaskState, event: TaskEvent): TransitionResult {
if (event.taskId !== current.taskId) return reject(current, 'TASK_ID_MISMATCH');
if (event.sequence <= current.version) return reject(current, 'DUPLICATE_OR_OLD_SEQUENCE');
if (event.sequence !== current.version + 1) return reject(current, 'SEQUENCE_GAP');
if (['completed', 'failed', 'cancelled'].includes(current.status)) {
return reject(current, 'TERMINAL_STATE');
}
const base = { ...current, version: event.sequence };
let next: TaskState | null = null;
let completionWonCancellationRace = false;
switch (event.type) {
case 'TaskStarted': {
if (current.status !== 'queued') return reject(current, 'ILLEGAL_TRANSITION');
if (event.attempt !== current.attempt + 1) return reject(current, 'ATTEMPT_MISMATCH');
next = {
...base,
status: 'running',
attempt: event.attempt,
workerId: event.workerId,
leaseUntilEpochMs: event.leaseUntilEpochMs,
cancellationSignalId: null,
};
break;
}
case 'CancellationRequested': {
if (current.status === 'queued') {
next = {
...base,
status: 'cancelled',
cancellationSignalId: event.signalId,
workerId: null,
leaseUntilEpochMs: null,
};
} else if (current.status === 'running') {
next = { ...base, status: 'cancelling', cancellationSignalId: event.signalId };
} else {
return reject(current, 'ILLEGAL_TRANSITION');
}
break;
}
case 'TaskCompleted': {
if (current.status !== 'running' && current.status !== 'cancelling') {
return reject(current, 'ILLEGAL_TRANSITION');
}
if (event.attempt !== current.attempt) return reject(current, 'ATTEMPT_MISMATCH');
if (event.artifactVerified !== true) return reject(current, 'ARTIFACT_NOT_VERIFIED');
completionWonCancellationRace = current.status === 'cancelling';
next = {
...base,
status: 'completed',
artifactSha256: event.artifactSha256,
workerId: null,
leaseUntilEpochMs: null,
};
break;
}
case 'AttemptFailed': {
if (current.status !== 'running') return reject(current, 'ILLEGAL_TRANSITION');
if (event.attempt !== current.attempt) return reject(current, 'ATTEMPT_MISMATCH');
if (event.nextAttempt !== current.attempt + 1) return reject(current, 'INVALID_NEXT_ATTEMPT');
next = {
...base,
status: 'queued',
workerId: null,
leaseUntilEpochMs: null,
terminalErrorCode: null,
};
break;
}
case 'TaskFailed': {
if (current.status !== 'running' && current.status !== 'cancelling') {
return reject(current, 'ILLEGAL_TRANSITION');
}
if (event.attempt !== current.attempt) return reject(current, 'ATTEMPT_MISMATCH');
next = {
...base,
status: 'failed',
workerId: null,
leaseUntilEpochMs: null,
terminalErrorCode: event.errorCode,
};
break;
}
case 'TaskCancelled': {
if (current.status !== 'cancelling') return reject(current, 'ILLEGAL_TRANSITION');
if (event.effectsReconciled !== true) return reject(current, 'EFFECTS_NOT_RECONCILED');
next = {
...base,
status: 'cancelled',
workerId: null,
leaseUntilEpochMs: null,
};
break;
}
}
return {
accepted: true,
next,
evidence: {
previousStatus: current.status,
nextStatus: next.status,
eventId: event.eventId,
completionWonCancellationRace,
},
};
}
function canonicalEventHash(event: TaskEvent): string {
const keys = Object.keys(event).sort();
const normalized = Object.fromEntries(keys.map((key) => [key, event[key as keyof TaskEvent]]));
return createHash('sha256').update(JSON.stringify(normalized)).digest('hex');
}
interface StoredEvent {
event: TaskEvent;
payloadSha256: string;
}
export interface TaskStore {
load(taskId: string): Promise<{ state: TaskState; events: StoredEvent[] }>;
findProcessedEvent(eventId: string): Promise<StoredEvent | null>;
appendAtomically(input: {
expectedVersion: number;
event: TaskEvent;
payloadSha256: string;
nextState: TaskState;
}): Promise<'committed' | 'version-conflict' | 'event-id-conflict'>;
}
export type ApplyResult =
| { status: 'committed'; state: TaskState }
| { status: 'duplicate'; state: TaskState }
| { status: 'rejected'; reason: RejectReason; state: TaskState }
| { status: 'conflict'; code: 'VERSION_CONFLICT' | 'EVENT_ID_PAYLOAD_CONFLICT'; state: TaskState };
export async function applyEvent(store: TaskStore, event: TaskEvent): Promise<ApplyResult> {
const payloadSha256 = canonicalEventHash(event);
const processed = await store.findProcessedEvent(event.eventId);
if (processed) {
const current = await store.load(event.taskId);
return processed.payloadSha256 === payloadSha256
? { status: 'duplicate', state: current.state }
: { status: 'conflict', code: 'EVENT_ID_PAYLOAD_CONFLICT', state: current.state };
}
const current = await store.load(event.taskId);
const result = transition(current.state, event);
if (!result.accepted) {
return { status: 'rejected', reason: result.reason, state: current.state };
}
const committed = await store.appendAtomically({
expectedVersion: current.state.version,
event,
payloadSha256,
nextState: result.next,
});
if (committed === 'version-conflict') {
return { status: 'conflict', code: 'VERSION_CONFLICT', state: (await store.load(event.taskId)).state };
}
if (committed === 'event-id-conflict') {
return { status: 'conflict', code: 'EVENT_ID_PAYLOAD_CONFLICT', state: (await store.load(event.taskId)).state };
}
return { status: 'committed', state: result.next };
}
export class InMemoryTaskStore implements TaskStore {
private readonly states = new Map<string, TaskState>();
private readonly eventsByTask = new Map<string, StoredEvent[]>();
private readonly eventsById = new Map<string, StoredEvent>();
constructor(taskIds: string[]) {
for (const taskId of taskIds) this.states.set(taskId, initialState(taskId));
}
async load(taskId: string): Promise<{ state: TaskState; events: StoredEvent[] }> {
const state = this.states.get(taskId);
if (!state) throw new Error(`unknown task: ${taskId}`);
return {
state: structuredClone(state),
events: structuredClone(this.eventsByTask.get(taskId) ?? []),
};
}
async findProcessedEvent(eventId: string): Promise<StoredEvent | null> {
return structuredClone(this.eventsById.get(eventId) ?? null);
}
async appendAtomically(input: {
expectedVersion: number;
event: TaskEvent;
payloadSha256: string;
nextState: TaskState;
}): Promise<'committed' | 'version-conflict' | 'event-id-conflict'> {
const current = this.states.get(input.event.taskId);
if (!current || current.version !== input.expectedVersion) return 'version-conflict';
const existing = this.eventsById.get(input.event.eventId);
if (existing) return 'event-id-conflict';
const stored = { event: structuredClone(input.event), payloadSha256: input.payloadSha256 };
this.eventsById.set(input.event.eventId, stored);
this.eventsByTask.set(input.event.taskId, [...(this.eventsByTask.get(input.event.taskId) ?? []), stored]);
this.states.set(input.event.taskId, structuredClone(input.nextState));
return 'committed';
}
}
export function replay(taskId: string, events: TaskEvent[]): TaskState {
let state = initialState(taskId);
const seen = new Set<string>();
for (const event of [...events].sort((a, b) => a.sequence - b.sequence)) {
if (seen.has(event.eventId)) continue;
seen.add(event.eventId);
const result = transition(state, event);
if (!result.accepted) {
throw new Error(`invalid history at ${event.eventId}: ${result.reason}`);
}
state = result.next;
}
return state;
}
async function selfTest(): Promise<void> {
const store = new InMemoryTaskStore(['task-42']);
const started: TaskEvent = {
eventId: 'evt-1', taskId: 'task-42', sequence: 1, type: 'TaskStarted',
attempt: 1, workerId: 'worker-a', leaseUntilEpochMs: 1000,
};
assert.equal((await applyEvent(store, started)).status, 'committed');
assert.equal((await applyEvent(store, started)).status, 'duplicate');
const cancel: TaskEvent = {
eventId: 'evt-2', taskId: 'task-42', sequence: 2, type: 'CancellationRequested',
signalId: 'sig-9', requester: 'user', reason: 'no longer needed',
};
assert.equal((await applyEvent(store, cancel)).state.status, 'cancelling');
const completed: TaskEvent = {
eventId: 'evt-3', taskId: 'task-42', sequence: 3, type: 'TaskCompleted',
attempt: 1, artifactSha256: 'a'.repeat(64), artifactVerified: true,
};
assert.equal((await applyEvent(store, completed)).state.status, 'completed');
const lateCancel: TaskEvent = {
eventId: 'evt-4', taskId: 'task-42', sequence: 4, type: 'TaskCancelled',
stoppedAtSequence: 3, effectsReconciled: true,
};
const rejected = await applyEvent(store, lateCancel);
assert.equal(rejected.status, 'rejected');
if (rejected.status === 'rejected') assert.equal(rejected.reason, 'TERMINAL_STATE');
const loaded = await store.load('task-42');
assert.deepEqual(replay('task-42', loaded.events.map((item) => item.event)), loaded.state);
}
await selfTest();

内存存储只用于确定性教学测试。真实数据库适配器必须用一个事务同时:验证版本、插入 Event ID 唯一记录、追加事件、更新当前状态。若这些步骤分开提交,就可能出现事件已存但状态未更新,或状态已更新但事件丢失。

简化 PostgreSQL Schema:

CREATE TABLE task_state (
task_id text PRIMARY KEY,
status text NOT NULL,
version bigint NOT NULL CHECK (version >= 0),
attempt integer NOT NULL CHECK (attempt >= 0),
state_json jsonb NOT NULL
);
CREATE TABLE task_event (
task_id text NOT NULL,
sequence bigint NOT NULL,
event_id text NOT NULL UNIQUE,
event_type text NOT NULL,
payload jsonb NOT NULL,
payload_sha256 text NOT NULL,
created_at timestamptz NOT NULL DEFAULT now(),
PRIMARY KEY (task_id, sequence),
FOREIGN KEY (task_id) REFERENCES task_state(task_id)
);

提交伪代码:

BEGIN;
SELECT version
FROM task_state
WHERE task_id = $task_id
FOR UPDATE;
-- 应用层再次检查version与转换结果。
INSERT INTO task_event (...)
VALUES (...);
UPDATE task_state
SET status = $next_status,
version = $next_version,
attempt = $next_attempt,
state_json = $next_state
WHERE task_id = $task_id
AND version = $expected_version;
-- 若UPDATE影响0行,ROLLBACK并返回VERSION_CONFLICT。
COMMIT;

使用 FOR UPDATE 和版本条件都可以控制竞争;具体选型取决于事务长度、冲突率与数据库。即使用行锁,也要保留唯一约束来防 Event ID 重复。应用层预检查提升错误质量,数据库约束负责最后防线。

快照(Snapshot)快照Snapshot某个已知版本上的完整状态副本,用于缩短恢复路径,但不能替代增量事件和校验。打开术语条目 → 保存某个事件序号处的状态,减少每次恢复需要折叠的事件数:

interface TaskSnapshot {
taskId: string;
throughSequence: number;
state: TaskState;
stateSha256: string;
reducerVersion: string;
}

恢复算法:

  1. 找到不超过目标序号的最新可信 Snapshot。
  2. 验证 Snapshot 的 Schema、Hash、taskId、state.version === throughSequence
  3. throughSequence + 1 开始按序读取事件。
  4. 对每个事件执行同一个纯转换函数;遇到 Gap 或非法转换立即停止。
  5. 将重建状态与当前状态行比较;不同则生成差异报告,不自动选择“看起来更新”的一方。

Snapshot 不是事实来源的替代品。若只保留 Snapshot 并删除事件,就失去完整回放能力;这可能仍是合理存储策略,但必须明确恢复与审计边界。事件保留也不是免费:Schema 演进、隐私删除、存储增长和重放时间都要管理。

旧事件不能因为代码升级就改变语义。常见策略:

  • 事件携带版本,Reducer 兼容旧版本;
  • 在读取边界把旧事件上转换成内部新格式;
  • 进行一次可验证迁移,保留迁移前后 Hash 和映射报告;
  • Snapshot 记录 reducerVersion,不兼容时从更早历史重建。

不要在 Event Payload 上就地修改历史然后称为“同一事件”。如果因合规要求必须删除字段,应记录明确的擦除或重写流程,并说明之后哪些审计能力不再存在。

“有六个状态,所以系统简单”是错误判断。真正影响实现和测试的至少有:

  • 合法转换数量;
  • 每条转换的 Guard 数量;
  • 并发命令的竞争组合;
  • 外部副作用与内部状态的一致性边界;
  • 事件乱序窗口;
  • 超时、重试和人工操作是否增加正交维度;
  • 状态与事件 Schema 的版本数量。

可以用一个工程计数帮助发现组合爆炸,但不要把它当成普适复杂度定律。假设状态集合为 S,事件类型集合为 E,完整枚举上界是:

N_pairs = |S| × |E|

本课 6 个状态、6 个事件类型,至少有 36 个状态—事件组合需要明确“接受或拒绝”。如果每条接受路径还要测试“正确版本、过期版本、重复 Event ID、序号 Gap”,测试维度进一步扩展。

更实用的指标:

transition coverage = 已测试的状态-事件分类 / 已声明的状态-事件分类
rejection coverage = 已测试的拒绝原因 / 全部拒绝原因
race coverage = 已测试的竞争对 / 风险清单中的竞争对

这些指标只说明测试覆盖了模型声明,不证明模型声明完整。新故障仍可能揭示缺失状态。

如果把任务状态写成:

running-with-cancel-requested-and-artifact-uploading

很快会得到笛卡尔积。可以拆成:

type OrthogonalTaskState = {
lifecycle: 'queued' | 'running' | 'terminal';
cancellation: 'none' | 'requested' | 'confirmed';
artifact: 'none' | 'writing' | 'verified' | 'failed';
};

但拆分后必须有跨区域不变量,例如 lifecycle=terminal.completed 要求 artifact=verified。正交状态减少名称爆炸,却增加组合约束。选择依据是各维度是否真的可以独立变化,以及是否能写出完整不变量。

任务卡在 running 且租约过期时,恢复器不能直接写 queued。它至少要回答:

  1. Worker 是否还活着?
  2. 外部副作用是否已经发生?
  3. Artifact 是否存在且属于当前 attempt?
  4. 最后一个权威事件是什么?
  5. 是否有未发布或乱序消息?

恢复命令可以产生不同事件:

LeaseExpiredObserved
ArtifactRecovered
AttemptFailed(retryable=true)
TaskFailed(errorCode=UNKNOWN_EFFECT)
HumanReviewRequested

本课核心状态机没有把所有恢复事件展开,是为了保持示例可读;真实产物应根据可观察事实增加事件,而不是让恢复器直接修改 status

故障检测证据状态机响应不应做什么恢复策略
同一 Event ID 重投Event ID 与 Payload Hash 均相同返回 duplicate,版本不变再次应用转换返回首次处理摘要并确认消息
同一 Event ID、不同 PayloadEvent ID 相同,Hash 不同EVENT_ID_PAYLOAD_CONFLICT任选一个 Payload隔离发送方与消息,人工调查协议破坏
事件序号跳跃sequence > version + 1SEQUENCE_GAP跳过缺失序号继续从权威事件存储补拉;超时后告警
过期命令条件更新影响 0 行VERSION_CONFLICT把 expectedVersion 改新后强写重新读取并重新执行业务 Guard
终态收到迟到事件当前状态属于终态TERMINAL_STATE把状态重新打开记录拒绝;需要新工作时创建新任务或显式恢复命令
Worker 租约过期当前时间超过租约,仅是观察不直接转状态立即重试外部副作用查询 Worker、Artifact、Ledger 后生成恢复事件
Snapshot 与事件重放不一致同序号状态 Hash 不同停止恢复选择字段更多的一方固定 Reducer 版本,重放并生成差异报告
事件存储成功、投影失败权威历史已推进,读模型滞后权威状态不回滚重新提交业务事件重放投影或从 Snapshot 重建
  1. 对初始任务提交 evt-1 TaskStarted
  2. 再提交完全相同事件;
  3. 断言第二次结果为 duplicate
  4. 断言状态仍为 running/version=1/attempt=1
  5. 把同一 eventIdworkerId 改掉后提交;
  6. 断言返回 EVENT_ID_PAYLOAD_CONFLICT
  7. 断言事件表中仍只有一条 evt-1

验收:字节语义相同的重投幂等返回,身份相同但内容不同的消息被隔离。

  1. 把任务推进到 running/version=8
  2. 构造两个基于 version 8 的命令;
  3. 让取消事务先提交到 cancelling/version=9
  4. 完成事务条件更新失败;
  5. 重新读取 version 9,再按 cancelling + TaskCompleted 规则计算;
  6. 断言最终为 completed/version=10,且 Transition Evidence 标记竞争;
  7. 再切换策略为“取消优先”,确认同一场景稳定拒绝完成,而不是由时序随机决定。

验收:两种策略都必须显式、确定、可测试。

  1. 正常提交 sequence 1;
  2. 先送达 sequence 3;
  3. 断言返回 SEQUENCE_GAP,状态不变;
  4. 补拉并提交 sequence 2;
  5. 再提交 sequence 3;
  6. 断言状态依次推进,没有丢弃任何事件;
  7. 构造 sequence 2 但新 Event ID,断言为历史冲突而不是普通重复。

验收:到达顺序不决定权威顺序,缺口不能被静默跳过。

  1. 在 version 50 保存 Snapshot;
  2. 修改 Snapshot 中 attempt,但不更新 Hash;
  3. 恢复器应在重放前拒绝 Snapshot;
  4. 再更新 Hash 模拟“内部一致但语义错误”的 Snapshot;
  5. 从事件 1 全量重放,与 Snapshot 对比;
  6. 差异应进入报告,不自动覆盖事件历史;
  7. 固定 reducerVersion 后重建可信 Snapshot。

验收:Snapshot 是加速器,不是跳过历史验证的许可证。

15. 决策表:当前状态表、事件历史还是二者都要

Section titled “15. 决策表:当前状态表、事件历史还是二者都要”
需求只存当前状态存事件 + 当前状态主要代价
简单 CRUD,历史无业务价值通常足够可能过度设计事件 Schema 与迁移成本
需要解释“为什么到这里”需要额外审计表事件天然提供因果记录日志保留和隐私治理
需要重建多个读模型需要另写变更日志可由历史投影重放时间和版本兼容
需要撤销保存前值或补偿信息仍需显式反向事件事件历史不会自动可逆
高频计数更新条件更新更直接事件量可能很大存储与聚合压力
跨服务流程每个服务维护本地状态事件可连接流程,但不提供全局事务顺序、幂等、对账复杂度

Event Sourcing 不是“有事件消息就算”。如果数据库只保存最终状态、消息只是通知,仍然是状态存储加事件发布;这完全可以是更合适的设计。关键是明确事实来源和恢复路径。

可测试任务状态机应满足:

  • □ 命令、事件、状态、观察有不同类型和命名语义;
  • □ 所有状态 × 事件组合都有接受或拒绝分类;
  • □ 终态不能被普通事件重新打开;
  • completed 强制要求已验证 Artifact;
  • cancelledcancelling 的事实边界明确;
  • □ 每个事件具有 Event ID、任务内 Sequence 与 Payload Hash;
  • □ 相同 Event ID 的相同 Payload 幂等返回,不增加版本;
  • □ 相同 Event ID 的不同 Payload 被分类为协议冲突;
  • □ 版本条件写入与事件追加处于一个事务;
  • □ Gap、过期命令、迟到事件和终态事件都有测试;
  • □ 回放器遇到非法历史会停止,不跳过;
  • □ Snapshot 记录序号、Hash 与 Reducer 版本;
  • □ 至少完成三个故障注入实验并保存 Transition Evidence;
  • □ Trace 中不包含密钥、真实账号数据或无必要的输入全文。
  1. Command 与 Event 的时间语义为什么不同?
  2. 为什么 CancellationRequested 不能直接等同于 TaskCancelled
  3. 如何从业务不变量推导状态,而不是先列状态名称?
  4. δ(state,event) 返回拒绝为什么与返回下一状态同样重要?
  5. Event ID 去重和业务语义去重分别防什么?
  6. 为什么过期写入不能简单更新 expectedVersion 后重试?
  7. 任务内 Sequence 与消息到达顺序有什么区别?
  8. 遇到 Sequence Gap 有哪些策略,各自边界是什么?
  9. 为什么 Snapshot 必须记录 Reducer 版本和 Hash?
  10. 什么时候只存当前状态更合适,什么时候需要完整事件历史?
  11. 状态数量为什么不能单独代表复杂度?
  12. 租约超时后,恢复器为什么要先观察副作用再推进状态?
  • Statecharts 论文用于理解层级、并发和事件驱动状态模型的表达能力;本课使用的是较小的有限状态机,不声称实现完整 Statecharts 语义。
  • Event Sourcing 材料用于区分事件历史、当前状态和投影;是否采用完整事件溯源仍取决于审计、恢复、隐私和运维成本。
  • PostgreSQL 事务隔离文档用于核对示例数据库的并发边界;不同数据库的锁、唯一约束、序列化失败和事务语义必须按其官方文档验证。
  • 所有 taskId、workerId、Event ID、时间和 Hash 都是抽象教学数据,没有真实账号或私有系统事实。