Skip to content

Commit 67d8e6f

Browse files
committed
📦 new (webhook): add composite fingerprint deduplication
1 parent b176182 commit 67d8e6f

2 files changed

Lines changed: 271 additions & 3 deletions

File tree

‎src/services/webhookService.test.ts‎

Lines changed: 233 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -28,11 +28,13 @@ vi.mock('../config/redis', () => ({
2828
connect: vi.fn(),
2929
lPush: vi.fn().mockResolvedValue(1),
3030
exists: vi.fn().mockResolvedValue(0),
31+
set: vi.fn().mockResolvedValue('OK'),
3132
setEx: vi.fn().mockResolvedValue('OK'),
3233
quit: vi.fn(),
3334
},
3435
redisEventConfig: {
3536
keyPrefix: 'webhook:event:',
37+
fingerprintPrefix: 'unthread:fp:',
3638
eventTtl: 259200,
3739
}
3840
}));
@@ -316,4 +318,235 @@ describe('WebhookService', () => {
316318
expect(detected).toBe('dashboard');
317319
});
318320
});
321+
322+
describe('Composite Fingerprint Deduplication', () => {
323+
it('should generate fingerprint from eventTimestamp, event type, and data.id', () => {
324+
const event = buildMessageEvent({
325+
eventTimestamp: 1772463244428,
326+
data: {
327+
id: 'T08DF0UA02H-C08DWG00P25-1772463242.918629',
328+
conversationId: 'conv-1',
329+
content: 'Hello there!',
330+
},
331+
});
332+
333+
const fingerprint = (service as any).generateFingerprint(event);
334+
expect(fingerprint).toBe('1772463244428:message_created:T08DF0UA02H-C08DWG00P25-1772463242.918629');
335+
});
336+
337+
it('should generate same fingerprint for retries with different eventIds', () => {
338+
const event1 = buildMessageEvent({
339+
eventId: '5fb567d5-7de5-4aab-a89b-37123a315df3',
340+
eventTimestamp: 1772463244428,
341+
data: {
342+
id: 'T08DF0UA02H-C08DWG00P25-1772463242.918629',
343+
conversationId: 'conv-1',
344+
content: 'Hello there!',
345+
},
346+
});
347+
348+
const event2 = buildMessageEvent({
349+
eventId: '94a4cbce-0772-4ad8-bd06-d20cc36c1818',
350+
eventTimestamp: 1772463244428,
351+
data: {
352+
id: 'T08DF0UA02H-C08DWG00P25-1772463242.918629',
353+
conversationId: 'conv-1',
354+
content: 'Hello there!',
355+
},
356+
});
357+
358+
const fp1 = (service as any).generateFingerprint(event1);
359+
const fp2 = (service as any).generateFingerprint(event2);
360+
expect(fp1).toBe(fp2);
361+
});
362+
363+
it('should generate different fingerprints for genuinely different messages', () => {
364+
const event1 = buildMessageEvent({
365+
eventTimestamp: 1772463244428,
366+
data: {
367+
id: 'msg-1',
368+
conversationId: 'conv-1',
369+
content: 'Hello there!',
370+
},
371+
});
372+
373+
const event2 = buildMessageEvent({
374+
eventTimestamp: 1772463693587,
375+
data: {
376+
id: 'msg-2',
377+
conversationId: 'conv-1',
378+
content: 'Hello there!',
379+
},
380+
});
381+
382+
const fp1 = (service as any).generateFingerprint(event1);
383+
const fp2 = (service as any).generateFingerprint(event2);
384+
expect(fp1).not.toBe(fp2);
385+
});
386+
387+
it('should return null when data.id is missing', () => {
388+
const event: UnthreadWebhookEvent = {
389+
event: 'message_created',
390+
eventId: 'test-id',
391+
eventTimestamp: Date.now(),
392+
webhookTimestamp: Date.now(),
393+
data: { conversationId: 'conv-1' },
394+
};
395+
396+
const fingerprint = (service as any).generateFingerprint(event);
397+
expect(fingerprint).toBeNull();
398+
});
399+
400+
it('should return null when data is missing', () => {
401+
const event: UnthreadWebhookEvent = {
402+
event: 'message_created',
403+
eventId: 'test-id',
404+
eventTimestamp: Date.now(),
405+
webhookTimestamp: Date.now(),
406+
};
407+
408+
const fingerprint = (service as any).generateFingerprint(event);
409+
expect(fingerprint).toBeNull();
410+
});
411+
412+
it('should generate fingerprint for conversation events', () => {
413+
const event: UnthreadWebhookEvent = {
414+
event: 'conversation_updated',
415+
eventId: 'test-conv-update',
416+
eventTimestamp: 1772463244428,
417+
webhookTimestamp: Date.now(),
418+
data: {
419+
id: 'conv-123',
420+
title: 'Updated Conversation',
421+
},
422+
};
423+
424+
const fingerprint = (service as any).generateFingerprint(event);
425+
expect(fingerprint).toBe('1772463244428:conversation_updated:conv-123');
426+
});
427+
});
428+
429+
describe('processEvent Integration - Dedup Flow', () => {
430+
// eslint-disable-next-line @typescript-eslint/no-explicit-any
431+
let redisClient: any;
432+
433+
beforeEach(async () => {
434+
const redis = await import('../config/redis');
435+
redisClient = redis.client;
436+
vi.mocked(redisClient.exists).mockReset().mockResolvedValue(0);
437+
vi.mocked(redisClient.set).mockReset().mockResolvedValue('OK');
438+
vi.mocked(redisClient.setEx).mockReset().mockResolvedValue('OK');
439+
vi.mocked(redisClient.lPush).mockReset().mockResolvedValue(1);
440+
vi.mocked(redisClient.connect).mockReset();
441+
});
442+
443+
it('should reject event when eventId already exists in Redis', async () => {
444+
// eventId check: already processed
445+
vi.mocked(redisClient.exists).mockResolvedValueOnce(1);
446+
447+
const event = buildMessageEvent({
448+
eventTimestamp: 1772463244428,
449+
data: {
450+
id: 'msg-dup-eventid',
451+
conversationId: 'conv-1',
452+
content: 'Hello',
453+
metadata: { event_payload: { conversationUpdates: {} } },
454+
},
455+
});
456+
457+
await service.processEvent(event);
458+
459+
// Should not attempt fingerprint claim or queue publish
460+
expect(redisClient.set).not.toHaveBeenCalled();
461+
expect(redisClient.lPush).not.toHaveBeenCalled();
462+
});
463+
464+
it('should reject event when fingerprint already exists (SET NX returns null)', async () => {
465+
// eventId check: not found
466+
vi.mocked(redisClient.exists).mockResolvedValueOnce(0);
467+
// fingerprint claim: already exists (SET NX fails)
468+
vi.mocked(redisClient.set).mockResolvedValueOnce(null);
469+
470+
const event = buildMessageEvent({
471+
eventTimestamp: 1772463244428,
472+
data: {
473+
id: 'msg-dup-fp',
474+
conversationId: 'conv-1',
475+
content: 'Hello',
476+
metadata: { event_payload: { conversationUpdates: {} } },
477+
},
478+
});
479+
480+
await service.processEvent(event);
481+
482+
// Should have attempted atomic fingerprint claim
483+
expect(redisClient.set).toHaveBeenCalledWith(
484+
expect.stringContaining('unthread:fp:'),
485+
'processed',
486+
expect.objectContaining({ EX: expect.any(Number), NX: true })
487+
);
488+
// Should not have published to queue
489+
expect(redisClient.lPush).not.toHaveBeenCalled();
490+
// Should have marked eventId for faster future lookups
491+
expect(redisClient.setEx).toHaveBeenCalledWith(
492+
expect.stringContaining('webhook:event:'),
493+
expect.any(Number),
494+
'processed'
495+
);
496+
});
497+
498+
it('should claim fingerprint atomically and publish event when new', async () => {
499+
// eventId check: not found
500+
vi.mocked(redisClient.exists).mockResolvedValueOnce(0);
501+
// fingerprint claim: success (SET NX returns OK)
502+
vi.mocked(redisClient.set).mockResolvedValueOnce('OK');
503+
504+
const event = buildMessageEvent({
505+
eventTimestamp: 1772463244428,
506+
data: {
507+
id: 'msg-new',
508+
conversationId: 'conv-1',
509+
content: 'Hello',
510+
metadata: { event_payload: { conversationUpdates: {} } },
511+
},
512+
});
513+
514+
await service.processEvent(event);
515+
516+
// Should have claimed fingerprint via SET NX
517+
expect(redisClient.set).toHaveBeenCalledWith(
518+
expect.stringContaining('unthread:fp:1772463244428:message_created:msg-new'),
519+
'processed',
520+
expect.objectContaining({ EX: expect.any(Number), NX: true })
521+
);
522+
// Should have published to queue
523+
expect(redisClient.lPush).toHaveBeenCalled();
524+
// Should have marked eventId as processed
525+
expect(redisClient.setEx).toHaveBeenCalled();
526+
});
527+
528+
it('should claim fingerprint before processing buffered or non-buffered events', async () => {
529+
// eventId check: not found
530+
vi.mocked(redisClient.exists).mockResolvedValueOnce(0);
531+
// fingerprint claim: success
532+
vi.mocked(redisClient.set).mockResolvedValueOnce('OK');
533+
534+
const event = buildMessageEvent({
535+
eventTimestamp: 1772463244428,
536+
data: {
537+
id: 'msg-order-test',
538+
conversationId: 'conv-1',
539+
content: 'Hello',
540+
metadata: { event_payload: { conversationUpdates: {} } },
541+
},
542+
});
543+
544+
await service.processEvent(event);
545+
546+
// Verify fingerprint SET NX was called before lPush
547+
const setCallOrder = vi.mocked(redisClient.set).mock.invocationCallOrder[0];
548+
const lPushCallOrder = vi.mocked(redisClient.lPush).mock.invocationCallOrder[0];
549+
expect(setCallOrder).toBeLessThan(lPushCallOrder);
550+
});
551+
});
319552
});

‎src/services/webhookService.ts‎

Lines changed: 38 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -48,24 +48,39 @@ export class WebhookService {
4848

4949
await this.initializeServices();
5050

51-
// Check for duplicate events
51+
// Check for duplicate events by eventId (exact retry)
5252
const eventExists = await this.redisService.eventExists(event.eventId);
5353
if (eventExists) {
5454
LogEngine.info(`Event already processed - duplicate detected: ${event.eventId}`);
5555
return;
5656
}
5757

58+
// Atomically claim fingerprint slot (retry with new eventId detection)
59+
// Uses SET NX to combine check+mark into a single atomic operation
60+
const fingerprint = this.generateFingerprint(event);
61+
if (fingerprint) {
62+
const claimed = await this.redisService.claimFingerprint(fingerprint);
63+
if (!claimed) {
64+
LogEngine.info(`Event already processed - fingerprint duplicate detected: ${event.eventId} (fp: ${fingerprint})`);
65+
// Mark eventId too so subsequent retries with same eventId are caught faster
66+
await this.redisService.markEventProcessed(event.eventId);
67+
return;
68+
}
69+
}
70+
5871
// Detect platform source (enhanced with file attachment correlation)
5972
const sourcePlatform = this.detectPlatformSource(event);
6073

6174
// Handle buffered events - they will be processed later via callback
6275
if (sourcePlatform === 'buffered') {
6376
// Mark buffered events as processed to prevent duplicate buffering on retries
77+
// Note: fingerprint already claimed atomically above via SET NX
6478
await this.redisService.markEventProcessed(event.eventId);
6579

6680
const processingTime = Date.now() - startTime;
6781
LogEngine.info('File attachment event buffered for correlation', {
6882
eventId: event.eventId,
83+
fingerprint,
6984
hasFiles: this.fileAttachmentCorrelation.hasFileAttachments(event),
7085
processingTime: `${processingTime}ms`
7186
});
@@ -78,6 +93,7 @@ export class WebhookService {
7893
const totalProcessingTime = Date.now() - startTime;
7994
LogEngine.debug(`Event processing completed`, {
8095
eventId: event.eventId,
96+
fingerprint,
8197
sourcePlatform,
8298
totalProcessingTime: `${totalProcessingTime}ms`
8399
});
@@ -282,7 +298,7 @@ export class WebhookService {
282298
const transformedEvent = this.transformEvent(event, sourcePlatform);
283299
await this.redisService.publishEvent(transformedEvent);
284300

285-
// Mark as processed
301+
// Mark eventId as processed
286302
await this.redisService.markEventProcessed(event.eventId);
287303

288304
} catch (error) {
@@ -301,4 +317,23 @@ export class WebhookService {
301317
destroy(): void {
302318
this.fileAttachmentCorrelation.destroy();
303319
}
304-
}
320+
321+
/**
322+
* Generate a composite fingerprint for retry deduplication.
323+
* Format: {eventTimestamp}:{event}:{data.id}
324+
*
325+
* This catches Unthread retries that assign a new eventId to the same logical event.
326+
* Returns null if insufficient data to generate a meaningful fingerprint.
327+
*/
328+
private generateFingerprint(event: UnthreadWebhookEvent): string | null {
329+
const eventTimestamp = event.eventTimestamp;
330+
const eventType = event.event;
331+
const dataId = event.data?.id;
332+
333+
if (!eventTimestamp || !eventType || !dataId) {
334+
return null;
335+
}
336+
337+
return `${eventTimestamp}:${eventType}:${dataId}`;
338+
}
339+
}

0 commit comments

Comments
 (0)