长流程真正棘手的地方,不是把步骤串起来,而是进程崩溃后还能判断应该从哪里继续。本文用 Node.js 与 SQLite 实现一个包含模型调用、工具执行、人工确认和发布的状态机,并处理租约、幂等、超时重试与人工恢复。
为什么不能只写一串 await
假设任务依次执行四步:调用模型生成草稿、运行文本分析工具、等待人工确认、发布结果。最直接的实现是一串 await,但只要进程在中间退出,内存中的进度就全部丢失。
更麻烦的是,恢复并不等于简单重跑:模型请求可能已经成功但结果尚未写盘;工具可能已经产生副作用;用户也可能连续点击两次“确认”。因此需要把执行模型拆成三个概念:
| 概念 | 作用 |
|---|---|
| 状态机 | 明确当前步骤及允许的下一步 |
| 检查点 | 将步骤结果持久化,重启后继续 |
| 幂等键 | 识别重复启动、重复确认和重复副作用 |
示例状态依次为 draft → tool → confirm → publish。工作流本身还有 running、waiting、done、failed 和 canceled 等状态。每完成一步,结果和下一步骤在同一次数据库更新中落盘。
SQLite 表与执行约束
示例使用 better-sqlite3,因为它提供真实、直接的事务 API。运行前准备项目:
mkdir recoverable-workflow && cd recoverable-workflow
npm init -y
npm pkg set type=module
npm install better-sqlite3
数据库包含三张表:workflows 保存状态和上下文;requests 记录外部请求 ID,拦截重复启动或确认;effects 保存工具与发布产生的本地副作用。lock_token 和 lock_until 构成执行租约,避免两个进程同时推进同一任务。
SQLite 适合单机服务、后台任务和中小规模内部工具。开启 WAL 后,读写并发会更友好,但它仍不是分布式协调系统;如果任务要跨多台机器高并发执行,应迁移到支持行锁的数据库或专用队列。
完整可运行的状态机
将下面内容保存为 app.js。模型调用使用可中止的模拟适配器,便于直接运行;接入真实模型时,只需替换 callModel,并继续接收 AbortSignal。
import Database from 'better-sqlite3';
import { randomUUID } from 'node:crypto';
import { setTimeout as delay } from 'node:timers/promises';
const db = new Database('workflow.db');
db.pragma('journal_mode = WAL');
db.exec(`
CREATE TABLE IF NOT EXISTS workflows (
id TEXT PRIMARY KEY,
status TEXT NOT NULL,
step TEXT NOT NULL,
input TEXT NOT NULL,
ctx TEXT NOT NULL DEFAULT '{}',
attempts INTEGER NOT NULL DEFAULT 0,
next_at INTEGER NOT NULL DEFAULT 0,
lock_token TEXT,
lock_until INTEGER NOT NULL DEFAULT 0,
error TEXT
);
CREATE TABLE IF NOT EXISTS requests (
request_id TEXT PRIMARY KEY,
workflow_id TEXT NOT NULL,
kind TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS effects (
effect_key TEXT PRIMARY KEY,
result TEXT NOT NULL
);
`);
const read = id => db.prepare('SELECT * FROM workflows WHERE id=?').get(id);
const merge = (w, patch) => JSON.stringify({ ...JSON.parse(w.ctx), ...patch });
const start = db.transaction((requestId, topic) => {
const old = db.prepare('SELECT workflow_id FROM requests WHERE request_id=?').get(requestId);
if (old) return old.workflow_id;
const id = randomUUID();
db.prepare(`INSERT INTO workflows(id,status,step,input) VALUES(?,'running','draft',?)`).run(id, topic);
db.prepare(`INSERT INTO requests VALUES(?,?,'start')`).run(requestId, id);
return id;
});
async function callModel(topic, signal) {
const ms = Number(process.env.MODEL_DELAY || 300);
await delay(ms, undefined, { signal });
return `关于“${topic}”的待审核草稿。`;
}
async function withTimeout(fn, ms) {
const controller = new AbortController();
const timer = setTimeout(() => controller.abort(), ms);
try { return await fn(controller.signal); }
finally { clearTimeout(timer); }
}
function guarded(sql, params, token) {
const result = db.prepare(`${sql} WHERE id=? AND lock_token=?`).run(...params, token);
if (result.changes !== 1) throw new Error('执行租约已失效');
}
async function run(id) {
const now = Date.now();
const token = randomUUID();
const claimed = db.prepare(`
UPDATE workflows SET lock_token=?, lock_until=?
WHERE id=? AND status='running' AND next_at<=? AND lock_until<?
`).run(token, now + 30000, id, now, now);
if (!claimed.changes) return console.log('当前没有可执行步骤');
let w = read(id);
try {
while (w.status === 'running') {
if (w.step === 'draft') {
const draft = await withTimeout(signal => callModel(w.input, signal), 1000);
if (process.env.CRASH_AFTER_MODEL === '1') process.exit(9);
guarded(`UPDATE workflows SET step='tool',ctx=?,attempts=0,error=NULL`, [merge(w, { draft })], token);
} else if (w.step === 'tool') {
db.transaction(() => {
const key = `${id}:tool`;
let effect = db.prepare('SELECT result FROM effects WHERE effect_key=?').get(key);
if (!effect) {
effect = { result: JSON.stringify({ characters: JSON.parse(w.ctx).draft.length }) };
db.prepare('INSERT INTO effects VALUES(?,?)').run(key, effect.result);
}
guarded(`UPDATE workflows SET step='confirm',ctx=?`, [merge(w, { analysis: JSON.parse(effect.result) })], token);
})();
} else if (w.step === 'confirm') {
guarded(`UPDATE workflows SET status='waiting',lock_token=NULL,lock_until=0`, [], token);
} else if (w.step === 'publish') {
db.transaction(() => {
const key = `${id}:publish`;
db.prepare('INSERT OR IGNORE INTO effects VALUES(?,?)').run(key, JSON.stringify({ publishedAt: new Date().toISOString() }));
guarded(`UPDATE workflows SET status='done',lock_token=NULL,lock_until=0`, [], token);
})();
}
w = read(id);
}
console.log(JSON.stringify(w, null, 2));
} catch (error) {
const attempts = w.attempts + 1;
const status = attempts >= 3 ? 'failed' : 'running';
const nextAt = status === 'failed' ? 0 : Date.now() + 1000 * 2 ** attempts;
db.prepare(`UPDATE workflows SET status=?,attempts=?,next_at=?,error=?,lock_token=NULL,lock_until=0 WHERE id=? AND lock_token=?`)
.run(status, attempts, nextAt, String(error.message), id, token);
console.error(`执行失败:${error.message};状态=${status}`);
}
}
const approve = db.transaction((id, requestId, accepted) => {
if (db.prepare('SELECT 1 FROM requests WHERE request_id=?').get(requestId)) return '重复确认已忽略';
const w = read(id);
if (!w || w.status !== 'waiting' || w.step !== 'confirm') throw new Error('任务不在待确认状态');
db.prepare(`INSERT INTO requests VALUES(?,?,'approve')`).run(requestId, id);
if (!accepted) db.prepare(`UPDATE workflows SET status='canceled' WHERE id=?`).run(id);
else db.prepare(`UPDATE workflows SET status='running',step='publish',next_at=0 WHERE id=?`).run(id);
return accepted ? '已确认,可继续运行' : '已取消';
});
const [command, ...args] = process.argv.slice(2);
if (command === 'start') console.log(start(args[0], args.slice(1).join(' ')));
else if (command === 'run') await run(args[0]);
else if (command === 'approve') console.log(approve(args[0], args[1], args[2] === 'yes'));
else if (command === 'status') console.log(JSON.stringify(read(args[0]), null, 2));
else console.log('用法: start <请求ID> <主题> | run <任务ID> | approve <任务ID> <请求ID> yes|no | status <任务ID>');
验证重启、重复提交与重试
先创建任务,记下返回的任务 ID:
node app.js start req-001 "SQLite 工作流恢复"
node app.js run <任务ID>
node app.js status <任务ID>
第一次 run 会执行模型和工具,然后停在 waiting。人工确认后再次运行:
node app.js approve <任务ID> approve-001 yes
node app.js run <任务ID>
重复使用 req-001 启动时,会得到原任务 ID;重复使用 approve-001 确认时,也不会再次推进状态。这要求调用方为每次逻辑操作生成稳定的请求 ID,而不是每次重试都生成新值。
可以用环境变量模拟崩溃和超时:
CRASH_AFTER_MODEL=1 node app.js run <任务ID>
node app.js run <任务ID>
MODEL_DELAY=2000 node app.js run <任务ID>
崩溃后需等待 30 秒租约过期再运行。超时会记录错误,按指数间隔重试,三次后进入 failed。生产环境通常还需要一个定时扫描器,持续查找 next_at <= 当前时间 的任务;若步骤可能超过租约时长,还应定期续租,而不是盲目增大租约。
幂等的边界与人工恢复
数据库事务只能保证 SQLite 内部原子性,不能让一次外部模型请求与本地写盘组成分布式事务。如果模型已经返回、进程却在保存检查点前退出,恢复后仍可能再次调用模型。若供应商支持幂等键,应传入类似 工作流ID:步骤名 的稳定键;不支持时,要接受“可能重复计算”,并确保后续只采用一个持久化结果。
具有外部副作用的工具更需要谨慎。例如发送邮件、创建工单或合并代码,应让目标系统接收幂等键,或采用发件箱模式:先在本地事务中写入待发送记录,再由独立派发器重试。不要把“步骤结果保存过”误认为“外部副作用一定只发生一次”。
进入 failed 后,人工恢复可以先查看 error 与 ctx,修复输入或外部依赖,再执行一条受审计的管理命令,将状态改回 running、清零 attempts 并设置合理的 step。管理接口必须校验允许的状态迁移,不应向操作者暴露任意 SQL。
总结
这套实现没有依赖特定 Agent 框架,核心要点是:
- 用显式状态机描述步骤,而不是依赖调用栈保存进度;
- 每一步完成后持久化结果,重启时从数据库状态继续;
- 用请求 ID、效果键和唯一约束处理重复提交;
- 用带所有权的租约避免并发执行,并让崩溃后的任务可重新领取;
- 为模型调用设置可中止超时,对临时错误退避重试;
- 区分数据库原子性与外部副作用幂等,必要时使用目标端幂等键或发件箱模式。
将步骤替换为资料检索、内容生成、代码扫描、人工审查和最终提交后,同一模型就能用于内容生产、代码审查及其他需要跨分钟甚至跨天运行的任务。
评论