Skip to content

Commit 61df305

Browse files
authored
Merge pull request #827 from timgit/dlq-redrive
DLQ redrive
2 parents 3f1980d + 76b7643 commit 61df305

21 files changed

Lines changed: 458 additions & 17 deletions

README.md

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -47,7 +47,7 @@ This will likely cater the most to teams already familiar with the simplicity of
4747
* Job dependency workflow orchestration
4848
* Cron scheduling, job deferral
4949
* Queue storage policies to support a variety of rate limiting, debouncing, and concurrency use cases
50-
* Priority queues, dead letter queues, automatic retries with exponential backoff
50+
* Priority queues, dead letter queues with redrive, automatic retries with exponential backoff
5151
* Pub/sub API for fan-out queue relationships
5252
* SQL support for non-Node.js runtimes for most operations
5353
* Serverless function compatible

docs/api/jobs.md

Lines changed: 42 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -319,10 +319,20 @@ Returns an array of jobs from a queue
319319
blocking: boolean,
320320
deadLetter: string,
321321
policy: string,
322-
output: object
322+
output: object,
323+
sourceName: string | null,
324+
sourceId: string | null,
325+
sourceCreatedOn: Date | null,
326+
sourceRetryCount: number | null
323327
}
324328
```
325329

330+
When a job is moved into a dead letter queue, the `source*` fields record where it
331+
came from: `sourceName` is the queue it originally failed on, `sourceId` is the id
332+
of the original job, `sourceCreatedOn` is the original job's creation time (so its
333+
true age survives the move), and `sourceRetryCount` is how many retries it consumed
334+
before being dead-lettered. These are `null` for jobs that were not dead-lettered.
335+
326336

327337
**Notes**
328338

@@ -354,6 +364,37 @@ Deletes a job by id.
354364

355365
Deletes a set of jobs by id.
356366

367+
### `redrive(name, options)`
368+
369+
Moves jobs out of a dead letter queue and re-creates them as fresh jobs on their original source queue. `name` is the
370+
dead letter queue to drain. Returns the number of jobs moved.
371+
372+
Each job is routed back to the queue it originally failed on (its `sourceName`),
373+
so a single dead letter queue that collects from many source queues fans back out
374+
correctly. Re-created jobs get a new id, a reset retry count, cleared output, and
375+
the destination queue's current retry, retention, and policy configuration. Only
376+
jobs that are not currently being processed (still in the `created`/`retry` state)
377+
are moved.
378+
379+
`options`:
380+
381+
- `destination`override queue to move all matched jobs into, instead of each
382+
job's original source queue. Required to redrive jobs that have no recorded
383+
source queue (e.g. jobs dead-lettered before this feature existed); such jobs are
384+
left in place otherwise.
385+
- `sourceName`only redrive jobs that originated from this source queue.
386+
- `limit`maximum number of jobs to move in this call, oldest first (default
387+
`1000`). Loop or schedule repeated calls to drain large dead letter queues at a
388+
controlled rate.
389+
390+
```js
391+
// drain a dead letter queue back to its source queues, 500 at a time
392+
let moved
393+
do {
394+
moved = await boss.redrive('email-dlq', { limit: 500 })
395+
} while (moved > 0)
396+
```
397+
357398
### `deleteQueuedJobs(name)`
358399

359400
Deletes all queued jobs in a queue.

docs/api/queues.md

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -41,7 +41,7 @@ Allowed policy values:
4141
4242
* **deadLetter**, string
4343
44-
When a job fails after all retries, if the queue has a `deadLetter` property, the job's payload will be copied into that queue, copying the same retention and retry configuration as the original job.
44+
When a job fails after all retries, if the queue has a `deadLetter` property, the job's payload will be copied into that queue, copying the same retention and retry configuration as the original job. The dead-lettered job also records where it came from via the `sourceName`, `sourceId`, `sourceCreatedOn`, and `sourceRetryCount` fields, which power [`redrive()`](jobs#redrivename-options) for moving jobs back to their source queue.
4545
4646
* **warningQueueSize**, int
4747

docs/index.md

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -59,7 +59,7 @@ This will likely cater the most to teams already familiar with the simplicity of
5959
* Job dependency workflow orchestration
6060
* Cron scheduling, job deferral
6161
* Queue storage policies to support a variety of rate limiting, debouncing, and concurrency use cases
62-
* Priority queues, dead letter queues, automatic retries with exponential backoff
62+
* Priority queues, dead letter queues with redrive, automatic retries with exponential backoff
6363
* Pub/sub API for fan-out queue relationships
6464
* SQL support for non-Node.js runtimes for most operations
6565
* Serverless function compatible

package-lock.json

Lines changed: 2 additions & 2 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

package.json

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
{
22
"name": "pg-boss",
3-
"version": "12.22.0",
3+
"version": "12.23.0",
44
"description": "Queueing jobs in Postgres from Node.js like a boss",
55
"type": "module",
66
"main": "./dist/index.js",
@@ -61,7 +61,7 @@
6161
"docs:readme": "node ./scripts/sync-readme.js"
6262
},
6363
"pgboss": {
64-
"schema": 33
64+
"schema": 34
6565
},
6666
"repository": {
6767
"type": "git",

packages/dashboard/app/routes/queues.$name.jobs.$jobId.tsx

Lines changed: 27 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -273,6 +273,33 @@ export default function JobDetail ({ loaderData }: Route.ComponentProps) {
273273
</div>
274274
)}
275275

276+
{(job.sourceName || job.sourceId) && (
277+
<div className="grid grid-cols-1 sm:grid-cols-2 lg:grid-cols-4 gap-x-6 gap-y-4">
278+
{job.sourceName && (
279+
<div>
280+
<dt className="pgb-eyebrow">Source Queue</dt>
281+
<dd className="mt-1 text-sm text-[var(--text-primary)]">
282+
<DbLink
283+
to={`/queues/${job.sourceName}`}
284+
className="font-mono text-xs text-primary-600 hover:text-primary-700 dark:text-primary-400 dark:hover:text-primary-300"
285+
>
286+
{job.sourceName}
287+
</DbLink>
288+
</dd>
289+
</div>
290+
)}
291+
{job.sourceId && (
292+
<ConfigItem label="Source Job ID" value={job.sourceId} mono />
293+
)}
294+
{job.sourceRetryCount !== null && job.sourceRetryCount !== undefined && (
295+
<ConfigItem label="Source Retry Count" value={job.sourceRetryCount} mono />
296+
)}
297+
{job.sourceCreatedOn && (
298+
<ConfigItem label="Source Created" value={formatDate(new Date(job.sourceCreatedOn))} mono />
299+
)}
300+
</div>
301+
)}
302+
276303
<div className="grid grid-cols-1 sm:grid-cols-2 lg:grid-cols-5 gap-x-6 gap-y-4">
277304
<ConfigItem
278305
label="Created"

packages/dashboard/tests/server/queries.test.ts

Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -757,6 +757,39 @@ describe('Job Queries', () => {
757757
expect(job!.priority).toBe(5)
758758
expect(job!.data).toEqual({ foo: 'bar' })
759759
})
760+
761+
it('returns null source-tracking fields for a job that was not dead-lettered', async () => {
762+
await createTestQueue('test-queue')
763+
const jobId = await sendTestJob('test-queue', { foo: 'bar' })
764+
765+
const job = await getJobById(ctx.connectionString, ctx.schema, 'test-queue', jobId!)
766+
767+
expect(job!.sourceName).toBeNull()
768+
expect(job!.sourceId).toBeNull()
769+
expect(job!.sourceCreatedOn).toBeNull()
770+
expect(job!.sourceRetryCount).toBeNull()
771+
})
772+
773+
it('exposes source-tracking fields for a dead-lettered job', async () => {
774+
await createTestQueue('test-dlq')
775+
await createTestQueue('test-queue', { deadLetter: 'test-dlq', retryLimit: 0 })
776+
777+
const jobId = await sendTestJob('test-queue', { foo: 'bar' })
778+
779+
// fail() routes the job to the dead letter queue synchronously
780+
await fetchTestJob('test-queue')
781+
await failTestJob('test-queue', jobId!)
782+
783+
// The dead-lettered copy is a new job in the DLQ with its own id
784+
const moved = await fetchTestJob('test-dlq')
785+
const dlqJob = await getJobById(ctx.connectionString, ctx.schema, 'test-dlq', moved!.id)
786+
787+
expect(dlqJob).not.toBeNull()
788+
expect(dlqJob!.sourceName).toBe('test-queue')
789+
expect(dlqJob!.sourceId).toBe(jobId)
790+
expect(dlqJob!.sourceCreatedOn).toBeTruthy()
791+
expect(dlqJob!.sourceRetryCount).toBe(0)
792+
})
760793
})
761794
})
762795

packages/proxy/src/contracts.ts

Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -93,6 +93,12 @@ export const findJobsOptionsSchema = z.object({
9393
queued: z.boolean().optional(),
9494
}) satisfies z.ZodType<types.HttpFindJobsOptions>
9595

96+
export const redriveOptionsSchema = z.object({
97+
destination: queueNameSchema.optional(),
98+
sourceName: queueNameSchema.optional(),
99+
limit: z.number().optional(),
100+
}) satisfies z.ZodType<types.HttpRedriveOptions>
101+
96102
export const insertOptionsSchema = z.object({
97103
returnId: z.boolean().optional(),
98104
}) satisfies z.ZodType<types.HttpInsertOptions>
@@ -162,6 +168,10 @@ export const jobWithMetadataSchema = jobSchemaBase.extend({
162168
policy: z.string(),
163169
deadLetter: z.string(),
164170
output: jsonRecordSchema,
171+
sourceName: z.string().nullable(),
172+
sourceId: z.string().nullable(),
173+
sourceCreatedOn: z.iso.datetime().nullable().transform((val) => val ? new Date(val) : null),
174+
sourceRetryCount: z.number().nullable(),
165175
}) satisfies z.ZodType<types.HttpJobWithMetadata>
166176

167177
export const commandResponseSchema = z.object({
@@ -380,6 +390,16 @@ export const deleteJobResponseSchema: z.ZodType<types.HttpDeleteJobResponse> = z
380390
result: commandResponseSchema
381391
})
382392

393+
export const redriveRequestSchema: z.ZodType<types.HttpRedriveRequest> = z.object({
394+
name: queueNameSchema,
395+
options: redriveOptionsSchema.optional()
396+
})
397+
398+
export const redriveResponseSchema: z.ZodType<types.HttpRedriveResponse> = z.object({
399+
ok: z.literal(true),
400+
result: z.number()
401+
})
402+
383403
export const deleteQueuedJobsRequestSchema: z.ZodType<types.HttpDeleteQueuedJobsRequest> = z.object({
384404
name: queueNameSchema
385405
})

packages/proxy/src/routes.ts

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -35,6 +35,8 @@ import {
3535
isInstalledResponseSchema,
3636
publishRequestSchema,
3737
publishResponseSchema,
38+
redriveRequestSchema,
39+
redriveResponseSchema,
3840
resumeRequestSchema,
3941
resumeResponseSchema,
4042
retryRequestSchema,
@@ -177,6 +179,7 @@ export const postMethods: RouteEntry[] = [
177179
post('jobs', 'resume', resumeRequestSchema, resumeResponseSchema, (body) => [body.name, body.id]),
178180
post('jobs', 'retry', retryRequestSchema, retryResponseSchema, (body) => [body.name, body.id]),
179181
post('jobs', 'deleteJob', deleteJobRequestSchema, deleteJobResponseSchema, (body) => [body.name, body.id]),
182+
post('jobs', 'redrive', redriveRequestSchema, redriveResponseSchema, (body) => withOptionalOptions([body.name], body.options)),
180183
post('jobs', 'deleteQueuedJobs', deleteQueuedJobsRequestSchema, deleteQueuedJobsResponseSchema, (body) => [body.name]),
181184
post('jobs', 'deleteStoredJobs', deleteStoredJobsRequestSchema, deleteStoredJobsResponseSchema, (body) => [body.name]),
182185
post('jobs', 'deleteAllJobs', deleteAllJobsRequestSchema, deleteAllJobsResponseSchema, (body) => (body.name ? [body.name] : [])),

0 commit comments

Comments
 (0)