Repository navigation
Expand file tree
/
Copy pathprocessor.js
More file actions
160 lines (140 loc) · 4.79 KB
/
Copy pathprocessor.js
File metadata and controls
160 lines (140 loc) · 4.79 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
/**
* pagingWrick Queue Processor
*
* This script processes messages from the Service Bus queue and sends Discord notifications.
* Can run locally, as a scheduled task, or in any Node.js environment.
*
* Authentication modes:
* - Managed Identity (Azure): Set SERVICEBUS_NAMESPACE (e.g., "sb-pagingwrick-dev-xxx.servicebus.windows.net")
* - Connection String (local): Set SERVICE_BUS_CONNECTION_STRING
*/
const { ServiceBusClient } = require('@azure/service-bus');
const { DefaultAzureCredential } = require('@azure/identity');
// Configuration from environment variables
const SERVICEBUS_NAMESPACE = process.env.SERVICEBUS_NAMESPACE;
const SERVICE_BUS_CONNECTION_STRING = process.env.SERVICE_BUS_CONNECTION_STRING;
const QUEUE_NAME = process.env.QUEUE_NAME || 'events';
const DISCORD_WEBHOOK_URL = process.env.DISCORD_WEBHOOK_URL;
// Determine authentication mode
const useMangedIdentity = !!SERVICEBUS_NAMESPACE;
if (!SERVICEBUS_NAMESPACE && !SERVICE_BUS_CONNECTION_STRING) {
console.error('❌ Either SERVICEBUS_NAMESPACE or SERVICE_BUS_CONNECTION_STRING is required');
console.error(' For Azure (managed identity): Set SERVICEBUS_NAMESPACE=sb-xxx.servicebus.windows.net');
console.error(' For local dev: Set SERVICE_BUS_CONNECTION_STRING=Endpoint=sb://...');
process.exit(1);
}
if (!DISCORD_WEBHOOK_URL) {
console.error('❌ DISCORD_WEBHOOK_URL environment variable is required');
process.exit(1);
}
/**
* Send Discord notification
*/
async function sendDiscordNotification(event) {
const severityEmoji = event.severity === 'high' ? '🔴' : '🟠';
const color = event.severity === 'high' ? 0xff0000 : 0xffa500;
const payload = {
embeds: [{
title: `${severityEmoji} ${event.severity.toUpperCase()} SEVERITY ALERT`,
description: event.message,
color: color,
fields: [
{
name: 'Source',
value: event.source,
inline: true
},
{
name: 'Severity',
value: event.severity.toUpperCase(),
inline: true
},
{
name: 'Timestamp',
value: new Date(event.timestamp).toLocaleString(),
inline: false
},
...(event.metadata ? [{
name: 'Additional Info',
value: '```json\n' + JSON.stringify(event.metadata, null, 2) + '\n```',
inline: false
}] : [])
],
timestamp: event.timestamp,
footer: {
text: 'pagingWrick notification service'
}
}]
};
const response = await fetch(DISCORD_WEBHOOK_URL, {
method: 'POST',
headers: {
'Content-Type': 'application/json'
},
body: JSON.stringify(payload)
});
if (!response.ok) {
throw new Error(`Discord webhook failed: ${response.status} ${response.statusText}`);
}
console.log('✅ Discord notification sent');
}
/**
* Process a single message
*/
async function processMessage(message) {
console.log(`\n📨 Processing message:`, message.body);
try {
const event = message.body;
// Validate event structure
if (!event.source || !event.severity || !event.message) {
console.error('❌ Invalid event schema:', event);
return;
}
console.log(` Source: ${event.source}`);
console.log(` Severity: ${event.severity}`);
console.log(` Message: ${event.message}`);
// Send Discord notification for all severities
console.log(` 🔔 ${event.severity.toUpperCase()} severity - sending Discord notification...`);
await sendDiscordNotification(event);
} catch (error) {
console.error('❌ Error processing message:', error);
throw error; // Re-throw to move to dead-letter queue
}
}
/**
* Main processor loop
*/
async function main() {
console.log('🚀 pagingWrick Queue Processor starting...');
console.log(` Queue: ${QUEUE_NAME}`);
console.log(` Auth mode: ${useMangedIdentity ? 'Managed Identity' : 'Connection String'}`);
if (useMangedIdentity) {
console.log(` Namespace: ${SERVICEBUS_NAMESPACE}`);
}
console.log(` Waiting for messages...\n`);
// Create ServiceBusClient with appropriate authentication
let client;
if (useMangedIdentity) {
const credential = new DefaultAzureCredential();
client = new ServiceBusClient(SERVICEBUS_NAMESPACE, credential);
} else {
client = new ServiceBusClient(SERVICE_BUS_CONNECTION_STRING);
}
const receiver = client.createReceiver(QUEUE_NAME);
// Subscribe to messages
receiver.subscribe({
processMessage: async (message) => {
await processMessage(message);
},
processError: async (error) => {
console.error('❌ Error from message receiver:', error);
}
});
// Keep running
console.log('✅ Processor is running. Press Ctrl+C to stop.\n');
}
// Run the processor
main().catch((error) => {
console.error('❌ Fatal error:', error);
process.exit(1);
});