Skip to content

Commit 292fe7c

Browse files
committed
refactor: streamline job handling and improve logging in queues
1 parent 07eeec2 commit 292fe7c

4 files changed

Lines changed: 36 additions & 44 deletions

File tree

example-queue-project/package.json

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -9,9 +9,9 @@
99
"lint": "tsc --noEmit",
1010
"dev": "tsx src/index.ts",
1111
"start": "node dist/index.js",
12-
"add-job": "tsx src/add-job.ts",
13-
"process-jobs": "tsx src/process-jobs.ts",
14-
"demo": "echo 'Run \"npm run process-jobs\" in one terminal, then \"npm run add-job\" in another!'"
12+
"push": "tsx src/index.ts --push",
13+
"run": "tsx src/index.ts --run",
14+
"demo": "echo 'Run \"npm run run\" in one terminal, then \"npm run push\" in another!'"
1515
},
1616
"dependencies": {
1717
"@muniter/queue": "file:..",

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

Lines changed: 15 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,8 @@
11
import { FileQueue } from "@muniter/queue";
22

33
interface EmailJobs {
4-
'welcome-email': { to: string; name: string };
5-
'notification': { to: string; subject: string; body: string };
4+
'welcome-email': { to: string; name: string };
5+
'notification': { to: string; subject: string; body: string };
66
}
77

88

@@ -28,4 +28,16 @@ emailQueue.onJob('notification', async (payload) => {
2828
console.log(`Sending notification email to ${to}: ${subject}`);
2929
await new Promise(resolve => setTimeout(resolve, 500));
3030
console.log(`Notification sent successfully`);
31-
});
31+
});
32+
33+
emailQueue.on('beforeExec', (event) => {
34+
console.log(`\n[emailQueue][${new Date().toISOString()}] Starting ${event.name} job ${event.id}...`);
35+
});
36+
37+
emailQueue.on('afterExec', (event) => {
38+
console.log(`[emailQueue][${new Date().toISOString()}] Email job ${event.id} (${event.name}) completed successfully`);
39+
});
40+
41+
emailQueue.on('afterError', (event) => {
42+
console.error(`[emailQueue][${new Date().toISOString()}] Email job ${event.id} (${event.name}) failed:`, event.error);
43+
});

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

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -34,4 +34,17 @@ generalQueue.onJob('generate-report', async (payload) => {
3434
}
3535

3636
console.log(`Report generated: ${type} for ${period}`);
37+
});
38+
39+
// Add event listeners for both queues
40+
generalQueue.on('beforeExec', (event) => {
41+
console.log(`\n[generalQueue][${new Date().toISOString()}] Starting ${event.name} job ${event.id}...`);
42+
});
43+
44+
generalQueue.on('afterExec', (event) => {
45+
console.log(`[generalQueue][${new Date().toISOString()}] Job ${event.id} (${event.name}) completed successfully`);
46+
});
47+
48+
generalQueue.on('afterError', (event) => {
49+
console.error(`[generalQueue][${new Date().toISOString()}] Job ${event.id} (${event.name}) failed:`, event.error);
3750
});

example-queue-project/src/index.ts

Lines changed: 5 additions & 38 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
1-
import { initializeDatabase, run } from './database.js';
1+
import { initializeDatabase } from './database.js';
22
import { emailQueue } from './email-queue.js';
33
import { generalQueue } from './general-queue.js';
44
import { parseArgs } from 'util';
@@ -37,53 +37,20 @@ async function push() {
3737
});
3838
console.log(`Image job added with ID: ${imageJobId}`);
3939

40-
const reportJobId = await generalQueue.addJob('generate-report', {
40+
await generalQueue.addJob('generate-report', {
4141
payload: {
4242
type: 'sales',
4343
period: 'Q4-2023'
4444
},
4545
delay: 2
4646
});
47-
console.log(`Delayed report job added with ID: ${reportJobId}`);
48-
49-
// Add event listeners for both queues
50-
generalQueue.on('beforeExec', (event) => {
51-
console.log(`\n[generalQueue][${new Date().toISOString()}] Starting ${event.name} job ${event.id}...`);
52-
});
53-
54-
generalQueue.on('afterExec', (event) => {
55-
console.log(`[generalQueue][${new Date().toISOString()}] Job ${event.id} (${event.name}) completed successfully`);
56-
});
57-
58-
generalQueue.on('afterError', (event) => {
59-
console.error(`[generalQueue][${new Date().toISOString()}] Job ${event.id} (${event.name}) failed:`, event.error);
60-
});
61-
62-
emailQueue.on('beforeExec', (event) => {
63-
console.log(`\n[emailQueue][${new Date().toISOString()}] Starting ${event.name} job ${event.id}...`);
64-
});
65-
66-
emailQueue.on('afterExec', (event) => {
67-
console.log(`[emailQueue][${new Date().toISOString()}] Email job ${event.id} (${event.name}) completed successfully`);
68-
});
69-
70-
emailQueue.on('afterError', (event) => {
71-
console.error(`[emailQueue][${new Date().toISOString()}] Email job ${event.id} (${event.name}) failed:`, event.error);
72-
});
73-
74-
console.log('\nStarting workers...');
75-
76-
process.on('SIGINT', () => {
77-
console.log('\nShutting down gracefully...');
78-
process.exit(0);
79-
});
8047
}
8148

8249
async function run() {
8350
console.log('Starting queue workers...');
8451
await Promise.allSettled([
85-
emailQueue.run(true),
86-
generalQueue.run(true)
52+
emailQueue.run(true, 1),
53+
generalQueue.run(true, 1),
8754
])
8855
console.log('Queue workers are running. Press Ctrl+C to exit.');
8956
}
@@ -105,7 +72,7 @@ if (import.meta.url === `file://${process.argv[1]}`) {
10572
console.log(' --run Run the queue workers (default: false)');
10673
process.exit(0);
10774
}
108-
75+
10976
if (args.values.push) {
11077
await push();
11178
} else if (args.values.run) {

0 commit comments

Comments
 (0)