Skip to content

Commit 0c57219

Browse files
committed
queue
1 parent a47ffb2 commit 0c57219

6 files changed

Lines changed: 34 additions & 39 deletions

File tree

example-queue-project/mongoose-example.ts

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
import mongoose from 'mongoose';
2-
import { createMongooseQueue, QueueJob } from '../src/adapters/mongoose.ts';
2+
import { MongooseQueue, QueueJob, createQueueModel } from '../src/drivers/mongoose.ts';
33

44
// Define your job types
55
interface MyJobs {
@@ -21,7 +21,8 @@ async function main() {
2121
console.log('Connected to MongoDB');
2222

2323
// Create a queue instance
24-
const queue = createMongooseQueue<MyJobs>('my-app');
24+
const model = createQueueModel();
25+
const queue = new MongooseQueue<MyJobs>({ model, name: 'my-app' });
2526

2627
// Set up job handlers
2728
queue.setHandlers({
@@ -78,7 +79,7 @@ async function customModelExample(): Promise<void> {
7879
const JobModel = mongoose.model('CustomJob', QueueJob.schema, 'custom_jobs');
7980

8081
// Create queue with custom model
81-
const queue = createMongooseQueue<MyJobs>('custom-app', JobModel);
82+
const queue = new MongooseQueue<MyJobs>({ model: JobModel, name: 'custom-app' });
8283

8384
// Set handlers and use as normal
8485
queue.setHandlers({
@@ -100,7 +101,8 @@ async function customModelExample(): Promise<void> {
100101
async function continuousProcessingExample(): Promise<void> {
101102
await mongoose.connect('mongodb://localhost:27017/queue-example');
102103

103-
const queue = createMongooseQueue<MyJobs>('continuous-queue');
104+
const model = createQueueModel();
105+
const queue = new MongooseQueue<MyJobs>({ model, name: 'continuous-queue' });
104106

105107
queue.setHandlers({
106108
'send-email': async (job) => {

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

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,11 +1,13 @@
1-
import { createSQLiteQueue } from "adapter-queue/sqlite";
1+
import { SQLiteQueue } from "adapter-queue/drivers/sqlite";
2+
import Database from "better-sqlite3";
23

34
export interface GeneralJobs {
45
'process-image': { url: string; width: number; height: number };
56
'generate-report': { type: string; period: string };
67
}
78

8-
export const generalQueue = createSQLiteQueue<GeneralJobs>('general-queue', 'queue.db');
9+
const db = new Database('queue.db');
10+
export const generalQueue = new SQLiteQueue<GeneralJobs>({ database: db, name: 'general-queue' });
911

1012
// Register job handlers for general queue
1113
generalQueue.setHandlers({

src/drivers/mongoose.ts

Lines changed: 0 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -220,14 +220,6 @@ export class MongooseQueue<TJobMap = Record<string, unknown>> extends DbQueue<TJ
220220
}
221221
}
222222

223-
// Factory function to create a Mongoose queue
224-
export function createMongooseQueue<TJobMap = Record<string, unknown>>(
225-
name: string,
226-
model?: Model<IQueueJob>
227-
): MongooseQueue<TJobMap> {
228-
const queueModel = model || createQueueModel();
229-
return new MongooseQueue<TJobMap>({ model: queueModel, name });
230-
}
231223

232224
// Create a default queue model
233225
export function createQueueModel(

src/drivers/sqlite.ts

Lines changed: 0 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -236,11 +236,6 @@ export class SQLiteQueue<T = Record<string, any>> extends DbQueue<T> {
236236
}
237237
}
238238

239-
// Convenience factory for better-sqlite3
240-
export function createSQLiteQueue<T = Record<string, any>>(name: string, filename: string, options?: Database.Options): SQLiteQueue<T> {
241-
const db = new Database(filename, options);
242-
return new SQLiteQueue<T>({ database: db, name });
243-
}
244239

245240
// Re-export for convenience
246241
export { DbQueue };

tests/all-queues.test.ts

Lines changed: 7 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -21,10 +21,11 @@ import {
2121
} from "@aws-sdk/client-sqs";
2222
import { InMemoryQueue } from "../src/drivers/memory.js";
2323
import { FileQueue } from "../src/drivers/file.js";
24-
import { createSQLiteQueue } from "../src/drivers/sqlite.js";
24+
import { SQLiteQueue } from "../src/drivers/sqlite.js";
25+
import Database from "better-sqlite3";
2526
import { RedisQueue } from "../src/drivers/redis.js";
2627
import { SqsQueue } from "../src/drivers/sqs.js";
27-
import { createMongooseQueue } from "../src/drivers/mongoose.js";
28+
import { MongooseQueue, createQueueModel } from "../src/drivers/mongoose.js";
2829
import mongoose from "mongoose";
2930
import type { Queue } from "../src/core/queue.js";
3031
import type { JobRequestFull } from "../src/interfaces/job.ts";
@@ -110,7 +111,8 @@ const drivers: Array<() => Promise<QueueDriverConfig> | QueueDriverConfig> = [
110111
},
111112
createQueue: async () => {
112113
// Use in-memory SQLite database for tests - much faster and no file cleanup needed
113-
return createSQLiteQueue<TestJobs>("test-queue", ":memory:");
114+
const db = new Database(":memory:");
115+
return new SQLiteQueue<TestJobs>({ database: db, name: "test-queue" });
114116
},
115117
cleanup: async (queue) => {
116118
// Clear all jobs from the in-memory database
@@ -267,7 +269,8 @@ const drivers: Array<() => Promise<QueueDriverConfig> | QueueDriverConfig> = [
267269
}
268270
},
269271
createQueue: async () => {
270-
return createMongooseQueue<TestJobs>("test-queue");
272+
const model = createQueueModel();
273+
return new MongooseQueue<TestJobs>({ model, name: "test-queue" });
271274
},
272275
cleanup: async () => {
273276
// Clean up MongoDB collections between tests

tests/drivers/mongoose.test.ts

Lines changed: 17 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,7 @@ import mongoose from 'mongoose';
33
import { GenericContainer } from 'testcontainers';
44
import type { StartedTestContainer } from 'testcontainers';
55
import {
6-
createMongooseQueue,
6+
MongooseQueue,
77
createQueueModel,
88
MongooseDatabaseAdapter,
99
QueueJobSchema
@@ -148,23 +148,24 @@ describe('Mongoose Adapter', () => {
148148
});
149149
});
150150

151-
describe('createMongooseQueue', () => {
151+
describe('MongooseQueue', () => {
152152
it('should create a queue with default model', () => {
153-
const queue = createMongooseQueue('test-queue');
153+
const model = createQueueModel();
154+
const queue = new MongooseQueue({ model, name: 'test-queue' });
154155
expect(queue).toBeDefined();
155156
expect(queue.name).toBe('test-queue');
156157
});
157158

158159
it('should create a queue with custom model', () => {
159-
const queue = createMongooseQueue('test-queue', testModel);
160+
const queue = new MongooseQueue({ model: testModel, name: 'test-queue' });
160161
expect(queue).toBeDefined();
161162
expect(queue.name).toBe('test-queue');
162163
});
163164

164165
it('should push and retrieve jobs', async () => {
165-
const queue = createMongooseQueue<{
166+
const queue = new MongooseQueue<{
166167
'test-job': { message: string };
167-
}>('test-queue', testModel);
168+
}>({ model: testModel, name: 'test-queue' });
168169

169170
let processedMessage = '';
170171
queue.setHandlers({
@@ -183,9 +184,9 @@ describe('Mongoose Adapter', () => {
183184
});
184185

185186
it('should handle job priorities', async () => {
186-
const queue = createMongooseQueue<{
187+
const queue = new MongooseQueue<{
187188
'priority-job': { priority: number };
188-
}>('test-queue', testModel);
189+
}>({ model: testModel, name: 'test-queue' });
189190

190191
// Push jobs with different priorities
191192
await queue.addJob('priority-job', { payload: { priority: 1 }, priority: 1 });
@@ -200,9 +201,9 @@ describe('Mongoose Adapter', () => {
200201
});
201202

202203
it('should handle delayed jobs', async () => {
203-
const queue = createMongooseQueue<{
204+
const queue = new MongooseQueue<{
204205
'delayed-job': { when: string };
205-
}>('test-queue', testModel);
206+
}>({ model: testModel, name: 'test-queue' });
206207

207208
// Push a delayed job
208209
await queue.addJob('delayed-job', { payload: { when: 'future' }, delaySeconds: 2 });
@@ -221,10 +222,10 @@ describe('Mongoose Adapter', () => {
221222

222223

223224
it('should handle job failures and retries through the queue system', async () => {
224-
const queue = createMongooseQueue<{
225+
const queue = new MongooseQueue<{
225226
'failing-job': { attemptNumber: number };
226227
'success-job': { data: string };
227-
}>('test-queue', testModel);
228+
}>({ model: testModel, name: 'test-queue' });
228229

229230
let attempts: number[] = [];
230231
let successfulJobs: string[] = [];
@@ -274,10 +275,10 @@ describe('Mongoose Adapter', () => {
274275
});
275276

276277
it('should handle TTR timeout and job recovery in real processing', async () => {
277-
const queue = createMongooseQueue<{
278+
const queue = new MongooseQueue<{
278279
'long-job': { duration: number };
279280
'quick-job': { data: string };
280-
}>('test-queue', testModel);
281+
}>({ model: testModel, name: 'test-queue' });
281282

282283
let jobExecutions: string[] = [];
283284

@@ -349,7 +350,7 @@ describe('Mongoose Adapter', () => {
349350

350351
describe('Integration with Mongoose features', () => {
351352
it('should work with Mongoose queries', async () => {
352-
const queue = createMongooseQueue('test-queue', testModel);
353+
const queue = new MongooseQueue({ model: testModel, name: 'test-queue' });
353354

354355
await queue.addJob('test-job', { payload: { data: 'test1' } });
355356
await queue.addJob('test-job', { payload: { data: 'test2' } });
@@ -364,7 +365,7 @@ describe('Mongoose Adapter', () => {
364365
});
365366

366367
it('should maintain Mongoose document structure', async () => {
367-
const queue = createMongooseQueue('test-queue', testModel);
368+
const queue = new MongooseQueue({ model: testModel, name: 'test-queue' });
368369

369370
const jobId = await queue.addJob('test-job', { payload: { test: true } });
370371

0 commit comments

Comments
 (0)