Skip to content

Commit 4c3bc4f

Browse files
committed
make queue name required
1 parent 0b1894c commit 4c3bc4f

25 files changed

Lines changed: 78 additions & 57 deletions

example-queue-project/src/email-queue.ts

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -8,15 +8,15 @@ interface EmailJobs {
88

99
// File-based queue (for local development)
1010
export const emailQueueFile = new FileQueue<EmailJobs>({
11+
name: "email-queue",
1112
path: "./email-queue",
1213
});
1314

1415
// SQS-based queue (for production) - simple API
1516
export const emailQueueSqs = createSQSQueue<EmailJobs>(
17+
"email-queue",
1618
"https://sqs.us-east-1.amazonaws.com/428011609647/test-queue",
17-
{
18-
profile: "javier",
19-
}
19+
"delete"
2020
);
2121

2222
export const emailQueue = emailQueueSqs;

example-queue-project/src/general-queue.ts

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,7 @@ export interface GeneralJobs {
55
'generate-report': { type: string; period: string };
66
}
77

8-
export const generalQueue = createSQLiteQueue<GeneralJobs>('queue.db');
8+
export const generalQueue = createSQLiteQueue<GeneralJobs>('general-queue', 'queue.db');
99

1010
// Register job handlers for general queue
1111
generalQueue.setHandlers({

example-queue-project/src/redis-queue.ts

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@ interface EmailJobs {
66
}
77

88
// Create Redis queue with simple API
9-
export const emailQueue = createRedisQueue<EmailJobs>('redis://localhost:6379');
9+
export const emailQueue = createRedisQueue<EmailJobs>('redis-email-queue', 'redis://localhost:6379');
1010

1111
// Register job handlers
1212
emailQueue.setHandlers({

src/adapters/mongodb.ts

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -185,33 +185,35 @@ export class MongoDatabaseAdapter implements DatabaseAdapter {
185185
// Main export - constructor pattern
186186
export class MongoQueue<T = Record<string, any>> extends DbQueue<T> {
187187
mongoAdapter: MongoDatabaseAdapter;
188-
constructor(config: { collection: MongoCollection }) {
188+
constructor(config: { collection: MongoCollection; name: string }) {
189189
const adapter = new MongoDatabaseAdapter(config.collection);
190-
super(adapter);
190+
super(adapter, { name: config.name });
191191
this.mongoAdapter = adapter;
192192
}
193193
}
194194

195195
// Convenience factory for MongoDB driver
196196
export function createMongoQueue<T = Record<string, any>>(
197+
name: string,
197198
client: MongoClient,
198199
database: string,
199200
collection: string = 'jobs'
200201
): MongoQueue<T> {
201202
const db = client.db(database);
202203
const col = db.collection(collection);
203-
return new MongoQueue<T>({ collection: col });
204+
return new MongoQueue<T>({ collection: col, name });
204205
}
205206

206207
// Convenience factory with connection string
207208
export async function createMongoQueueFromUrl<T = Record<string, any>>(
209+
name: string,
208210
url: string,
209211
database: string,
210212
collection: string = 'jobs'
211213
): Promise<MongoQueue<T>> {
212214
const client = new MongoClient(url);
213215
await client.connect();
214-
return createMongoQueue<T>(client, database, collection);
216+
return createMongoQueue<T>(name, client, database, collection);
215217
}
216218

217219
// Re-export for convenience

src/adapters/redis.ts

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -309,16 +309,16 @@ export class RedisDatabaseAdapter implements DatabaseAdapter {
309309

310310
// Main export - constructor pattern for database-backed Redis queue
311311
export class RedisQueue<T = Record<string, any>> extends DbQueue<T> {
312-
constructor(config: { client: RedisClient; keyPrefix?: string }) {
312+
constructor(config: { client: RedisClient; keyPrefix?: string; name: string }) {
313313
const adapter = new RedisDatabaseAdapter(config.client, config.keyPrefix);
314-
super(adapter);
314+
super(adapter, { name: config.name });
315315
}
316316
}
317317

318318
// Convenience factory for node-redis
319-
export function createRedisQueue<T = Record<string, any>>(url?: string): RedisQueue<T> {
319+
export function createRedisQueue<T = Record<string, any>>(name: string, url?: string): RedisQueue<T> {
320320
const client = createClient(url ? { url } : {}) as any;
321-
return new RedisQueue<T>({ client });
321+
return new RedisQueue<T>({ client, name });
322322
}
323323

324324
// Re-export for convenience

src/adapters/sqlite.ts

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -227,16 +227,16 @@ export class SQLiteDatabaseAdapter implements DatabaseAdapter {
227227

228228
// Main export - constructor pattern
229229
export class SQLiteQueue<T = Record<string, any>> extends DbQueue<T> {
230-
constructor(config: { database: SQLiteDatabase }) {
230+
constructor(config: { database: SQLiteDatabase; name: string }) {
231231
const adapter = new SQLiteDatabaseAdapter(config.database);
232-
super(adapter);
232+
super(adapter, { name: config.name });
233233
}
234234
}
235235

236236
// Convenience factory for better-sqlite3
237-
export function createSQLiteQueue<T = Record<string, any>>(filename: string, options?: Database.Options): SQLiteQueue<T> {
237+
export function createSQLiteQueue<T = Record<string, any>>(name: string, filename: string, options?: Database.Options): SQLiteQueue<T> {
238238
const db = new Database(filename, options);
239-
return new SQLiteQueue<T>({ database: db });
239+
return new SQLiteQueue<T>({ database: db, name });
240240
}
241241

242242
// Re-export for convenience

src/adapters/sqs.ts

Lines changed: 9 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -7,19 +7,24 @@ export interface SQSConfig extends SQSClientConfig {
77

88
// Main export - constructor pattern
99
export class SQSQueue<T = Record<string, any>> extends SqsQueue<T> {
10-
constructor(config: { client: SQSClient; queueUrl: string }) {
11-
super(config.client, config.queueUrl);
10+
constructor(config: { client: SQSClient; queueUrl: string; name: string; onFailure: "delete" | "leaveInQueue" }) {
11+
super(config.client, config.queueUrl, { name: config.name, onFailure: config.onFailure });
1212
}
1313
}
1414

1515
// Convenience factory for AWS SDK v3
16-
export function createSQSQueue<T = Record<string, any>>(queueUrl: string, sqsConfig?: SQSConfig): SQSQueue<T> {
16+
export function createSQSQueue<T = Record<string, any>>(
17+
name: string,
18+
queueUrl: string,
19+
onFailure: "delete" | "leaveInQueue",
20+
sqsConfig?: SQSConfig
21+
): SQSQueue<T> {
1722
const client = new SQSClient({
1823
region: sqsConfig?.region || process.env.AWS_REGION || 'us-east-1',
1924
...sqsConfig
2025
});
2126

22-
return new SQSQueue<T>({ client, queueUrl });
27+
return new SQSQueue<T>({ client, queueUrl, name, onFailure });
2328
}
2429

2530
// Re-export for convenience

src/cli/worker.ts

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -84,15 +84,15 @@ async function main(): Promise<void> {
8484
process.exit(1);
8585
}
8686

87-
queue = new SqsQueue(config.sqsClient, config.queueUrl);
87+
queue = new SqsQueue(config.sqsClient, config.queueUrl, { name: 'cli-queue', onFailure: 'delete' });
8888
} else {
8989
if (!config.dbAdapter) {
9090
console.error('Error: Database adapter must be provided when using DB driver');
9191
console.error('This CLI is a template. You need to provide your own database adapter instance.');
9292
process.exit(1);
9393
}
9494

95-
queue = new DbQueue(config.dbAdapter);
95+
queue = new DbQueue(config.dbAdapter, { name: 'cli-queue' });
9696
}
9797

9898
const worker = new Worker(queue, {

src/core/queue.ts

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -39,7 +39,7 @@ export abstract class Queue<TJobMap = Record<string, any>, TJobRequest extends B
3939
protected ttrDefault = 300;
4040
protected plugins: QueuePlugin[];
4141
protected pluginDisposers: Array<() => Promise<void>> = [];
42-
public readonly name?: string;
42+
public readonly name: string;
4343

4444
/**
4545
* Registry of job handlers mapping job names to their handler functions.
@@ -62,14 +62,14 @@ export abstract class Queue<TJobMap = Record<string, any>, TJobRequest extends B
6262
* Creates a new Queue instance.
6363
*
6464
* @param options - Configuration options
65+
* @param options.name - Required name for the queue
6566
* @param options.ttrDefault - Default time-to-run for jobs in seconds (default: 300)
66-
* @param options.name - Optional name for the queue
6767
* @param options.plugins - Array of plugins to use with this queue
6868
*/
69-
constructor(options: QueueOptions = {}) {
69+
constructor(options: QueueOptions) {
7070
super();
71-
if (options.ttrDefault) this.ttrDefault = options.ttrDefault;
7271
this.name = options.name;
72+
if (options.ttrDefault) this.ttrDefault = options.ttrDefault;
7373
this.plugins = options.plugins || [];
7474
}
7575

@@ -235,7 +235,7 @@ export abstract class Queue<TJobMap = Record<string, any>, TJobRequest extends B
235235
if (this.pluginDisposers.length === 0) {
236236
for (const plugin of this.plugins) {
237237
if (plugin.init) {
238-
const dispose = await plugin.init({ queue: this as any, queueName: this.name });
238+
const dispose = await plugin.init({ queue: this as any });
239239
if (dispose) {
240240
this.pluginDisposers.push(dispose);
241241
disposers.push(dispose);

src/drivers/db.ts

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@ import type { QueueOptions } from '../interfaces/plugin.ts';
66
export class DbQueue<TJobMap = Record<string, any>> extends Queue<TJobMap, DbJobRequest<any>> {
77
constructor(
88
private db: DatabaseAdapter,
9-
options: QueueOptions = {}
9+
options: QueueOptions
1010
) {
1111
super(options);
1212
}

0 commit comments

Comments
 (0)