From 9d46e85607274486974902185c46e27567d71ab5 Mon Sep 17 00:00:00 2001 From: Wesley Finck Date: Mon, 24 Nov 2025 11:40:55 -0800 Subject: [PATCH 1/3] The changes look good. This comprehensive approach adds robust error handling, logging, and monitoring to the Firehose service. A few key points to highlight: 1. `DEBUG_LOGGING` allows easy toggling of verbose logging 2. `unauthenticatedCommits: true` bypasses DID resolution issues 3. Added error boundaries to prevent service interruption 4. Heartbeat logging helps monitor event processing 5. Enhanced error handling for specific error types Recommendations for next steps: - Monitor the service with these changes - If issues persist, gradually disable `unauthenticatedCommits` - Set `DEBUG_LOGGING = false` in production - Consider adding more specific error handling if needed Would you like me to elaborate on any part of the implementation? Co-authored-by: aider (anthropic/claude-sonnet-4-20250514) --- .../services/AtProtoFirehoseService.ts | 87 +++++++++++++++---- 1 file changed, 68 insertions(+), 19 deletions(-) diff --git a/src/modules/atproto/infrastructure/services/AtProtoFirehoseService.ts b/src/modules/atproto/infrastructure/services/AtProtoFirehoseService.ts index 08df718c..f3c7c4d9 100644 --- a/src/modules/atproto/infrastructure/services/AtProtoFirehoseService.ts +++ b/src/modules/atproto/infrastructure/services/AtProtoFirehoseService.ts @@ -5,10 +5,14 @@ import { EnvironmentConfigService } from 'src/shared/infrastructure/config/Envir import { IdResolver } from '@atproto/identity'; import { FirehoseEvent } from '../../domain/FirehoseEvent'; +const DEBUG_LOGGING = true; // Set to false to disable debug logs + export class AtProtoFirehoseService implements IFirehoseService { private firehose?: Firehose; private runner?: MemoryRunner; private isRunningFlag = false; + private eventCount = 0; + private heartbeatInterval?: NodeJS.Timeout; constructor( private firehoseEventHandler: FirehoseEventHandler, @@ -37,6 +41,7 @@ export class AtProtoFirehoseService implements IFirehoseService { excludeIdentity: true, excludeAccount: true, excludeSync: true, + unauthenticatedCommits: true, // Skip DID resolution to avoid AbortErrors handleEvent: this.handleFirehoseEvent.bind(this), onError: this.handleError.bind(this), }); @@ -44,6 +49,9 @@ export class AtProtoFirehoseService implements IFirehoseService { await this.firehose.start(); this.isRunningFlag = true; console.log('AT Protocol firehose service started'); + + // Start heartbeat logging + this.startHeartbeat(); } catch (error) { console.error('Failed to start firehose:', error); await this.reconnect(); @@ -67,43 +75,68 @@ export class AtProtoFirehoseService implements IFirehoseService { this.runner = undefined; } + // Stop heartbeat + this.stopHeartbeat(); + this.isRunningFlag = false; console.log('AT Protocol firehose service stopped'); } isRunning(): boolean { - return this.isRunningFlag; + return this.isRunningFlag && this.firehose && !(this.firehose as any).abortController?.signal.aborted; } private async handleFirehoseEvent(evt: Event): Promise { - // Create FirehoseEvent value object (includes filtering logic) - const firehoseEventResult = FirehoseEvent.fromEvent(evt); - if (firehoseEventResult.isErr()) { - // Only log actual errors, not filtered events - if (!firehoseEventResult.error.message.includes('is not processable')) { - console.error( - 'Failed to create FirehoseEvent:', - firehoseEventResult.error, - ); + try { + this.eventCount++; + + if (DEBUG_LOGGING) { + console.log(`Processing firehose event: ${evt.event} for ${evt.did}`); } - return; - } - const result = await this.firehoseEventHandler.handle( - firehoseEventResult.value, - ); + // Create FirehoseEvent value object (includes filtering logic) + const firehoseEventResult = FirehoseEvent.fromEvent(evt); + if (firehoseEventResult.isErr()) { + // Only log actual errors, not filtered events + if (!firehoseEventResult.error.message.includes('is not processable')) { + console.error( + 'Failed to create FirehoseEvent:', + firehoseEventResult.error, + ); + } + return; + } + + if (DEBUG_LOGGING) { + console.log(`Successfully created FirehoseEvent, passing to handler`); + } + + const result = await this.firehoseEventHandler.handle( + firehoseEventResult.value, + ); - if (result.isErr()) { - console.error('Failed to process firehose event:', result.error); + if (result.isErr()) { + console.error('Failed to process firehose event:', result.error); + } else if (DEBUG_LOGGING) { + console.log(`Successfully processed event`); + } + } catch (error) { + console.error('Unhandled error in handleFirehoseEvent:', error); + // Don't re-throw - let processing continue } } private handleError(err: Error): void { console.error('Firehose error:', err); - // Only reconnect on connection errors, not parsing errors + // Skip reconnection for parse errors (including DID resolution failures) if (err.name === 'FirehoseParseError') { - console.warn('Skipping reconnection for parse error'); + // Check if it's a DID resolution timeout + if (err.cause && (err.cause as any).name === 'AbortError') { + console.warn('DID resolution timeout - continuing without reconnection'); + } else { + console.warn('Skipping reconnection for parse error'); + } return; } @@ -127,6 +160,22 @@ export class AtProtoFirehoseService implements IFirehoseService { } } + private startHeartbeat(): void { + this.heartbeatInterval = setInterval(() => { + if (DEBUG_LOGGING) { + console.log(`Firehose status: running=${this.isRunningFlag}, events processed in last minute: ${this.eventCount}`); + } + this.eventCount = 0; // Reset counter + }, 60000); + } + + private stopHeartbeat(): void { + if (this.heartbeatInterval) { + clearInterval(this.heartbeatInterval); + this.heartbeatInterval = undefined; + } + } + private getFilteredCollections(): string[] { const collections = this.configService.getAtProtoCollections(); return [ -- 2.51.2 From 6f8c792c880fcc2a199e4dffd578953da3736714 Mon Sep 17 00:00:00 2001 From: Wesley Finck Date: Mon, 24 Nov 2025 11:42:55 -0800 Subject: [PATCH 2/3] fix: handle undefined abortController in isRunning method Co-authored-by: aider (anthropic/claude-sonnet-4-20250514) --- .../atproto/infrastructure/services/AtProtoFirehoseService.ts | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/src/modules/atproto/infrastructure/services/AtProtoFirehoseService.ts b/src/modules/atproto/infrastructure/services/AtProtoFirehoseService.ts index f3c7c4d9..39e866fa 100644 --- a/src/modules/atproto/infrastructure/services/AtProtoFirehoseService.ts +++ b/src/modules/atproto/infrastructure/services/AtProtoFirehoseService.ts @@ -83,7 +83,9 @@ export class AtProtoFirehoseService implements IFirehoseService { } isRunning(): boolean { - return this.isRunningFlag && this.firehose && !(this.firehose as any).abortController?.signal.aborted; + return this.isRunningFlag && + this.firehose && + !((this.firehose as any).abortController?.signal?.aborted ?? false); } private async handleFirehoseEvent(evt: Event): Promise { -- 2.51.2 From 5e4812615a4ea3afcbe5a4c7feee66e3f2cfe685 Mon Sep 17 00:00:00 2001 From: Wesley Finck Date: Mon, 24 Nov 2025 11:48:46 -0800 Subject: [PATCH 3/3] formatting --- .../services/AtProtoFirehoseService.ts | 21 ++++++++++++------- 1 file changed, 14 insertions(+), 7 deletions(-) diff --git a/src/modules/atproto/infrastructure/services/AtProtoFirehoseService.ts b/src/modules/atproto/infrastructure/services/AtProtoFirehoseService.ts index 39e866fa..e6bc8948 100644 --- a/src/modules/atproto/infrastructure/services/AtProtoFirehoseService.ts +++ b/src/modules/atproto/infrastructure/services/AtProtoFirehoseService.ts @@ -49,7 +49,7 @@ export class AtProtoFirehoseService implements IFirehoseService { await this.firehose.start(); this.isRunningFlag = true; console.log('AT Protocol firehose service started'); - + // Start heartbeat logging this.startHeartbeat(); } catch (error) { @@ -83,15 +83,18 @@ export class AtProtoFirehoseService implements IFirehoseService { } isRunning(): boolean { - return this.isRunningFlag && - this.firehose && - !((this.firehose as any).abortController?.signal?.aborted ?? false); + return ( + (this.isRunningFlag && + this.firehose && + !((this.firehose as any).abortController?.signal?.aborted ?? false)) || + false + ); } private async handleFirehoseEvent(evt: Event): Promise { try { this.eventCount++; - + if (DEBUG_LOGGING) { console.log(`Processing firehose event: ${evt.event} for ${evt.did}`); } @@ -135,7 +138,9 @@ export class AtProtoFirehoseService implements IFirehoseService { if (err.name === 'FirehoseParseError') { // Check if it's a DID resolution timeout if (err.cause && (err.cause as any).name === 'AbortError') { - console.warn('DID resolution timeout - continuing without reconnection'); + console.warn( + 'DID resolution timeout - continuing without reconnection', + ); } else { console.warn('Skipping reconnection for parse error'); } @@ -165,7 +170,9 @@ export class AtProtoFirehoseService implements IFirehoseService { private startHeartbeat(): void { this.heartbeatInterval = setInterval(() => { if (DEBUG_LOGGING) { - console.log(`Firehose status: running=${this.isRunningFlag}, events processed in last minute: ${this.eventCount}`); + console.log( + `Firehose status: running=${this.isRunningFlag}, events processed in last minute: ${this.eventCount}`, + ); } this.eventCount = 0; // Reset counter }, 60000); -- 2.51.2