Skip to content

Tasks:持久化任务

源码版本v2026.6.11

职责

task-executor.tsdetached-task-runtime.ts 构成 OpenClaw 的"任务 (task)"子系统:把任何可能跨进程、跨会话、跨重启的异步工作建模成一条带生命周期的记录,无论它在哪个 runtime(subagent / acp / cli / cron)里跑,都共享同一份状态机 (state machine)、同一份 SQLite 持久化 (persistence)、同一份取消与回收语义。CronService 负责调度时机,task-executor.ts 负责"一条记录怎么从 queued 走到 succeeded/failed/timed_out/cancelled/lost"。

任务不是线程也不是 worker——它只是一个状态对象 (TaskRecord,TaskRecord:116-146) 加上对应的 runtime 自己实现的生命周期钩子。runtime(subagent/acp/cli/cron,TaskRuntime:5)决定怎么真正执行,task 系统只管:存记录、转状态、投递回原会话、定期清理过期记录。

设计动机

为什么 cron 不够,还要单独抽一层 task?

  1. 跨 runtime:cron 是定时,ACP 是 IDE 触发的 prompt,subagent 是 agent 主循环派生的子任务,CLI 是用户在终端敲的命令。四种 runtime 的触发时机完全不同,但它们都要面对同样的问题:进程崩溃后怎么知道某条工作没跑完?gateway 重启后怎么回收?结果怎么回到原会话?把这些问题抽成 task registry,每种 runtime 只实现自己的 DetachedTaskLifecycleRuntime(DEFAULT_DETACHED_TASK_LIFECYCLE_RUNTIME:35-49),其他逻辑共享。
  2. flow 概念:很多任务属于一个更大的 flow(比如一个 agent run 触发三个子任务,主任务取消时三个子任务也要取消)。task-flow-registry(src/tasks/task-flow-registry.ts) 给任务一个父 flow,可以批量取消、批量查状态。对于 detached ACP/subagent run,executor 自动给它包一层"one-task flow"(isOneTaskFlowEligible + ensureSingleTaskFlow:46-92),让这些单次 run 也能享受 flow 的状态/重试面。
  3. delivery 状态:任务结束之后,结果不一定要回到原会话——可能用户已经关了 IDE、可能 cron 是 isolated session 没有原会话。TaskDeliveryStatus(TaskDeliveryStatus:16-23)单独跟踪这条状态:pending/delivered/session_queued/failed/parent_missing/not_applicable,把"任务执行成功"和"结果送到客户端"解耦。

关键文件

数据流

任务的状态机只有 7 个状态(TaskStatus:7-14):

typescript
export type TaskStatus =
  | "queued"       // 已落库,等 runtime 拉起
  | "running"      // runtime 报告开始执行
  | "succeeded"    // 终态:成功
  | "failed"       // 终态:失败(带 error)
  | "timed_out"    // 终态:超时
  | "cancelled"    // 终态:用户/系统取消
  | "lost";        // 终态:sweeper 在 grace 后仍未观察到活跃,标记丢失

createQueuedTaskRun / createRunningTaskRun 是创建入口,二者都会调 ensureSingleTaskFlow(ensureSingleTaskFlow:56):

typescript
function isOneTaskFlowEligible(task: TaskRecord): boolean {
  if (task.parentFlowId?.trim() || task.scopeKind !== "session") return false;
  if (task.deliveryStatus === "not_applicable") return false;
  return task.runtime === "acp" || task.runtime === "subagent";
}

function ensureSingleTaskFlow(params: { task; requesterOrigin? }): TaskRecord {
  if (!isOneTaskFlowEligible(params.task)) return params.task;
  try {
    const flow = createTaskFlowForTask({ task: params.task, requesterOrigin: params.requesterOrigin });
    if (!flow) return params.task;
    const linked = linkTaskToFlowById({ taskId: params.task.taskId, flowId: flow.flowId });
    if (!linked) { deleteTaskFlowRecordById(flow.flowId); return params.task; }
    if (linked.parentFlowId !== flow.flowId) { deleteTaskFlowRecordById(flow.flowId); return linked; }
    return linked;
  } catch (error) {
    log.warn("Failed to create one-task flow for detached run", { ... });
    return params.task;
  }
}

关键判断:scopeKind === "session"runtime ∈ {acp, subagent}deliveryStatus !== "not_applicable"——这三条满足时才给它自动包一层 flow。系统级任务(scopeKind: "system")和 cron 任务有自己单独的 cancel/delivery 路径,不重走 flow。失败时只 warn 不抛,保证任务创建本身不会因为 flow 包装失败被阻塞。

detached-task-runtime.ts(getDetachedTaskLifecycleRuntime:47)对外暴露统一 API,但具体实现可以插件替换:

typescript
const DEFAULT_DETACHED_TASK_LIFECYCLE_RUNTIME: DetachedTaskLifecycleRuntime = {
  createQueuedTaskRun: createQueuedTaskRunFromExecutor,
  createRunningTaskRun: createRunningTaskRunFromExecutor,
  startTaskRunByRunId: startTaskRunByRunIdFromExecutor,
  recordTaskRunProgressByRunId: recordTaskRunProgressByRunIdFromExecutor,
  finalizeTaskRunByRunId: finalizeTaskRunByRunIdFromExecutor,
  completeTaskRunByRunId: completeTaskRunByRunIdFromExecutor,
  failTaskRunByRunId: failTaskRunByRunIdFromExecutor,
  setDetachedTaskDeliveryStatusByRunId: setDetachedTaskDeliveryStatusByRunIdFromExecutor,
  cancelDetachedTaskRunById: cancelDetachedTaskRunByIdInCore,
};

getDetachedTaskLifecycleRuntime() 优先返回插件注册的实现,否则用默认——这让 plugin 可以接管整个 task 生命周期(例如把 task 跑到独立 worker 进程)而不必改 task-executor。

边界与失败

  • 进程级索引的 symbol-keyed globalThis:getTaskRegistryProcessState()(getTaskRegistryProcessState:18)用 Symbol.for("openclaw.taskRegistry.state") 挂到 globalThis,这样模块热重载(dev watch)或测试隔离时索引仍然共享——否则不同模块实例会各自维护一份 in-memory 索引,落库到 SQLite 的任务和内存索引会失配。
  • sweeper 批次 yield:SWEEP_YIELD_BATCH_SIZE = 25(SWEEP_YIELD_BATCH_SIZE:83),每处理 25 条任务就让出事件循环,避免大批量清理时阻塞主线程。
  • stale running 判定:TASK_STALE_RUNNING_MS = 30 * 60_000——running 状态超过 30 分钟无任何事件更新(包括 recordTaskRunProgressByRunId),sweeper 才会开始回收。这是给慢任务一个宽松窗口,不会误杀正常长跑。
  • recovery hook 容错:tryRecoverTaskBeforeMarkLost(tryRecoverTaskBeforeMarkLost:134)里 try/catch 包住整个 hook 调用——任何 plugin 实现 throw、返回非法对象、超过 5s 阈值,都只 warn,然后继续走 mark-lost。理由:recovery 是 best-effort,不能让坏插件阻塞清理。
  • retention 双轨:终态任务保留 7 天,lost 任务只保留 1 天(resolveTaskRetentionMs:5-9)。lost 是"sweeper 推测出来的终态",可信度低于 runtime 显式上报的 succeeded/failed,所以保留期短——既给排查窗口,又不让"幽灵任务"长期占用 DB。
  • one-task flow 失败降级:ensureSingleTaskFlow 任何一步抛错都返回原 task(不重新包装),日志里记 warn。任务能创建成功比 flow 包装完整更重要。
  • scopeKind 区分:scopeKind: "system" 的任务没有 requesterSessionKey(system-scope requesterSessionKey:92-94),ownerKey 是唯一查找锚——这避免了系统任务和某个 session 的 ownerKey 撞车。
  • cron-task-cancel settlement grace:gateway 重启时 abortActiveCronTaskRuns(abortActiveCronTaskRunSettlementGrace:50)会启动一个 60 秒的"沉降宽限期",让被 abort 的 promise 有时间在 finally 里清完资源、写完终态——直接 clear 表会让 cleanup 走丢。

小结

Tasks 子系统把"一条异步工作"抽象成 7 状态的 TaskRecord,四种 runtime(subagent/acp/cli/cron)共用一套状态机 + SQLite 持久化 + sweeper 回收。它和 Cron:定时任务 是双向关系:cron 通过 registerActiveCronTaskRun 把自己挂到 task 进程级表里,task 系统又通过 cron-task-cancel.ts 给 cron 提供"重启时打断活跃 run"的能力。ACP 启动的 detached prompt 走的也是同一套(见 ACP:IDE 桥接),只是 runtime 标成 "acp"。session 关联状态和 task delivery 回投会用到 Memory 文件 里的 session store 路径。