Skip to content

Commit fa42a9f

Browse files
committed
updates
1 parent 8bc6a9d commit fa42a9f

4 files changed

Lines changed: 108 additions & 79 deletions

File tree

package.json

Lines changed: 4 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -15,8 +15,6 @@
1515
"lint": "tsc --noEmit",
1616
"dev": "tsc --watch",
1717
"prepublishOnly": "npm run build && npm test",
18-
"queue:worker": "node dist/cli/worker.js",
19-
"queue:worker:isolate": "node dist/cli/worker.js --isolate",
2018
"release": "node scripts/release.ts"
2119
},
2220
"keywords": [
@@ -47,10 +45,8 @@
4745
"devDependencies": {
4846
"@aws-sdk/client-sqs": "^3.831.0",
4947
"@types/better-sqlite3": "7.6.13",
50-
"@types/mongodb": "^4.0.7",
5148
"@types/node": "^20.0.0",
5249
"better-sqlite3": "^12.0.0",
53-
"mongodb": "^6.17.0",
5450
"mongoose": "^8.0.0",
5551
"redis": "^5.5.6",
5652
"testcontainers": "^11.0.3",
@@ -74,10 +70,6 @@
7470
"import": "./dist/src/adapters/sqs.js",
7571
"types": "./dist/src/adapters/sqs.d.ts"
7672
},
77-
"./mongodb": {
78-
"import": "./dist/src/adapters/mongodb.js",
79-
"types": "./dist/src/adapters/mongodb.d.ts"
80-
},
8173
"./plugins/ecs-protection-manager": {
8274
"import": "./dist/src/plugins/ecs-protection-manager.js",
8375
"types": "./dist/src/plugins/ecs-protection-manager.d.ts"
@@ -89,12 +81,15 @@
8981
"./mongoose": {
9082
"import": "./dist/src/adapters/mongoose.js",
9183
"types": "./dist/src/adapters/mongoose.d.ts"
84+
},
85+
"./worker": {
86+
"import": "./dist/src/worker/worker.js",
87+
"types": "./dist/src/worker/worker.d.ts"
9288
}
9389
},
9490
"peerDependencies": {
9591
"@aws-sdk/client-sqs": "^3.0.0",
9692
"better-sqlite3": "^9.0.0",
97-
"mongodb": "^6.0.0",
9893
"mongoose": "^7.0.0 || ^8.0.0",
9994
"redis": "^4.0.0"
10095
},

pnpm-lock.yaml

Lines changed: 0 additions & 57 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

src/index.ts

Lines changed: 0 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -1,17 +1,4 @@
11
export { Queue } from './core/queue.ts';
2-
export { DbQueue } from './drivers/db.ts';
3-
export { FileQueue } from './drivers/file.ts';
4-
export { InMemoryQueue } from './drivers/memory.ts';
5-
export { Worker, runWorker } from './worker/worker.ts';
6-
export {
7-
MongooseQueue,
8-
createMongooseQueue,
9-
createQueueModel,
10-
QueueJob,
11-
MongooseDatabaseAdapter,
12-
QueueJobSchema
13-
} from './adapters/mongoose.ts';
14-
export type { IQueueJob } from './adapters/mongoose.ts';
152

163
export type {
174
JobStatus,

tests/adapters/mongoose.test.ts

Lines changed: 104 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -218,6 +218,110 @@ describe('Mongoose Adapter', () => {
218218
const job2 = await queue.mongooseAdapter.reserveJob(60);
219219
expect(job2).toBeDefined();
220220
});
221+
222+
223+
it('should handle job failures and retries through the queue system', async () => {
224+
const queue = createMongooseQueue<{
225+
'failing-job': { attemptNumber: number };
226+
'success-job': { data: string };
227+
}>('test-queue', testModel);
228+
229+
let attempts: number[] = [];
230+
let successfulJobs: string[] = [];
231+
232+
queue.setHandlers({
233+
'failing-job': async (job) => {
234+
attempts.push(job.payload.attemptNumber);
235+
236+
if (job.payload.attemptNumber <= 2) {
237+
throw new Error(`Intentional failure on attempt ${job.payload.attemptNumber}`);
238+
}
239+
240+
// Success on attempt 3
241+
},
242+
'success-job': async (job) => {
243+
successfulJobs.push(job.payload.data);
244+
}
245+
});
246+
247+
// Add jobs
248+
const failingJob1 = await queue.addJob('failing-job', { payload: { attemptNumber: 1 } });
249+
const failingJob2 = await queue.addJob('failing-job', { payload: { attemptNumber: 2 } });
250+
const failingJob3 = await queue.addJob('failing-job', { payload: { attemptNumber: 3 } });
251+
const successJob = await queue.addJob('success-job', { payload: { data: 'test-data' } });
252+
253+
// Process all jobs
254+
await queue.run();
255+
256+
// Check results
257+
expect(attempts).toEqual([1, 2, 3]);
258+
expect(successfulJobs).toEqual(['test-data']);
259+
260+
// Check final job statuses
261+
const job1Doc = await testModel.findById(failingJob1);
262+
const job2Doc = await testModel.findById(failingJob2);
263+
const job3Doc = await testModel.findById(failingJob3);
264+
const successDoc = await testModel.findById(successJob);
265+
266+
expect(job1Doc?.status).toBe('failed');
267+
expect(job2Doc?.status).toBe('failed');
268+
expect(job3Doc?.status).toBe('done');
269+
expect(successDoc?.status).toBe('done');
270+
271+
// Check error messages
272+
expect(job1Doc?.errorMessage).toContain('Intentional failure on attempt 1');
273+
expect(job2Doc?.errorMessage).toContain('Intentional failure on attempt 2');
274+
});
275+
276+
it('should handle TTR timeout and job recovery in real processing', async () => {
277+
const queue = createMongooseQueue<{
278+
'long-job': { duration: number };
279+
'quick-job': { data: string };
280+
}>('test-queue', testModel);
281+
282+
let jobExecutions: string[] = [];
283+
284+
queue.setHandlers({
285+
'long-job': async (job) => {
286+
jobExecutions.push(`long-job-start-${job.payload.duration}`);
287+
// Simulate a job that takes longer than its TTR
288+
await new Promise(resolve => setTimeout(resolve, job.payload.duration));
289+
jobExecutions.push(`long-job-end-${job.payload.duration}`);
290+
},
291+
'quick-job': async (job) => {
292+
jobExecutions.push(`quick-job-${job.payload.data}`);
293+
}
294+
});
295+
296+
// Add a job with very short TTR that will timeout
297+
const longJobId = await queue.addJob('long-job', {
298+
payload: { duration: 2000 }, // 2 seconds
299+
ttr: 1 // 1 second TTR - will timeout
300+
});
301+
302+
// Add a quick job
303+
const quickJobId = await queue.addJob('quick-job', {
304+
payload: { data: 'test' }
305+
});
306+
307+
// Process queue - long job will timeout, quick job should complete
308+
await queue.run();
309+
310+
// Verify states after first run
311+
let longJobDoc = await testModel.findById(longJobId);
312+
let quickJobDoc = await testModel.findById(quickJobId);
313+
314+
// Quick job should be done, long job should be in some intermediate state
315+
expect(quickJobDoc?.status).toBe('done');
316+
317+
// Check what actually got executed
318+
expect(jobExecutions).toContain('quick-job-test');
319+
expect(jobExecutions).toContain('long-job-start-2000');
320+
321+
// The long job may or may not have finished depending on timing
322+
// but it should have been processed at least once
323+
expect(longJobDoc).toBeDefined();
324+
}, 10000); // Longer timeout for this test
221325
});
222326

223327
describe('createQueueModel', () => {

0 commit comments

Comments
 (0)