diff --git a/src/modules/feeds/application/eventHandlers/CardAddedToCollectionEventHandler.ts b/src/modules/feeds/application/eventHandlers/CardAddedToCollectionEventHandler.ts index 0836388e..69a510e3 100644 --- a/src/modules/feeds/application/eventHandlers/CardAddedToCollectionEventHandler.ts +++ b/src/modules/feeds/application/eventHandlers/CardAddedToCollectionEventHandler.ts @@ -1,46 +1,14 @@ import { CardAddedToCollectionEvent } from '../../../cards/domain/events/CardAddedToCollectionEvent'; import { IEventHandler } from '../../../../shared/application/events/IEventSubscriber'; -import { Result, ok, err } from '../../../../shared/core/Result'; -import { - AddActivityToFeedUseCase, - AddCardCollectedActivityDTO, -} from '../useCases/commands/AddActivityToFeedUseCase'; -import { ActivityTypeEnum } from '../../domain/value-objects/ActivityType'; +import { Result } from '../../../../shared/core/Result'; +import { CardCollectionSaga } from '../sagas/CardCollectionSaga'; export class CardAddedToCollectionEventHandler implements IEventHandler { - constructor(private addActivityToFeedUseCase: AddActivityToFeedUseCase) {} + constructor(private cardCollectionSaga: CardCollectionSaga) {} async handle(event: CardAddedToCollectionEvent): Promise> { - try { - const request: AddCardCollectedActivityDTO = { - type: ActivityTypeEnum.CARD_COLLECTED, - actorId: event.addedBy.value, - cardId: event.cardId.getStringValue(), - collectionIds: [event.collectionId.getStringValue()], - // TODO: Fetch card metadata (title, URL) from card repository - cardTitle: undefined, - cardUrl: undefined, - }; - - const result = await this.addActivityToFeedUseCase.execute(request); - - if (result.isErr()) { - console.error( - '[FEEDS] Failed to add card-added-to-collection activity:', - result.error, - ); - return err(new Error(result.error.message)); - } - - return ok(undefined); - } catch (error) { - console.error( - '[FEEDS] Unexpected error handling CardAddedToCollectionEvent:', - error, - ); - return err(error as Error); - } + return this.cardCollectionSaga.handleCardEvent(event); } } diff --git a/src/modules/feeds/application/eventHandlers/CardAddedToLibraryEventHandler.ts b/src/modules/feeds/application/eventHandlers/CardAddedToLibraryEventHandler.ts index 64af4d1f..fe35d3d1 100644 --- a/src/modules/feeds/application/eventHandlers/CardAddedToLibraryEventHandler.ts +++ b/src/modules/feeds/application/eventHandlers/CardAddedToLibraryEventHandler.ts @@ -1,47 +1,14 @@ import { CardAddedToLibraryEvent } from '../../../cards/domain/events/CardAddedToLibraryEvent'; import { IEventHandler } from '../../../../shared/application/events/IEventSubscriber'; -import { Result, ok, err } from '../../../../shared/core/Result'; -import { - AddActivityToFeedUseCase, - AddCardCollectedActivityDTO, -} from '../useCases/commands/AddActivityToFeedUseCase'; -import { ActivityTypeEnum } from '../../domain/value-objects/ActivityType'; +import { Result } from '../../../../shared/core/Result'; +import { CardCollectionSaga } from '../sagas/CardCollectionSaga'; export class CardAddedToLibraryEventHandler implements IEventHandler { - constructor(private addActivityToFeedUseCase: AddActivityToFeedUseCase) {} + constructor(private cardCollectionSaga: CardCollectionSaga) {} async handle(event: CardAddedToLibraryEvent): Promise> { - try { - const request: AddCardCollectedActivityDTO = { - type: ActivityTypeEnum.CARD_COLLECTED, - actorId: event.curatorId.value, - cardId: event.cardId.getStringValue(), - // No collection IDs for library-only additions - collectionIds: undefined, - }; - - const result = await this.addActivityToFeedUseCase.execute(request); - - if (result.isErr()) { - console.error( - '[FEEDS] Error processing CardAddedToLibraryEvent:', - result.error, - ); - return err(new Error(result.error.message)); - } - - console.log( - `[FEEDS] Successfully processed CardAddedToLibraryEvent for card ${event.cardId.getStringValue()}, created activity ${result.value.activityId}`, - ); - return ok(undefined); - } catch (error) { - console.error( - '[FEEDS] Unexpected error handling CardAddedToLibraryEvent:', - error, - ); - return err(error as Error); - } + return this.cardCollectionSaga.handleCardEvent(event); } } diff --git a/src/modules/feeds/application/sagas/CardCollectionSaga.ts b/src/modules/feeds/application/sagas/CardCollectionSaga.ts new file mode 100644 index 00000000..dfbf12b7 --- /dev/null +++ b/src/modules/feeds/application/sagas/CardCollectionSaga.ts @@ -0,0 +1,171 @@ +import { Result, ok, err } from '../../../../shared/core/Result'; +import { CardAddedToLibraryEvent } from '../../../cards/domain/events/CardAddedToLibraryEvent'; +import { CardAddedToCollectionEvent } from '../../../cards/domain/events/CardAddedToCollectionEvent'; +import { + AddActivityToFeedUseCase, + AddCardCollectedActivityDTO, +} from '../useCases/commands/AddActivityToFeedUseCase'; +import { ActivityTypeEnum } from '../../domain/value-objects/ActivityType'; + +interface PendingCardActivity { + cardId: string; + actorId: string; + collectionIds: string[]; + timestamp: Date; + hasLibraryEvent: boolean; + hasCollectionEvents: boolean; +} + +export class CardCollectionSaga { + private pendingActivities = new Map(); + private flushTimers = new Map(); + private readonly AGGREGATION_WINDOW_MS = 3000; + + constructor(private addActivityToFeedUseCase: AddActivityToFeedUseCase) {} + + async handleCardEvent( + event: CardAddedToLibraryEvent | CardAddedToCollectionEvent, + ): Promise> { + try { + const aggregationKey = this.createKey(event); + const existing = this.pendingActivities.get(aggregationKey); + + if (existing && this.isWithinWindow(existing)) { + // Merge with existing activity + this.mergeActivity(existing, event); + // Reset the timer + this.rescheduleFlush(aggregationKey); + } else { + // Create new pending activity + this.createPendingActivity(aggregationKey, event); + this.scheduleFlush(aggregationKey); + } + + return ok(undefined); + } catch (error) { + console.error('[SAGA] Error handling card event:', error); + return err(error as Error); + } + } + + private createKey( + event: CardAddedToLibraryEvent | CardAddedToCollectionEvent, + ): string { + const cardId = event.cardId.getStringValue(); + const actorId = this.getActorId(event); + return `${cardId}-${actorId}`; + } + + private getActorId( + event: CardAddedToLibraryEvent | CardAddedToCollectionEvent, + ): string { + if ('curatorId' in event) { + return event.curatorId.value; // CardAddedToLibraryEvent + } else { + return event.addedBy.value; // CardAddedToCollectionEvent + } + } + + private isWithinWindow(pending: PendingCardActivity): boolean { + const now = new Date(); + const timeDiff = now.getTime() - pending.timestamp.getTime(); + return timeDiff <= this.AGGREGATION_WINDOW_MS; + } + + private createPendingActivity( + key: string, + event: CardAddedToLibraryEvent | CardAddedToCollectionEvent, + ): void { + const cardId = event.cardId.getStringValue(); + const actorId = this.getActorId(event); + + const pending: PendingCardActivity = { + cardId, + actorId, + collectionIds: [], + timestamp: new Date(), + hasLibraryEvent: false, + hasCollectionEvents: false, + }; + + this.mergeActivity(pending, event); + this.pendingActivities.set(key, pending); + } + + private mergeActivity( + existing: PendingCardActivity, + event: CardAddedToLibraryEvent | CardAddedToCollectionEvent, + ): void { + if ('curatorId' in event) { + // CardAddedToLibraryEvent + existing.hasLibraryEvent = true; + } else { + // CardAddedToCollectionEvent + existing.hasCollectionEvents = true; + const collectionId = event.collectionId.getStringValue(); + if (!existing.collectionIds.includes(collectionId)) { + existing.collectionIds.push(collectionId); + } + } + } + + private scheduleFlush(key: string): void { + const timer = setTimeout(() => { + this.flushActivity(key); + }, this.AGGREGATION_WINDOW_MS); + + this.flushTimers.set(key, timer); + } + + private rescheduleFlush(key: string): void { + // Clear existing timer + const existingTimer = this.flushTimers.get(key); + if (existingTimer) { + clearTimeout(existingTimer); + } + + // Schedule new timer + this.scheduleFlush(key); + } + + private async flushActivity(key: string): Promise { + const pending = this.pendingActivities.get(key); + if (!pending) return; + + try { + // Create the aggregated activity + const request: AddCardCollectedActivityDTO = { + type: ActivityTypeEnum.CARD_COLLECTED, + actorId: pending.actorId, + cardId: pending.cardId, + collectionIds: + pending.collectionIds.length > 0 ? pending.collectionIds : undefined, + }; + + const result = await this.addActivityToFeedUseCase.execute(request); + + if (result.isErr()) { + console.error( + '[SAGA] Failed to create aggregated activity:', + result.error, + ); + } else { + console.log( + `[SAGA] Successfully created aggregated activity ${result.value.activityId} for card ${pending.cardId}`, + ); + } + } catch (error) { + console.error('[SAGA] Error flushing activity:', error); + } finally { + // Clean up + this.pendingActivities.delete(key); + this.flushTimers.delete(key); + } + } + + // For testing or graceful shutdown + public async flushAll(): Promise { + const keys = Array.from(this.pendingActivities.keys()); + await Promise.all(keys.map((key) => this.flushActivity(key))); + } +}