On this page
Follow one queueWhat the columns meanInspect itUpstream source map and remaining simplificationsMatching queue ownership and recovery
Matching owns three distinct persisted resources:
| Resource | Toy table | Purpose |
|---|---|---|
| Physical queue metadata | task_queues |
Fence old owners and remember acknowledged progress |
| Spooled deliveries | tasks |
Hold tasks until Matching has handled their delivery |
| Queue-family configuration | task_queue_user_data |
Store versioning and scheduling configuration independently of queue lifetime |
A physical queue key includes namespace, family, workflow/activity type, partition, and deployment version. Each physical queue has its own range ID and task sequence. Reloading one queue does not take ownership of other versions or the other task type.
Follow one queue
- Load:
Partitions.loadcreates a backlog manager for the physical queue.loadBacklogreads its metadata and conditionally incrementsrange_id. A new queue starts at range 1. - Spool: if an arriving task cannot sync-match a waiting worker, Matching appends a
tasksrow. In one transaction, persistence checks the range ID, allocates the next per-queue task ID, updateslast_task_id, and inserts the task. - Read: the backlog manager reads tasks in task-ID order and tracks a process-local
readLevel. Scheduling may dispatch these tasks out of order. - Dispatch: Matching asks History to record the task start. A retryable failure leaves the delivery pending. A successful start, an already-started task, or obsolete work lets Matching acknowledge that delivery. This acknowledges delivery, not completion of the workflow or activity.
- Checkpoint: the acknowledgement level advances through the fully handled prefix. Persistence checks the range ID, saves
ack_level, and deletes tasks at or below it in one transaction. - Unload/restart: loaded managers, read levels, and out-of-order acknowledgements disappear. Metadata and tasks remain. The next manager increments the range and resumes reading above the saved acknowledgement level.
Sync matching can skip the tasks insert entirely. A queue metadata row can therefore exist with ack_level = last_task_id = 0 even after it delivered work. An empty tasks table also does not imply completed workflows: History owns their execution state.
What the columns mean
| Column/state | Durable? | Meaning |
|---|---|---|
range_id |
Yes | Monotonically increasing ownership fence for this physical queue. An old range cannot append, checkpoint, or delete tasks after takeover. |
ack_level |
Yes | Highest task ID through which all deliveries have been handled. |
last_task_id |
Yes | Highest allocated task ID. This toy allocates consecutive IDs during each append. |
readLevel |
No | Highest task ID read by the currently loaded manager. |
| Outstanding acknowledgements | No | Which read tasks completed dispatch, including completions above a gap. |
Partition generation |
No | Identifies a manager instance within this process. It is separate from the persisted range ID. |
The range is a fencing token, not a timed lease. Unloading does not delete metadata or reset the range. Persisted metadata alone also cannot tell you whether a manager is currently loaded.
Why a dispatched task can remain in tasks
Suppose tasks 1, 2, and 3 were read. Task 2 dispatches first:
Read level: 3
Completed dispatch: {2}
Saved acknowledgement level: 0
Persisted task rows: 1, 2, 3
Task 1 then dispatches. The prefix through 2 is complete, so the saved acknowledgement advances to 2 and rows 1 and 2 are deleted. Task 3 remains.
If the process crashes after dispatching 2 but before handling 1, the next owner rereads 2. History rejects its repeated start if it has already started or become obsolete. If a started task was never received by a worker, History's timeout processing schedules another attempt. This is why replaying a delivery is safe without claiming exactly-once delivery.
The toy's Matching.backlog, backlog gauges, and SQL counts report stored rows. They include dispatched rows retained above acknowledgement gaps. They are not an exact count of tasks still awaiting dispatch. The lifecycle walkthrough separately prints the loaded manager's outstanding count.
Inspect it
Run pnpm matching:lifecycle for the memory-backed walkthrough. Its output includes persistedQueue and loadedQueueProgress before and after unload/reload.
For PostgreSQL, run the normal server and stock SDK demo, or pnpm recovery with other toy servers stopped. Startup applies migrations automatically. Migration 003 adds queue ownership metadata; migration 004 carries History version directives and relocates old auto backlog to default storage, allocating fresh destination IDs without collisions. Stop the old server before upgrading. Refresh your database client's table list to see task_queues.
SELECT q.physical_queue,
q.range_id, q.ack_level, q.last_task_id,
count(t.task_id) AS stored_task_rows
FROM task_queues q
LEFT JOIN tasks t USING (physical_queue)
GROUP BY q.physical_queue
ORDER BY q.physical_queue;
SELECT physical_queue, task_id,
data->'task'->>'runId' AS run_id,
data->'task'->>'kind' AS kind,
data->'task'->'versionDirective' AS version_directive
FROM tasks
ORDER BY physical_queue, task_id;
Restart the server and poll the same queue again: range_id increases while the saved acknowledgement remains intact. A workflow and an activity queue with the same family name have separate rows. pnpm db:inspect WORKFLOW_ID also prints both Matching tables.
Upstream source map and remaining simplifications
| Toy component | Upstream source |
|---|---|
| Partition-owned physical queue lifetimes | task_queue_partition_manager.go, physical_task_queue_manager.go |
| Conditional acquisition and queue writes | db.go, common/persistence/sql/task_queues.go |
| Ordered read/acknowledgement levels | ack_manager.go, backlog_manager.go |
| Allocating task IDs | task_writer.go |
| Newer priority/fairness backlog storage | pri_backlog_manager.go, fair_backlog_manager.go |
This is the compact ordered-acknowledgement model; upstream also has priority subqueues and fairness progress structures. We do not create task_queues_v2 / tasks_v2 tables without implementing their semantics.
Other deliberate differences:
- Upstream's older task writer allocates task-ID blocks from ranges, which can leave unused ID gaps. The toy persists a consecutive counter per queue; range IDs only fence ownership.
- Upstream can checkpoint periodically and garbage-collect task rows separately. The toy checkpoints and deletes the prefix together after each handled delivery.
- The toy reads the full physical backlog instead of paginating/prefetching it. Dispatch operations are serialized within one Matching service.
- Sticky queue expiry, queue retention/GC, persisted fairness passes, and automatic distributed placement are not implemented. Queue rows remain for inspection.
- Fencing is enforced across independent database connections. Membership routes each partition to one Matching host, and a host that loses a partition unloads it. Membership is static: two processes configured with different host lists for the same queues can still take the queue from each other repeatedly.
- Migration 004 and memory restore upgrade legacy deployment routing while preserving tasks and fencing affected queues. Older exported memory snapshots using
deliveryIdare not an import format for the current schema.
test/matching-persistence.test.ts covers unload, takeover, acknowledgement gaps, and replay after History accepted a start. pnpm test:postgres runs the ownership contract across independent pools, checks migration of old backlog rows, and retains the existing full-process SIGKILL recovery checks.
Queue-family user data has its own expected-version comparison in persistence. Queue range fences protect backlog ownership; they do not protect task_queue_user_data. Independent manager tests verify both conditional insertion and update, with a typed conflict for the loser.