Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 6 additions & 4 deletions eventProcessors.js
Original file line number Diff line number Diff line change
Expand Up @@ -38,12 +38,12 @@ async function loadFromDatabase(pool) {
//eventProcessors = [];
// Register each processor from database
for (const row of result.rows) {
console.log(row);
// console.log(row);
registerProcessor(row.id, row.event_type, row.table_name, row.field_mappings, row.field_verification);
logger.debug('Event processor loaded', { id: row.id });
}

console.log(eventProcessors);

// console.log(eventProcessors);
logger.info(`Loaded ${result.rows.length} event processors from database`);
} catch (err) {
console.error(err);
Expand Down Expand Up @@ -234,7 +234,9 @@ function registerProcessor(id, eventType, tableName, fieldMappings, fieldVerific
const eventMid = event?.mid || 'unknown';
const eventId = event?.edata?.eks?.target?.id || event?.object?.id || 'unknown';
logger.error(`Error processing ${eventType} event (mid: ${eventMid}, id: ${eventId}): ${err.message}`);
logger.error(`Skipping bad record. Event data: ${JSON.stringify(event).substring(0, 500)}...`);
const safeEventMeta = { eid: event?.eid, ets: event?.ets };
logger.error(`Skipping bad record. Event meta: ${JSON.stringify(safeEventMeta)}`);
// logger.error(`Skipping bad record. Event data: ${JSON.stringify(event).substring(0, 500)}...`);
// Don't throw - allow batch to continue processing other events
return { success: false, skipped: true, error: err.message };
}
Expand Down
6 changes: 4 additions & 2 deletions index.js
Original file line number Diff line number Diff line change
Expand Up @@ -679,8 +679,9 @@ async function processTelemetryLogs(batchId = `batch_${Date.now()}`) {
const eventType = event.eid;
const eventUid = event.uid || "unknown";
const eventMid = event.mid || "unknown";
const maskedId = eventUid ? `${eventUid.substring(0, 4)}***` : 'unknown';
logger.debug(
`[${batchId}] [Log ${logIndex + 1}] [Event ${eventIndex + 1}/${events.length}] Processing event type: ${eventType}, uid: ${eventUid}, mid: ${eventMid}`,
`[${batchId}] [Log ${logIndex + 1}] [Event ${eventIndex + 1}/${events.length}] Processing event type: ${eventType}, uid: ${maskedId}, mid: ${eventMid}`,
);

let eventProcessed = false;
Expand Down Expand Up @@ -757,8 +758,9 @@ async function processTelemetryLogs(batchId = `batch_${Date.now()}`) {
logger.warn(
`[${batchId}] [Log ${logIndex + 1}] [Event ${eventIndex + 1}] mid: ${eventMid} - No processor matched for event type: ${eventType} - sending to dead letter queue`,
);
const maskedId = eventUid ? `${eventUid.substring(0, 4)}***` : 'unknown';
logger.debug(
`[${batchId}] [Log ${logIndex + 1}] [Event ${eventIndex + 1}] mid: ${eventMid} - Dead letter event uid: ${eventUid}, channel: ${event.channel || "unknown"}`,
`[${batchId}] [Log ${logIndex + 1}] [Event ${eventIndex + 1}] mid: ${eventMid} - Dead letter event uid: ${maskedId}, channel: ${event.channel || "unknown"}`,
);

await client.query(
Expand Down