Skip to content

Commit 65f1fb9

Browse files
committed
dealy seconds
1 parent 54b9fc5 commit 65f1fb9

21 files changed

Lines changed: 103 additions & 61 deletions

src/adapters/mongodb.ts

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -60,10 +60,10 @@ export class MongoDatabaseAdapter implements DatabaseAdapter {
6060
const res = await this.col.insertOne({
6161
payload,
6262
ttr: meta.ttr ?? 300,
63-
delay: meta.delay ?? 0,
63+
delaySeconds: meta.delaySeconds ?? 0,
6464
priority: meta.priority ?? 0,
6565
pushTime: now,
66-
delayTime: meta.delay ? new Date(now.getTime() + meta.delay * 1000) : null,
66+
delayTime: meta.delaySeconds ? new Date(now.getTime() + meta.delaySeconds * 1000) : null,
6767
status: 'waiting',
6868
attempt: 0
6969
});
@@ -115,7 +115,7 @@ export class MongoDatabaseAdapter implements DatabaseAdapter {
115115
payload: Buffer.from(doc.payload?.buffer || doc.payload || []), // Handle Binary/Buffer/empty
116116
meta: {
117117
ttr: doc.ttr,
118-
delay: doc.delay,
118+
delaySeconds: doc.delaySeconds,
119119
priority: doc.priority,
120120
pushedAt: doc.pushTime,
121121
reservedAt: now

src/adapters/redis.ts

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -77,10 +77,10 @@ export class RedisDatabaseAdapter implements DatabaseAdapter {
7777
id: jobId,
7878
payload: payload.toString('base64'),
7979
ttr: (meta.ttr || 300).toString(),
80-
delay: (meta.delay || 0).toString(),
80+
delay_seconds: (meta.delaySeconds || 0).toString(),
8181
priority: (meta.priority || 0).toString(),
8282
push_time: now.toString(),
83-
delay_time: meta.delay ? (now + meta.delay * 1000).toString() : '',
83+
delay_time: meta.delaySeconds ? (now + meta.delaySeconds * 1000).toString() : '',
8484
status: 'waiting',
8585
attempt: '0'
8686
};
@@ -89,9 +89,9 @@ export class RedisDatabaseAdapter implements DatabaseAdapter {
8989
await this.client.hSet(this.getJobKey(jobId), jobData);
9090

9191
// Add to waiting queue based on delay
92-
if (meta.delay && meta.delay > 0) {
92+
if (meta.delaySeconds && meta.delaySeconds > 0) {
9393
// Add to delayed set with execution time as score
94-
const executeAt = Math.floor((now + meta.delay * 1000) / 1000);
94+
const executeAt = Math.floor((now + meta.delaySeconds * 1000) / 1000);
9595
await this.client.zAdd(`${this.keyPrefix}:delayed`, {
9696
score: executeAt,
9797
value: jobId
@@ -157,7 +157,7 @@ export class RedisDatabaseAdapter implements DatabaseAdapter {
157157
payload: Buffer.from(jobData.payload, 'base64'),
158158
meta: {
159159
ttr: parseInt(jobData.ttr || '300'),
160-
delay: parseInt(jobData.delay || '0'),
160+
delaySeconds: parseInt(jobData.delay_seconds || '0'),
161161
priority: parseInt(jobData.priority || '0'),
162162
pushedAt: new Date(parseInt(jobData.push_time || '0')),
163163
reservedAt: new Date(now * 1000)

src/adapters/sqlite.ts

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -33,7 +33,7 @@ export class SQLiteDatabaseAdapter implements DatabaseAdapter {
3333
id INTEGER PRIMARY KEY AUTOINCREMENT,
3434
payload BLOB NOT NULL,
3535
ttr INTEGER DEFAULT 300,
36-
delay INTEGER DEFAULT 0,
36+
delay_seconds INTEGER DEFAULT 0,
3737
priority INTEGER DEFAULT 0,
3838
push_time INTEGER NOT NULL,
3939
delay_time INTEGER,
@@ -62,18 +62,18 @@ export class SQLiteDatabaseAdapter implements DatabaseAdapter {
6262
const now = new Date();
6363
const stmt = this.db.prepare(`
6464
INSERT INTO jobs (
65-
payload, ttr, delay, priority, push_time,
65+
payload, ttr, delay_seconds, priority, push_time,
6666
delay_time, status
6767
) VALUES (?, ?, ?, ?, ?, ?, ?)
6868
`);
6969

7070
const result = stmt.run(
7171
payload,
7272
meta.ttr || 300,
73-
meta.delay || 0,
73+
meta.delaySeconds || 0,
7474
meta.priority || 0,
7575
now.getTime(),
76-
meta.delay ? now.getTime() + meta.delay * 1000 : null,
76+
meta.delaySeconds ? now.getTime() + meta.delaySeconds * 1000 : null,
7777
'waiting'
7878
);
7979

@@ -143,7 +143,7 @@ export class SQLiteDatabaseAdapter implements DatabaseAdapter {
143143
payload: job.payload,
144144
meta: {
145145
ttr: job.ttr,
146-
delay: job.delay,
146+
delaySeconds: job.delay_seconds,
147147
priority: job.priority,
148148
pushedAt: new Date(job.push_time),
149149
reservedAt: new Date(now)

src/core/queue.ts

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -97,7 +97,7 @@ export abstract class Queue<TJobMap = Record<string, any>, TJobRequest extends B
9797
* await queue.addJob('backup', {
9898
* payload: { path: '/data' },
9999
* ttr: 3600,
100-
* delay: 60
100+
* delaySeconds: 60
101101
* });
102102
* ```
103103
*/
@@ -109,7 +109,7 @@ export abstract class Queue<TJobMap = Record<string, any>, TJobRequest extends B
109109

110110
const meta: JobMeta = {
111111
ttr: options.ttr ?? this.ttrDefault,
112-
delay: (options as any).delay ?? 0,
112+
delaySeconds: (options as any).delaySeconds ?? 0,
113113
priority: (options as any).priority ?? 0,
114114
pushedAt: new Date()
115115
};
@@ -423,7 +423,7 @@ export abstract class Queue<TJobMap = Record<string, any>, TJobRequest extends B
423423
* Pushes a new message to the queue storage backend.
424424
*
425425
* @param payload - Serialized job data as string
426-
* @param meta - Job metadata including TTR, delay, priority
426+
* @param meta - Job metadata including TTR, delaySeconds, priority
427427
* @returns Promise resolving to unique job ID
428428
* @protected
429429
* @abstract

src/drivers/db.ts

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -26,9 +26,14 @@ export class DbQueue<TJobMap = Record<string, any>> extends Queue<TJobMap, DbJob
2626
return null;
2727
}
2828

29+
const payload = record.payload.toString();
30+
// Extract job name from payload
31+
const jobData = JSON.parse(payload);
32+
2933
return {
3034
id: record.id,
31-
payload: record.payload.toString(),
35+
name: jobData.name,
36+
payload,
3237
meta: record.meta
3338
};
3439
}

src/drivers/file.ts

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -66,7 +66,7 @@ export class FileQueue<TJobMap = Record<string, any>> extends Queue<TJobMap, Fil
6666
}
6767

6868
const ttr = meta.ttr ?? 300;
69-
const delay = meta.delay ?? 0;
69+
const delay = meta.delaySeconds ?? 0;
7070

7171
if (delay === 0) {
7272
data.waiting.push([jobId, ttr]);
@@ -131,8 +131,13 @@ export class FileQueue<TJobMap = Record<string, any>> extends Queue<TJobMap, Fil
131131
if (reserved) {
132132
const jobPath = path.join(this.path, `job${reserved.id}.data`);
133133
const payload = await fs.readFile(jobPath, 'utf8');
134+
135+
// Extract job name from payload
136+
const jobData = JSON.parse(payload);
137+
134138
return {
135139
id: reserved.id,
140+
name: jobData.name,
136141
payload,
137142
meta: {
138143
ttr: reserved.ttr

src/drivers/memory.ts

Lines changed: 8 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -59,7 +59,7 @@ export interface InMemoryQueueOptions extends QueueOptions {
5959
* await queue.addJob('my-job', {
6060
* payload: { data: 'test' },
6161
* priority: 5,
62-
* delay: 10
62+
* delaySeconds: 10
6363
* });
6464
*
6565
* await queue.run(true, 1);
@@ -94,15 +94,15 @@ export class InMemoryQueue<TJobMap = Record<string, any>> extends Queue<TJobMap,
9494
this.jobs.set(id, job);
9595

9696
// Handle delay
97-
if (meta.delay && meta.delay > 0) {
98-
job.delayTime = Date.now() + (meta.delay * 1000);
97+
if (meta.delaySeconds && meta.delaySeconds > 0) {
98+
job.delayTime = Date.now() + (meta.delaySeconds * 1000);
9999
job.status = 'waiting';
100100

101101
// Schedule job to become available after delay
102102
const timeout = setTimeout(() => {
103103
this.delayedJobs.delete(id);
104104
this.addToWaitingQueue(id);
105-
}, meta.delay * 1000);
105+
}, meta.delaySeconds * 1000);
106106

107107
this.delayedJobs.set(id, timeout);
108108
} else {
@@ -147,8 +147,12 @@ export class InMemoryQueue<TJobMap = Record<string, any>> extends Queue<TJobMap,
147147

148148
this.ttrTimeouts.set(jobId, ttrTimeout);
149149

150+
// Extract job name from payload
151+
const jobData = JSON.parse(job.payload);
152+
150153
return {
151154
id: jobId,
155+
name: jobData.name,
152156
payload: job.payload,
153157
meta: job.meta
154158
};

src/drivers/redis.ts

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -146,9 +146,9 @@ export class RedisQueue<TJobMap = Record<string, any>> extends Queue<TJobMap, Re
146146
const message = `${ttr};${payload}`;
147147
await this.redis.hset(this.messagesKey, id, message);
148148

149-
if (meta.delay && meta.delay > 0) {
149+
if (meta.delaySeconds && meta.delaySeconds > 0) {
150150
// Add to delayed set with execution time as score
151-
const executeAt = now + meta.delay;
151+
const executeAt = now + meta.delaySeconds;
152152
await this.redis.zadd(this.delayedKey, executeAt, id);
153153
} else {
154154
// Add to waiting list (FIFO queue)
@@ -207,8 +207,12 @@ export class RedisQueue<TJobMap = Record<string, any>> extends Queue<TJobMap, Re
207207
// Increment attempt counter
208208
await this.redis.hincrby(this.attemptsKey, id, 1);
209209

210+
// Extract job name from payload
211+
const jobData = JSON.parse(payloadStr);
212+
210213
return {
211214
id,
215+
name: jobData.name,
212216
payload: payloadStr,
213217
meta: {
214218
ttr,

src/drivers/sqs.ts

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -61,7 +61,7 @@ export class SqsQueue<TJobMap = Record<string, any>> extends Queue<
6161
const command = new SendMessageCommand({
6262
QueueUrl: this.queueUrl,
6363
MessageBody: payload,
64-
DelaySeconds: meta.delay || 0,
64+
DelaySeconds: meta.delaySeconds || 0,
6565
MessageAttributes: messageAttributes,
6666
});
6767

@@ -116,8 +116,12 @@ export class SqsQueue<TJobMap = Record<string, any>> extends Queue<
116116
await this.client.send(visibilityCommand);
117117
}
118118

119+
// Extract job name from payload
120+
const jobData = JSON.parse(payload);
121+
119122
return {
120123
id: message.MessageId,
124+
name: jobData.name,
121125
payload,
122126
meta: {
123127
...meta,

src/interfaces/job.ts

Lines changed: 8 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@ export interface JobMeta {
66
/** Time to run - number of seconds to run the job */
77
ttr?: number;
88
/** Number of seconds to delay job execution from now */
9-
delay?: number;
9+
delaySeconds?: number;
1010
/** Job priority - higher numbers = higher priority (processed first) */
1111
priority?: number;
1212
/** Job pushed at */
@@ -27,7 +27,7 @@ export interface JobContext<T> {
2727
id: string;
2828
/** Job payload */
2929
payload: T;
30-
/** Job meta: ttr, delay, priority, pushedAt, reservedAt, doneAt, receiptHandle */
30+
/** Job meta: ttr, delaySeconds, priority, pushedAt, reservedAt, doneAt, receiptHandle */
3131
meta: JobMeta;
3232
/** Job pushed at */
3333
pushedAt?: Date;
@@ -54,7 +54,7 @@ export interface QueueMessage {
5454
name: string;
5555
/** Job payload */
5656
payload: string;
57-
/** Job meta: ttr, delay, priority, pushedAt, reservedAt, doneAt, receiptHandle */
57+
/** Job meta: ttr, delaySeconds, priority, pushedAt, reservedAt, doneAt, receiptHandle */
5858
meta: JobMeta;
5959
}
6060

@@ -81,7 +81,7 @@ export interface BaseJobOptions {
8181
// Full options interface (for internal use)
8282
export interface JobOptions extends BaseJobOptions {
8383
/** Number of seconds to delay job execution from now */
84-
delay?: number;
84+
delaySeconds?: number;
8585
/** Job priority - higher numbers = higher priority (processed first) */
8686
priority?: number;
8787
}
@@ -91,24 +91,24 @@ export interface DbJobOptions extends BaseJobOptions {
9191
// DB adapters may or may not support delay/priority - we allow them for flexibility
9292
// The specific DatabaseAdapter implementation determines actual support
9393
/** Number of seconds to delay job execution from now. Support varies by database adapter. */
94-
delay?: number;
94+
delaySeconds?: number;
9595
/** Job priority - higher numbers = higher priority. Support varies by database adapter. */
9696
priority?: number;
9797
}
9898

9999
export interface SqsJobOptions extends BaseJobOptions {
100100
/** Number of seconds to delay job execution from now (0-900 seconds max for SQS) */
101-
delay?: number;
101+
delaySeconds?: number;
102102
}
103103

104104
export interface FileJobOptions extends BaseJobOptions {
105105
/** Number of seconds to delay job execution from now */
106-
delay?: number;
106+
delaySeconds?: number;
107107
}
108108

109109
export interface InMemoryJobOptions extends BaseJobOptions {
110110
/** Number of seconds to delay job execution from now */
111-
delay?: number;
111+
delaySeconds?: number;
112112
/** Job priority - higher numbers = higher priority (processed first) */
113113
priority?: number;
114114
}

0 commit comments

Comments
 (0)