Skip to content

Commit 76b7643

Browse files
committed
updates to dashboard and proxy for DLQ redrive and such
1 parent b4041ac commit 76b7643

8 files changed

Lines changed: 105 additions & 0 deletions

File tree

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] : [])),

packages/proxy/src/types.ts

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -39,6 +39,8 @@ export type HttpFindJobsOptions = Omit<types.FindJobsOptions, 'db' | 'data'> & {
3939
data?: HttpJsonRecord
4040
}
4141

42+
export type HttpRedriveOptions = Omit<types.RedriveOptions, 'db'>
43+
4244
export type HttpInsertOptions = Omit<types.InsertOptions, 'db'>
4345

4446
export type HttpCompleteOptions = Omit<types.CompleteOptions, 'db'>
@@ -198,6 +200,16 @@ export type HttpDeleteJobRequest = HttpCancelRequest
198200

199201
export type HttpDeleteJobResponse = HttpCancelResponse
200202

203+
export type HttpRedriveRequest = {
204+
name: HttpQueueName
205+
options?: HttpRedriveOptions
206+
}
207+
208+
export type HttpRedriveResponse = {
209+
ok: true
210+
result: number
211+
}
212+
201213
export type HttpDeleteQueuedJobsRequest = {
202214
name: HttpQueueName
203215
}

packages/proxy/test/apiCallsTest.ts

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -440,6 +440,11 @@ describe('proxy api routes', () => {
440440
body: { name: 'queue', id: '1' },
441441
expected: ['queue', '1']
442442
},
443+
{
444+
method: 'redrive',
445+
body: { name: 'dlq', options: { destination: 'dest', sourceName: 'src', limit: 50 } },
446+
expected: ['dlq', { destination: 'dest', sourceName: 'src', limit: 50 }]
447+
},
443448
{
444449
method: 'deleteQueuedJobs',
445450
body: { name: 'queue' },

packages/proxy/test/typeDriftTest.ts

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -52,6 +52,10 @@ describe('Zod schema / HTTP type key drift', () => {
5252
expect(assertKeysMatch<AssertKeysMatch<SchemaOutput<typeof contracts.findJobsOptionsSchema>, httpTypes.HttpFindJobsOptions>>()).toBe(true)
5353
})
5454

55+
it('redriveOptionsSchema matches HttpRedriveOptions', () => {
56+
expect(assertKeysMatch<AssertKeysMatch<SchemaOutput<typeof contracts.redriveOptionsSchema>, httpTypes.HttpRedriveOptions>>()).toBe(true)
57+
})
58+
5559
it('insertOptionsSchema matches HttpInsertOptions', () => {
5660
expect(assertKeysMatch<AssertKeysMatch<SchemaOutput<typeof contracts.insertOptionsSchema>, httpTypes.HttpInsertOptions>>()).toBe(true)
5761
})

src/index.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -496,6 +496,7 @@ export type {
496496
QueueOptions,
497497
QueuePolicy,
498498
QueueResult,
499+
RedriveOptions,
499500
Request,
500501
Schedule,
501502
ScheduleOptions,

0 commit comments

Comments
 (0)