Tasks:持久化任务
职责
task-executor.ts 和 detached-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?
- 跨 runtime:cron 是定时,ACP 是 IDE 触发的 prompt,subagent 是 agent 主循环派生的子任务,CLI 是用户在终端敲的命令。四种 runtime 的触发时机完全不同,但它们都要面对同样的问题:进程崩溃后怎么知道某条工作没跑完?gateway 重启后怎么回收?结果怎么回到原会话?把这些问题抽成 task registry,每种 runtime 只实现自己的
DetachedTaskLifecycleRuntime(DEFAULT_DETACHED_TASK_LIFECYCLE_RUNTIME:35-49),其他逻辑共享。 - 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 的状态/重试面。 - delivery 状态:任务结束之后,结果不一定要回到原会话——可能用户已经关了 IDE、可能 cron 是 isolated session 没有原会话。
TaskDeliveryStatus(TaskDeliveryStatus:16-23)单独跟踪这条状态:pending/delivered/session_queued/failed/parent_missing/not_applicable,把"任务执行成功"和"结果送到客户端"解耦。
关键文件
task-registry.types.ts:1-146—TaskRuntime/TaskStatus/TaskDeliveryStatus/TaskNotifyPolicy/TaskRecord全部类型声明。task-executor.ts:46-127—ensureSingleTaskFlow+createQueuedTaskRun+createRunningTaskRun,任务创建入口(含 one-task flow 包装)。detached-task-runtime.ts:32-66—DetachedTaskLifecycleRuntime接口、默认 runtime、插件注册机制。tryRecoverTaskBeforeMarkLost:134-171— 回收 hook:best-effort,5 秒阈值告警,任何抛错都 fallback 到 mark-lost。task-registry.process-state.ts:1-32— 进程级 in-memory 索引(tasks/taskIdsByRunId/taskIdsByOwnerKey/taskIdsByParentFlowId/taskIdsByRelatedSessionKey/tasksWithPendingDelivery),Symbol.for("openclaw.taskRegistry.state")挂到 globalThis。task-registry.store.sqlite.ts:48-77— Kysely 映射的task_runs+task_delivery_state表,镜像到openclaw-state.db。task-registry.maintenance.ts:73-100— 后台 sweeper,周期 60s,TASK_STALE_RUNNING_MS=30*60_000,TASK_RECONCILE_GRACE_MS=5*60_000。task-retention.ts:4-29— 终态任务保留 7 天,lost 任务只保留 1 天。src/tasks/task-flow-registry.store.sqlite.ts— flow 持久化,支持parentFlowId级联。cron-task-cancel.ts:12-77— 进程级 active cron task run 取消句柄,settlement grace 60s。task-registry.reconcile.ts— 公共 facade:reconcileInspectableTasks/reconcileTaskLookupToken。
数据流
任务的状态机只有 7 个状态(TaskStatus:7-14):
export type TaskStatus =
| "queued" // 已落库,等 runtime 拉起
| "running" // runtime 报告开始执行
| "succeeded" // 终态:成功
| "failed" // 终态:失败(带 error)
| "timed_out" // 终态:超时
| "cancelled" // 终态:用户/系统取消
| "lost"; // 终态:sweeper 在 grace 后仍未观察到活跃,标记丢失createQueuedTaskRun / createRunningTaskRun 是创建入口,二者都会调 ensureSingleTaskFlow(ensureSingleTaskFlow:56):
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,但具体实现可以插件替换:
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 路径。