长流程真正棘手的地方,不是把步骤串起来,而是进程崩溃后还能判断应该从哪里继续。本文用 Node.js 与 SQLite 实现一个包含模型调用、工具执行、人工确认和发布的状态机,并处理租约、幂等、超时重试与人工恢复。

为什么不能只写一串 await

假设任务依次执行四步:调用模型生成草稿、运行文本分析工具、等待人工确认、发布结果。最直接的实现是一串 await,但只要进程在中间退出,内存中的进度就全部丢失。

更麻烦的是,恢复并不等于简单重跑:模型请求可能已经成功但结果尚未写盘;工具可能已经产生副作用;用户也可能连续点击两次“确认”。因此需要把执行模型拆成三个概念:

概念作用
状态机明确当前步骤及允许的下一步
检查点将步骤结果持久化,重启后继续
幂等键识别重复启动、重复确认和重复副作用

示例状态依次为 draft → tool → confirm → publish。工作流本身还有 runningwaitingdonefailedcanceled 等状态。每完成一步,结果和下一步骤在同一次数据库更新中落盘。

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_tokenlock_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 后,人工恢复可以先查看 errorctx,修复输入或外部依赖,再执行一条受审计的管理命令,将状态改回 running、清零 attempts 并设置合理的 step。管理接口必须校验允许的状态迁移,不应向操作者暴露任意 SQL。

总结

这套实现没有依赖特定 Agent 框架,核心要点是:

  • 用显式状态机描述步骤,而不是依赖调用栈保存进度;
  • 每一步完成后持久化结果,重启时从数据库状态继续;
  • 用请求 ID、效果键和唯一约束处理重复提交;
  • 用带所有权的租约避免并发执行,并让崩溃后的任务可重新领取;
  • 为模型调用设置可中止超时,对临时错误退避重试;
  • 区分数据库原子性与外部副作用幂等,必要时使用目标端幂等键或发件箱模式。

将步骤替换为资料检索、内容生成、代码扫描、人工审查和最终提交后,同一模型就能用于内容生产、代码审查及其他需要跨分钟甚至跨天运行的任务。