持久化任务
用 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 workerProducer 调用前配置 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 |
requested | Running 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 测试、授权服务测试及测试指南。跳过数据库检查不构成持久性验证。