Durable tasks
Run ordinary async handlers through the PostgreSQL queue with explicit ownership, retries, cooperative cancellation and worker drain.
@lenso/tasks is an optional durable queue for Bun and PostgreSQL. Producers enqueue work; separately started workers call ordinary async services. The delivered provider is PostgreSQL, using pg-boss 12.37.0. There is no D1, Cloudflare Queues, cron, workflow/DAG engine or default HTTP administration API.
This guide covers published @lenso/tasks@0.2.0 with Core 0.2.0, audited against the Tasks implementation at 549b987 and actual tarball. See installation for the matching CLI, Manage and dependency versions.
Follow the existing report example
examples/tasks contains producer, worker, migration and authorized CLI/MCP entries. Its business result is in task_example_reports; a stable reportId drives a primary-key upsert. Queue results are safe summaries, not the authoritative business table.
The example task wraps that service directly:
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 and createReportService are the example's existing exports. The service checks cancellation before its write. Its upsert keeps one row per report but permits a later submission with different rows to replace the summary; it does not detect conflicting business input.
Provision, produce, then consume
@lenso/tasks/postgres exports both createPostgresTaskProvider and migratePostgresTaskQueue. Invoke the migrator only from an explicit provisioning entry against an already-created, authorized database. Runtime import/startup never installs or upgrades schema. Storage/Tasks SQL files being present in a tarball does not make their internal paths public exports.
For the reviewed checkout's Tasks example, after building framework packages:
# In examples/tasks; DATABASE_URL is supplied for a dedicated local DB.export TASK_QUEUE_NAME=reportsbun run migrate# Worker terminal; requires the same database and queue name.bun run workerConfigure TASK_SESSION and TASK_AUTH_SOURCE_MODULE before producer calls. The module is a trusted application entry exporting connectTaskAuth(), returning a real verifying Auth source and optional owned-resource cleanup. It is not a tool argument. A source that interprets arbitrary tokens as subject IDs is not an authentication implementation. Workers/migrations need no user session.
From the checkout root, inspect and use the real declared operations:
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# Substitute the returned jobId, keeping the same trusted launch identity.printf '%s\n' '{"jobId":"<returned-job-id>"}' | bun run cli call tasks query --root examples/tasks --stdin --jsonThe query returns state, attempt, maxAttempts and cancelRequested, or null. It does not report worker health. Consult the existing worker/supervisor's safe logs for that evidence.
Keep management and transport selections independent
The existing createTasksOperations factory returns { plugin, operations, manage }. Manage selects submit/query/cancel/retry; CLI selects those plus report. Named mcpOperations and the actual stdio host allowlist select only the four management operations, so report is not exposed through the default MCP entry.
import { selectManageOperations } from "@lenso/manage";import { taskOperations } from "./src/plugin";
const queryOnly = selectManageOperations(taskOperations.manage, ["query"]);// queryOnly[0] is the same declaration used by the application's CLI.An explicit Manage adapter selects these declarations, borrows the already-running exact Tasks plugin, binds entry policy and implements caller-specific canList. This does not start a worker or give administrative queue access. The present Tasks methods are single-input and authenticate from the configured launch source/evidence; adding an adapter does not turn them into per-request identities. A remote multi-user entry needs an explicitly designed contextual service boundary, not a mutable shared session or an actor in JSON.
Manage remains finite invocation: submit returns an owned job ID, query checks durable ownership, and cancel/retry preserve their existing queue semantics. A tool request cancellation does not call tasks.cancel. Confirmation/approval metadata, if deliberately added, requires trusted entry callbacks and never substitutes for owner checks.
Separate queue and resource ownership
Use createTaskQueue({ provider, tasks }) for a standalone entry, or createTaskPlugin({ id, tasks, connect, worker? }) for Lenso lifecycle integration. Omitting worker creates a producer-only plugin. Different queueName values separate durable queues; plugin IDs separate in-process instances. Workers sharing a queue compete for jobs and must register the full task set. concurrency is local to each worker.
The provider accepts its own connectionString or a caller-owned pg.Pool, never both. A borrowed pool is not closed during normal cleanup, failure or migration. This is the node-postgres driver on Bun, not a Bun SQL pool cast. A separate Drizzle business resource can retain its native driver.
Observe worker.done. On shutdown, call worker.stop() or stop({ abort: true }), await actual handlers, close the queue, then close other owned clients. The plugin automatically uses its registered disposer for early close(), rollback and app stop. Do not close a DB while handlers still use it.
Payloads, results and authorization
The same Standard Schema v1 validates input before enqueue and execution. Persisted input is raw plain JSON so transformations do not compound; schemas should be deterministic and side-effect free. Inputs and validated values are finite, acyclic JSON, at most 64 KiB UTF-8 and depth 32. Dates, BigInt, undefined, class instances, accessors and sparse arrays are invalid.
Handler values are discarded unless result explicitly projects a safe JSON summary, at most 16 KiB. Queue error codes are handler-failed, invalid-input, invalid-result and aborted, without exception text or stacks.
Queue methods are trusted infrastructure APIs. The report example persists immutable (realmId, subjectId) ownership for each reportId and maps queue/job IDs to it. It authenticates and authorizes before query/cancel/retry/report access. Payload JSON contains no actor or credential. Ownership reservation, enqueue and mapping are separate transactions: a crash can leave an inaccessible orphan job. Failure does not relax policy.
Retries are at-least-once delivery
maxAttempts: 3 includes the first attempt. Automatic retry uses the declared delay/backoff; backoff timing is jittered, not an exact execution schedule. maxDelaySeconds requires backoff: true. Manual retry(jobId) accepts only final failure, preserves job/input/counter, and adds one attempt; it returns false for other states or a pruned job.
Queue deduplication is scoped to queue + key. The same task/input returns the original ID; conflicting task/input is rejected. It does not make external effects exactly-once. A worker can commit a report, die before acknowledgement, and execute again. Fence/idempotently apply effects at their actual business boundary; an attempt number is not a business idempotency key.
Cancellation records intent
cancel(jobId) result | Meaning |
|---|---|
cancelled | Pending claim prevented atomically |
requested | Running job has a durable request; handler may still run |
terminal | Retained job has already ended |
missing | No retained job in this queue |
A running job stays running with cancelRequested: true until its handler settles. Propagate/check AbortSignal before further side effects. Cancellation does not roll back a committed write. A terminal status cannot prove an old lease-losing attempt has physically stopped.
timeoutMs, lease loss and worker abort are cooperative and can lead to retry; they are not durable user cancellation. An abort-ignoring handler keeps its slot and can delay shutdown indefinitely. Lease expiry can let another process execute while the old handler unwinds. CLI/MCP request cancellation is yet another boundary; see agents.
Retention and recovery
Defaults include a 1,000 ms polling interval, 30-second heartbeat, 900-second attempt expiry and 7-day retention. Heartbeats do not extend attempt expiry. A live provider/supervisor is required to recover stale claims; a database alone does not run the JavaScript recovery loop.
After pruning, get is null, cancellation is missing, and retry is false. Dedup mappings remain; an expired mapping is rejected instead of silently creating a new job. There is no history archive or automatic dedup pruning. Missing status is not proof a job never ran.
Before recovery, inspect the authorized query, durable owner, attempt budget, business result and safe worker evidence. Retry only an eligible failure after checking idempotency. On ambiguous acknowledgement, query state before replaying. Missing schema is a provisioning error, not permission to run a migration automatically.
See real PostgreSQL tests, authorized service tests, and testing. Skipped database checks do not establish durability.