Skip to content

Commit 91dda28

Browse files
committed
feat: Add the ability to delay a dequeued job
1 parent c1af40b commit 91dda28

12 files changed

Lines changed: 485 additions & 11 deletions

.gitignore

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -128,3 +128,4 @@ dist
128128
.yarn/build-state.yml
129129
.yarn/install-state.gz
130130
.pnp.*
131+
.claude
Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,2 @@
1+
ALTER TABLE `tasks` ADD `availableAt` integer;--> statement-breakpoint
2+
CREATE INDEX `tasks_available_at_idx` ON `tasks` (`availableAt`);
Lines changed: 158 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,158 @@
1+
{
2+
"version": "6",
3+
"dialect": "sqlite",
4+
"id": "3cf1b002-16e9-4bf1-a137-e7fc10fee5b2",
5+
"prevId": "563ac9a6-ff12-4aae-9fe3-f72bf3251c1e",
6+
"tables": {
7+
"tasks": {
8+
"name": "tasks",
9+
"columns": {
10+
"id": {
11+
"name": "id",
12+
"type": "integer",
13+
"primaryKey": true,
14+
"notNull": true,
15+
"autoincrement": true
16+
},
17+
"queue": {
18+
"name": "queue",
19+
"type": "text",
20+
"primaryKey": false,
21+
"notNull": true,
22+
"autoincrement": false
23+
},
24+
"payload": {
25+
"name": "payload",
26+
"type": "text",
27+
"primaryKey": false,
28+
"notNull": true,
29+
"autoincrement": false
30+
},
31+
"createdAt": {
32+
"name": "createdAt",
33+
"type": "integer",
34+
"primaryKey": false,
35+
"notNull": true,
36+
"autoincrement": false
37+
},
38+
"availableAt": {
39+
"name": "availableAt",
40+
"type": "integer",
41+
"primaryKey": false,
42+
"notNull": false,
43+
"autoincrement": false
44+
},
45+
"status": {
46+
"name": "status",
47+
"type": "text",
48+
"primaryKey": false,
49+
"notNull": true,
50+
"autoincrement": false,
51+
"default": "'pending'"
52+
},
53+
"expireAt": {
54+
"name": "expireAt",
55+
"type": "integer",
56+
"primaryKey": false,
57+
"notNull": false,
58+
"autoincrement": false
59+
},
60+
"allocationId": {
61+
"name": "allocationId",
62+
"type": "text",
63+
"primaryKey": false,
64+
"notNull": true,
65+
"autoincrement": false
66+
},
67+
"numRunsLeft": {
68+
"name": "numRunsLeft",
69+
"type": "integer",
70+
"primaryKey": false,
71+
"notNull": true,
72+
"autoincrement": false
73+
},
74+
"maxNumRuns": {
75+
"name": "maxNumRuns",
76+
"type": "integer",
77+
"primaryKey": false,
78+
"notNull": true,
79+
"autoincrement": false
80+
},
81+
"idempotencyKey": {
82+
"name": "idempotencyKey",
83+
"type": "text",
84+
"primaryKey": false,
85+
"notNull": false,
86+
"autoincrement": false
87+
},
88+
"priority": {
89+
"name": "priority",
90+
"type": "integer",
91+
"primaryKey": false,
92+
"notNull": true,
93+
"autoincrement": false,
94+
"default": 0
95+
}
96+
},
97+
"indexes": {
98+
"tasks_queue_idx": {
99+
"name": "tasks_queue_idx",
100+
"columns": ["queue"],
101+
"isUnique": false
102+
},
103+
"tasks_status_idx": {
104+
"name": "tasks_status_idx",
105+
"columns": ["status"],
106+
"isUnique": false
107+
},
108+
"tasks_expire_at_idx": {
109+
"name": "tasks_expire_at_idx",
110+
"columns": ["expireAt"],
111+
"isUnique": false
112+
},
113+
"tasks_num_runs_left_idx": {
114+
"name": "tasks_num_runs_left_idx",
115+
"columns": ["numRunsLeft"],
116+
"isUnique": false
117+
},
118+
"tasks_max_num_runs_idx": {
119+
"name": "tasks_max_num_runs_idx",
120+
"columns": ["maxNumRuns"],
121+
"isUnique": false
122+
},
123+
"tasks_allocation_id_idx": {
124+
"name": "tasks_allocation_id_idx",
125+
"columns": ["allocationId"],
126+
"isUnique": false
127+
},
128+
"tasks_priority_idx": {
129+
"name": "tasks_priority_idx",
130+
"columns": ["priority"],
131+
"isUnique": false
132+
},
133+
"tasks_available_at_idx": {
134+
"name": "tasks_available_at_idx",
135+
"columns": ["availableAt"],
136+
"isUnique": false
137+
},
138+
"tasks_queue_idempotencyKey_unique": {
139+
"name": "tasks_queue_idempotencyKey_unique",
140+
"columns": ["queue", "idempotencyKey"],
141+
"isUnique": true
142+
}
143+
},
144+
"foreignKeys": {},
145+
"compositePrimaryKeys": {},
146+
"uniqueConstraints": {}
147+
}
148+
},
149+
"enums": {},
150+
"_meta": {
151+
"schemas": {},
152+
"tables": {},
153+
"columns": {}
154+
},
155+
"internal": {
156+
"indexes": {}
157+
}
158+
}

src/drizzle/meta/_journal.json

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,13 @@
2222
"when": 1752135917070,
2323
"tag": "0002_add_priority",
2424
"breakpoints": true
25+
},
26+
{
27+
"idx": 3,
28+
"version": "6",
29+
"when": 1756665292211,
30+
"tag": "0003_add_available_at",
31+
"breakpoints": true
2532
}
2633
]
2734
}

src/index.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,3 +9,4 @@ export type {
99
export { Runner } from "./runner";
1010

1111
export type { DequeuedJob, DequeuedJobError } from "./types";
12+
export { RetryAfterError } from "./types";

src/options.ts

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,8 @@ export interface EnqueueOptions {
1313
numRetries?: number;
1414
idempotencyKey?: string;
1515
priority?: number;
16+
/** delay the job by this many milliseconds. */
17+
delayMs?: number;
1618
}
1719

1820
export interface RunnerFuncs<T> {

src/queue.test.ts

Lines changed: 127 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -428,4 +428,131 @@ describe("SqliteQueue", () => {
428428
failed: 0,
429429
});
430430
});
431+
432+
test("availableAt scheduling", async () => {
433+
const queue = new SqliteQueue<Work>(
434+
"available-at-queue",
435+
buildDBClient(":memory:", { runMigrations: true }),
436+
{
437+
defaultJobArgs: {
438+
numRetries: 0,
439+
},
440+
keepFailedJobs: false,
441+
},
442+
);
443+
444+
// Enqueue job with availableAt in the future
445+
const job = await queue.enqueue({ increment: 1 }, { delayMs: 5000 });
446+
447+
// Job should not be available for dequeue yet
448+
const dequeuedJob = await queue.attemptDequeue({ timeoutSecs: 1 });
449+
expect(dequeuedJob).toBeNull();
450+
451+
{
452+
// Enqueue a new job available now
453+
const enqueuedJob = await queue.enqueue({ increment: 2 });
454+
// Should be able to dequeue the available job
455+
const availableJob = await queue.attemptDequeue({ timeoutSecs: 1 });
456+
expect(availableJob).not.toBeNull();
457+
expect(enqueuedJob!.id).toBe(availableJob!.id);
458+
}
459+
});
460+
461+
test("availableAt jobs become available after the scheduled time", async () => {
462+
const queue = new SqliteQueue<Work>(
463+
"available-later-queue",
464+
buildDBClient(":memory:", { runMigrations: true }),
465+
{
466+
defaultJobArgs: {
467+
numRetries: 0,
468+
},
469+
keepFailedJobs: false,
470+
},
471+
);
472+
473+
// Enqueue job with availableAt very soon
474+
const enqueuedJob = await queue.enqueue(
475+
{ increment: 1 },
476+
{ delayMs: 1000 },
477+
);
478+
479+
{
480+
// Shouldn't be able to dequeue the job yet
481+
const dequeuedJob = await queue.attemptDequeue({ timeoutSecs: 0 });
482+
expect(dequeuedJob).toBeNull();
483+
}
484+
485+
// Wait for the time to pass
486+
await new Promise((resolve) => setTimeout(resolve, 1500));
487+
488+
// Should now be able to dequeue the job
489+
const dequeuedJob = await queue.attemptDequeue({ timeoutSecs: 0 });
490+
expect(dequeuedJob).not.toBeNull();
491+
expect(enqueuedJob!.id).toBe(dequeuedJob!.id);
492+
});
493+
494+
test("finalize with refund retry increments numRunsLeft", async () => {
495+
const queue = new SqliteQueue<Work>(
496+
"refund-retry-queue",
497+
buildDBClient(":memory:", { runMigrations: true }),
498+
{
499+
defaultJobArgs: {
500+
numRetries: 1,
501+
},
502+
keepFailedJobs: true,
503+
},
504+
);
505+
506+
// Enqueue a job that will have 3 total attempts (initial + 2 retries)
507+
const job = await queue.enqueue({ increment: 1 });
508+
expect(job).not.toBeUndefined();
509+
510+
// Dequeue the job
511+
const dequeuedJob = await queue.attemptDequeue({ timeoutSecs: 30 });
512+
expect(dequeuedJob).not.toBeNull();
513+
expect(dequeuedJob!.maxNumRuns).toBe(2); // 1 initial + 1 retries
514+
expect(dequeuedJob!.numRunsLeft).toBe(1);
515+
516+
// Mark as pending_retry without refund - should consume one retry
517+
await queue.finalize(
518+
dequeuedJob!.id,
519+
dequeuedJob!.allocationId,
520+
"pending_retry",
521+
new Date(),
522+
false, // refundRetry = false
523+
);
524+
525+
// Dequeue again to check numRunsLeft
526+
const dequeuedJob2 = await queue.attemptDequeue({ timeoutSecs: 30 });
527+
expect(dequeuedJob2).not.toBeNull();
528+
expect(dequeuedJob2!.numRunsLeft).toBe(0); // One retry consumed
529+
530+
// Mark as pending_retry WITH refund - should NOT consume a retry
531+
await queue.finalize(
532+
dequeuedJob2!.id,
533+
dequeuedJob2!.allocationId,
534+
"pending_retry",
535+
new Date(),
536+
true, // refundRetry = true
537+
);
538+
539+
// Dequeue again to check numRunsLeft
540+
const dequeuedJob3 = await queue.attemptDequeue({ timeoutSecs: 30 });
541+
expect(dequeuedJob3).not.toBeNull();
542+
expect(dequeuedJob3!.numRunsLeft).toBe(0); // Same as before, retry was refunded
543+
544+
// Complete the job
545+
await queue.finalize(
546+
dequeuedJob3!.id,
547+
dequeuedJob3!.allocationId,
548+
"completed",
549+
);
550+
551+
expect(await queue.stats()).toEqual({
552+
pending: 0,
553+
running: 0,
554+
pending_retry: 0,
555+
failed: 0,
556+
});
557+
});
431558
});

src/queue.ts

Lines changed: 17 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
import assert from "node:assert";
2-
import { and, asc, count, eq, gt, lt, or } from "drizzle-orm";
2+
import { and, asc, count, eq, gt, isNull, lt, lte, or, sql } from "drizzle-orm";
33

44
import { buildDBClient } from "./db";
55
import { EnqueueOptions, SqliteQueueOptions } from "./options";
@@ -41,6 +41,7 @@ export class SqliteQueue<T> {
4141
const numRetries =
4242
opts.numRetries ?? this.options.defaultJobArgs.numRetries;
4343
const priority = opts.priority ?? 0;
44+
const availableAt = new Date(Date.now() + (opts.delayMs ?? 0));
4445
const [job] = await this.db
4546
.insert(tasksTable)
4647
.values({
@@ -51,6 +52,7 @@ export class SqliteQueue<T> {
5152
allocationId: generateAllocationId(),
5253
idempotencyKey: opts.idempotencyKey,
5354
priority: priority,
55+
availableAt,
5456
})
5557
.onConflictDoNothing({
5658
target: [tasksTable.queue, tasksTable.idempotencyKey],
@@ -90,6 +92,10 @@ export class SqliteQueue<T> {
9092
and(
9193
eq(tasksTable.queue, this.queueName),
9294
gt(tasksTable.numRunsLeft, 0),
95+
or(
96+
lte(tasksTable.availableAt, new Date()),
97+
isNull(tasksTable.availableAt),
98+
),
9399
or(
94100
// Not picked by a worker yet
95101
eq(tasksTable.status, "pending"),
@@ -143,6 +149,8 @@ export class SqliteQueue<T> {
143149
id: number,
144150
alloctionId: string,
145151
status: "completed" | "pending_retry" | "failed",
152+
availableAt: Date = new Date(),
153+
refundRetry = false,
146154
) {
147155
if (
148156
status === "completed" ||
@@ -156,7 +164,14 @@ export class SqliteQueue<T> {
156164
} else {
157165
await this.db
158166
.update(tasksTable)
159-
.set({ status: status, expireAt: null })
167+
.set({
168+
status: status,
169+
expireAt: null,
170+
availableAt,
171+
numRunsLeft: refundRetry
172+
? sql<number>`${tasksTable.numRunsLeft} + 1`
173+
: sql<number>`${tasksTable.numRunsLeft}`,
174+
})
160175
.where(
161176
and(eq(tasksTable.id, id), eq(tasksTable.allocationId, alloctionId)),
162177
);

0 commit comments

Comments
 (0)