跳到正文

持久化任务

用 PostgreSQL queue 执行普通 async handler,明确归属、重试、协作取消与 worker drain。

@lenso/tasks 是 Bun 与 PostgreSQL 上的可选持久队列。Producer 入队,显式启动的 worker 调用普通 async 服务。当前交付的 provider 使用 pg-boss 12.37.0;不包含 D1、Cloudflare Queues、cron、workflow/DAG 引擎或默认 HTTP 管理 API。

本页对应已发布 @lenso/tasks@0.2.0 与 Core 0.2.0,按 Tasks 实现 549b987及真实 tarball 审计。匹配 CLI、Manage 与依赖版本见安装。

使用现有 report 示例

examples/tasks 包含 producer、worker、显式迁移和授权 CLI/MCP 入口。真实业务结果写入 task_example_reports,稳定 reportId 用于主键 upsert。Queue result 是安全摘要,不替代业务表。

示例直接封装现有服务:

import { defineTask } from "@lenso/tasks";import { reportInput, type createReportService } from "./report-service";
export function createReportTask(service: ReturnType<typeof createReportService>) {  return defineTask({    name: "generate-report",    input: reportInput,    maxAttempts: 3,    retry: { delaySeconds: 2, backoff: true, maxDelaySeconds: 10 },    async handler(input, { attempt, signal, logger }) {      logger?.info({ event: "report-started" }, "Report generation started");      return service.generate(input, { attempt, signal });    },    result: (returned) => ({ sum: returned.sum, count: returned.count }),  });}

reportInput 与 createReportService 是该例子已有导出。服务在写入前检查取消。Upsert 为每个 report 保留一行,但后续使用同 reportId 和不同 rows 会覆盖摘要;这不是业务输入冲突检测。

先 provisioning,再生产和消费

@lenso/tasks/postgres 公开导出 createPostgresTaskProvider 与 migratePostgresTaskQueue。Migrator 只从显式 provisioning 入口对已创建、授权修改的数据库执行。运行时导入或启动不会安装、升级 schema。Storage/Tasks 的 SQL 在 tarball 中并不意味着内部路径是 public exports。

构建框架包后,在经过核对的源码 checkout 中使用:

# 在 examples/tasks 内;DATABASE_URL 已指向独立本地数据库。export TASK_QUEUE_NAME=reportsbun run migrate# Worker 终端必须使用同一数据库和 queue name。bun run worker

Producer 调用前配置 TASK_SESSION 与 TASK_AUTH_SOURCE_MODULE。后者是受信任的应用模块,导出 connectTaskAuth(),返回真实验证身份的 Auth source 和可选自有资源 cleanup。它不是 tool 参数。把任意 token 解释为 subject ID 不构成身份验证。Worker 与迁移无需用户 session。

在源码根目录检查并调用真正声明的操作:

bun run cli inspect tasks submit --root examples/tasks --jsonprintf '%s\n' '{"reportId":"daily-report","rows":[1,2]}' |  bun run cli call tasks submit --root examples/tasks --stdin --json# 使用返回的 jobId,保持同一受信任启动身份。printf '%s\n' '{"jobId":"<returned-job-id>"}' |  bun run cli call tasks query --root examples/tasks --stdin --json

该 query 返回 state、attempt、maxAttempts、cancelRequested 或 null,不代表 worker 健康检查。Worker/supervisor 的安全日志才提供相应运行证据。

管理与 transport 独立选择

现有 createTasksOperations 返回 { plugin, operations, manage }。Manage 选择 submit/query/cancel/retry,CLI 另包含 report。Named mcpOperations 与实际 stdio host allowlist 仅选择四项管理操作,默认 MCP 不暴露 report。

import { selectManageOperations } from "@lenso/manage";import { taskOperations } from "./src/plugin";
const queryOnly = selectManageOperations(taskOperations.manage, ["query"]);// queryOnly[0] 与应用 CLI 使用同一声明。

显式 Manage adapter选择这些声明,借用已运行的精确 Tasks plugin,绑定入口策略并实现当前调用者的 canList。它不启动 worker、不授予管理员 queue 权限。当前 Tasks method 仍为单输入,从配置的启动 source/evidence 认证;增加 adapter 不会变成逐请求身份。远程多用户入口需要显式设计 contextual service 边界,不能用共享可变 session 或 JSON actor 替代。

Manage 仍执行有限调用:submit 返回授权 job ID,query 检查持久 owner,cancel/retry 保留既有队列语义。Tool 请求取消不会调用 tasks.cancel。若显式增加 confirmation/approval metadata,必须提供可信入口 callback,也不能替代 owner 检查。

区分 queue 与资源归属

独立入口使用 createTaskQueue({ provider, tasks });需要 Lenso 生命周期时使用 createTaskPlugin({ id, tasks, connect, worker? })。省略 worker 是 producer-only 插件。不同 queueName 隔离持久队列,不同 plugin ID 隔离进程内实例。同队列的 worker 竞争任务,应注册完整 task 集合;concurrency 只作用于当前 worker。

Provider 接受自有 connectionString 或调用方持有的 pg.Pool,不能同时提供。借用 pool 不会在正常清理、失败或迁移时关闭。这是 Bun 上的 node-postgres 驱动,不能把 Bun SQL pool 强转成 pg.Pool。业务 Drizzle 资源仍可使用自己的原生驱动。

观察 worker.done。关闭时调用 worker.stop() 或 stop({ abort: true }),等待实际 handler 结束,再关闭 queue 和其他自有 client。插件的提前 close()、回滚及 app stop 共享已登记 disposer;handler 还在使用 DB 时不能先关闭 DB。

Payload、结果与授权

同一个 Standard Schema v1 在入队和执行前验证。持久化保存原始 plain JSON,避免 transform 跨进程重复叠加;schema 应确定且无副作用。输入和验证输出须是有限、无环 JSON,最多 64 KiB UTF-8、深度 32。Date、BigInt、undefined、class instance、accessor 和稀疏数组均无效。

Handler 返回值默认丢弃;只有显式 result 投影才保存安全摘要,最多 16 KiB。Queue 错误为 handler-failed、invalid-input、invalid-result、aborted,不存储原始 exception 文本或 stack。

Queue 方法是受信任的基础设施 API。Report 示例为每个 reportId 持久保存不可变 (realmId, subjectId) owner,并将 queue/job ID 映射到它。Query/cancel/retry/report 访问先验证身份和归属;payload 不含 actor 或凭据。归属 reservation、enqueue 与 mapping 是分开的事务,崩溃可能留下无法访问的孤儿任务;失败不能放宽策略。

重试是至少一次投递

maxAttempts: 3 包括第一次。自动重试使用声明的 delay/backoff,带 jitter,并非精确执行时间;maxDelaySeconds 要求 backoff: true。手动 retry(jobId) 只接受最终 failed,保留 job/input/计数并增加一次机会;其他状态或已清理 job 返回 false。

去重按 queue + key 生效:相同 task/input 返回原 ID,冲突输入被拒绝。它不能使外部副作用恰好一次。Worker 可能提交 report 后、acknowledgement 前死亡,随后再次执行。应在真实业务边界执行幂等或 fencing;attempt 不是业务幂等 key。

取消记录意图

cancel(jobId) 结果含义
cancelled原子阻止 pending job 被 claim
requestedRunning job 已持久记录取消请求;handler 仍可能运行
terminal保留中的 job 已结束
missing当前 queue 没有保留的 job

运行中的 job 在 handler 结束前保持 running 与 cancelRequested: true。传播并检查 AbortSignal,再执行后续副作用;取消不回滚已提交写入。Terminal 状态不能证明旧的丢失 lease 的 attempt 已物理停止。

timeoutMs、lease 丢失和 worker abort 是协作机制,可能产生重试,不是持久用户取消。忽略 abort 的 handler 继续占用 slot,可能无限延迟关闭。Lease 过期时,其他进程可能开始执行而旧 handler 仍在退出。CLI/MCP 请求取消又是另一个边界,见 Agents。

Retention 与恢复

默认轮询 1,000 ms、heartbeat 30 秒、attempt expiry 900 秒、retention 7 天。Heartbeat 不延长 attempt expiry。恢复旧 claim 需要存活的 provider/supervisor;数据库不会在所有进程离线时自己运行 JavaScript 恢复循环。

清理后 get 返回 null,取消返回 missing,retry 返回 false。Dedup mapping 保留;过期 mapping 被拒绝,不会静默产生新任务。没有历史 archive 或自动 dedup 清理。缺少状态不证明任务从未运行。

恢复前检查应用已有授权 query、持久 owner、attempt budget、业务结果与安全 worker 证据。确认幂等后只重试符合条件的最终失败。Acknowledgement 不确定时先查状态,不盲目 replay。缺少 schema 是 provisioning 错误,不是自动迁移的许可。

参阅真实 PostgreSQL 测试、授权服务测试及测试指南。跳过数据库检查不构成持久性验证。