diff --git a/fly.development.toml b/fly.development.toml index e1e8f972..7031d388 100644 --- a/fly.development.toml +++ b/fly.development.toml @@ -16,6 +16,7 @@ primary_region = 'yyz' processes = ['web'] [[vm]] + processes = ['web'] memory = '512mb' cpu_kind = 'shared' cpus = 1 diff --git a/fly.production.toml b/fly.production.toml index 219cc56e..f9cf8bcc 100644 --- a/fly.production.toml +++ b/fly.production.toml @@ -17,6 +17,7 @@ primary_region = 'yyz' processes = ['web'] [[vm]] + processes = ['web'] memory = '1gb' cpu_kind = 'shared' cpus = 1 diff --git a/src/modules/cards/tests/integration/BullMQEventSystem.integration.test.ts b/src/modules/cards/tests/integration/BullMQEventSystem.integration.test.ts index 86f37c87..1681597f 100644 --- a/src/modules/cards/tests/integration/BullMQEventSystem.integration.test.ts +++ b/src/modules/cards/tests/integration/BullMQEventSystem.integration.test.ts @@ -3,13 +3,17 @@ import Redis from 'ioredis'; import { BullMQEventPublisher } from '../../../../shared/infrastructure/events/BullMQEventPublisher'; import { BullMQEventSubscriber } from '../../../../shared/infrastructure/events/BullMQEventSubscriber'; import { CardAddedToLibraryEvent } from '../../domain/events/CardAddedToLibraryEvent'; +import { CardAddedToCollectionEvent } from '../../domain/events/CardAddedToCollectionEvent'; import { CardId } from '../../domain/value-objects/CardId'; import { CuratorId } from '../../domain/value-objects/CuratorId'; +import { CollectionId } from '../../domain/value-objects/CollectionId'; import { IEventHandler } from '../../../../shared/application/events/IEventSubscriber'; import { ok, err } from '../../../../shared/core/Result'; import { EventNames } from '../../../../shared/infrastructure/events/EventConfig'; import { Queue } from 'bullmq'; import { QueueNames } from 'src/shared/infrastructure/events/QueueConfig'; +import { CardCollectionSaga } from '../../../feeds/application/sagas/CardCollectionSaga'; +import { RedisSagaStateStore } from '../../../feeds/infrastructure/RedisSagaStateStore'; describe('BullMQ Event System Integration', () => { let redisContainer: StartedRedisContainer; @@ -282,4 +286,100 @@ describe('BullMQ Event System Integration', () => { await eventsQueue.close(); }, 10000); }); + + describe('Redis-Based Saga Integration', () => { + it('should handle distributed saga state across multiple workers', async () => { + // Arrange - Create two saga instances (simulating multiple workers) + const mockUseCase = { + execute: jest.fn().mockResolvedValue(ok({ activityId: 'test-activity' })), + } as any; + + const stateStore = new RedisSagaStateStore(redis); + const saga1 = new CardCollectionSaga(mockUseCase, stateStore); + const saga2 = new CardCollectionSaga(mockUseCase, stateStore); + + // Create test events for same card/user (should be aggregated) + const cardId = CardId.createFromString('saga-test-card').unwrap(); + const curatorId = CuratorId.create('did:plc:sagatest').unwrap(); + + const libraryEvent = CardAddedToLibraryEvent.create( + cardId, + curatorId, + ).unwrap(); + const collectionEvent = CardAddedToCollectionEvent.create( + cardId, + CollectionId.createFromString('test-collection').unwrap(), + curatorId, + ).unwrap(); + + // Act - Process events with different saga instances + const result1 = await saga1.handleCardEvent(libraryEvent); + const result2 = await saga2.handleCardEvent(collectionEvent); + + // Assert - Both operations succeeded + expect(result1.isOk()).toBe(true); + expect(result2.isOk()).toBe(true); + + // Wait for aggregation window + await new Promise((resolve) => setTimeout(resolve, 3500)); + + // Assert - Only one aggregated activity was created + expect(mockUseCase.execute).toHaveBeenCalledTimes(1); + + const call = mockUseCase.execute.mock.calls[0][0]; + expect(call.cardId).toBe(cardId.getStringValue()); + expect(call.actorId).toBe(curatorId.value); + expect(call.collectionIds).toContain('test-collection'); + }, 15000); + }); + + describe('Multi-Queue Event Routing', () => { + it('should route events to multiple queues', async () => { + // Arrange - Create subscribers for different queues + const feedsSubscriber = new BullMQEventSubscriber(redis, { + queueName: QueueNames.FEEDS, + }); + const searchSubscriber = new BullMQEventSubscriber(redis, { + queueName: QueueNames.SEARCH, + }); + + const feedsHandler = { + handle: jest.fn().mockResolvedValue(ok(undefined)), + }; + const searchHandler = { + handle: jest.fn().mockResolvedValue(ok(undefined)), + }; + + await feedsSubscriber.subscribe( + EventNames.CARD_ADDED_TO_LIBRARY, + feedsHandler, + ); + await searchSubscriber.subscribe( + EventNames.CARD_ADDED_TO_LIBRARY, + searchHandler, + ); + + await feedsSubscriber.start(); + await searchSubscriber.start(); + + // Act - Publish single event + const event = CardAddedToLibraryEvent.create( + CardId.createFromString('multi-queue-card').unwrap(), + CuratorId.create('did:plc:multiuser').unwrap(), + ).unwrap(); + + await publisher.publishEvents([event]); + + // Wait for processing + await new Promise((resolve) => setTimeout(resolve, 2000)); + + // Assert - Event processed by both queues + expect(feedsHandler.handle).toHaveBeenCalledTimes(1); + expect(searchHandler.handle).toHaveBeenCalledTimes(1); + + // Cleanup + await feedsSubscriber.stop(); + await searchSubscriber.stop(); + }, 15000); + }); }); diff --git a/src/modules/feeds/application/sagas/CardCollectionSaga.ts b/src/modules/feeds/application/sagas/CardCollectionSaga.ts index dfbf12b7..48fe92b9 100644 --- a/src/modules/feeds/application/sagas/CardCollectionSaga.ts +++ b/src/modules/feeds/application/sagas/CardCollectionSaga.ts @@ -6,6 +6,7 @@ import { AddCardCollectedActivityDTO, } from '../useCases/commands/AddActivityToFeedUseCase'; import { ActivityTypeEnum } from '../../domain/value-objects/ActivityType'; +import { ISagaStateStore } from './ISagaStateStore'; interface PendingCardActivity { cardId: string; @@ -17,37 +18,84 @@ interface PendingCardActivity { } export class CardCollectionSaga { - private pendingActivities = new Map(); - private flushTimers = new Map(); private readonly AGGREGATION_WINDOW_MS = 3000; + private readonly REDIS_KEY_PREFIX = 'saga:feed'; - constructor(private addActivityToFeedUseCase: AddActivityToFeedUseCase) {} + constructor( + private addActivityToFeedUseCase: AddActivityToFeedUseCase, + private stateStore: ISagaStateStore, + ) {} async handleCardEvent( event: CardAddedToLibraryEvent | CardAddedToCollectionEvent, ): Promise> { + const aggregationKey = this.createKey(event); + + const lockAcquired = await this.acquireLock(aggregationKey); + if (!lockAcquired) { + return ok(undefined); // Another worker is processing + } + try { - const aggregationKey = this.createKey(event); - const existing = this.pendingActivities.get(aggregationKey); + const existing = await this.getPendingActivity(aggregationKey); if (existing && this.isWithinWindow(existing)) { - // Merge with existing activity this.mergeActivity(existing, event); - // Reset the timer - this.rescheduleFlush(aggregationKey); + await this.setPendingActivity(aggregationKey, existing); } else { - // Create new pending activity - this.createPendingActivity(aggregationKey, event); - this.scheduleFlush(aggregationKey); + const newActivity = this.createNewPendingActivity(event); + await this.setPendingActivity(aggregationKey, newActivity); + await this.scheduleFlush(aggregationKey); } return ok(undefined); - } catch (error) { - console.error('[SAGA] Error handling card event:', error); - return err(error as Error); + } finally { + await this.releaseLock(aggregationKey); } } + // Key helpers + private getPendingKey(aggregationKey: string): string { + return `${this.REDIS_KEY_PREFIX}:pending:${aggregationKey}`; + } + + private getLockKey(aggregationKey: string): string { + return `${this.REDIS_KEY_PREFIX}:lock:${aggregationKey}`; + } + + // State management + private async getPendingActivity( + aggregationKey: string, + ): Promise { + const data = await this.stateStore.get(this.getPendingKey(aggregationKey)); + return data ? JSON.parse(data) : null; + } + + private async setPendingActivity( + aggregationKey: string, + activity: PendingCardActivity, + ): Promise { + const key = this.getPendingKey(aggregationKey); + const ttlSeconds = Math.ceil(this.AGGREGATION_WINDOW_MS / 1000) + 5; + await this.stateStore.setex(key, ttlSeconds, JSON.stringify(activity)); + } + + private async deletePendingActivity(aggregationKey: string): Promise { + await this.stateStore.del(this.getPendingKey(aggregationKey)); + } + + // Distributed locking + private async acquireLock(aggregationKey: string): Promise { + const lockKey = this.getLockKey(aggregationKey); + const lockTtl = Math.ceil(this.AGGREGATION_WINDOW_MS / 1000) + 10; + const result = await this.stateStore.set(lockKey, '1', 'EX', lockTtl, 'NX'); + return result === 'OK'; + } + + private async releaseLock(aggregationKey: string): Promise { + await this.stateStore.del(this.getLockKey(aggregationKey)); + } + private createKey( event: CardAddedToLibraryEvent | CardAddedToCollectionEvent, ): string { @@ -72,10 +120,9 @@ export class CardCollectionSaga { return timeDiff <= this.AGGREGATION_WINDOW_MS; } - private createPendingActivity( - key: string, + private createNewPendingActivity( event: CardAddedToLibraryEvent | CardAddedToCollectionEvent, - ): void { + ): PendingCardActivity { const cardId = event.cardId.getStringValue(); const actorId = this.getActorId(event); @@ -89,7 +136,7 @@ export class CardCollectionSaga { }; this.mergeActivity(pending, event); - this.pendingActivities.set(key, pending); + return pending; } private mergeActivity( @@ -109,31 +156,20 @@ export class CardCollectionSaga { } } - private scheduleFlush(key: string): void { - const timer = setTimeout(() => { - this.flushActivity(key); + private async scheduleFlush(aggregationKey: string): Promise { + setTimeout(async () => { + await this.flushActivity(aggregationKey); }, 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; + private async flushActivity(aggregationKey: string): Promise { + const lockAcquired = await this.acquireLock(aggregationKey); + if (!lockAcquired) return; try { - // Create the aggregated activity + const pending = await this.getPendingActivity(aggregationKey); + if (!pending) return; + const request: AddCardCollectedActivityDTO = { type: ActivityTypeEnum.CARD_COLLECTED, actorId: pending.actorId, @@ -142,30 +178,10 @@ export class CardCollectionSaga { 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); + await this.addActivityToFeedUseCase.execute(request); } finally { - // Clean up - this.pendingActivities.delete(key); - this.flushTimers.delete(key); + await this.deletePendingActivity(aggregationKey); + await this.releaseLock(aggregationKey); } } - - // 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))); - } } diff --git a/src/modules/feeds/application/sagas/ISagaStateStore.ts b/src/modules/feeds/application/sagas/ISagaStateStore.ts new file mode 100644 index 00000000..88b158bd --- /dev/null +++ b/src/modules/feeds/application/sagas/ISagaStateStore.ts @@ -0,0 +1,28 @@ +export interface ISagaStateStore { + /** + * Get the value for a key + */ + get(key: string): Promise; + + /** + * Set a key-value pair with expiration in seconds + */ + setex(key: string, ttlSeconds: number, value: string): Promise; + + /** + * Delete a key + */ + del(key: string): Promise; + + /** + * Set a key with options (for distributed locking) + * Returns 'OK' if successful, null otherwise + */ + set( + key: string, + value: string, + mode: 'EX', + ttl: number, + flag: 'NX', + ): Promise<'OK' | null>; +} diff --git a/src/modules/feeds/infrastructure/RedisSagaStateStore.ts b/src/modules/feeds/infrastructure/RedisSagaStateStore.ts new file mode 100644 index 00000000..aabde728 --- /dev/null +++ b/src/modules/feeds/infrastructure/RedisSagaStateStore.ts @@ -0,0 +1,29 @@ +import Redis from 'ioredis'; +import { ISagaStateStore } from '../application/sagas/ISagaStateStore'; + +export class RedisSagaStateStore implements ISagaStateStore { + constructor(private redis: Redis) {} + + async get(key: string): Promise { + return this.redis.get(key); + } + + async setex(key: string, ttlSeconds: number, value: string): Promise { + await this.redis.setex(key, ttlSeconds, value); + } + + async del(key: string): Promise { + await this.redis.del(key); + } + + async set( + key: string, + value: string, + mode: 'EX', + ttl: number, + flag: 'NX', + ): Promise<'OK' | null> { + const result = await this.redis.set(key, value, mode, ttl, flag); + return result === 'OK' ? 'OK' : null; + } +} diff --git a/src/shared/infrastructure/events/BullMQEventPublisher.ts b/src/shared/infrastructure/events/BullMQEventPublisher.ts index 4014da51..38703315 100644 --- a/src/shared/infrastructure/events/BullMQEventPublisher.ts +++ b/src/shared/infrastructure/events/BullMQEventPublisher.ts @@ -51,21 +51,13 @@ export class BullMQEventPublisher implements IEventPublisher { } private getTargetQueues(eventName: EventName): QueueName[] { - // Route events to appropriate queues - // For now, all events go to feeds queue - // Future: route different events to different queues switch (eventName) { case EventNames.CARD_ADDED_TO_LIBRARY: - return [QueueNames.FEEDS]; - // Future: return [QueueNames.FEEDS, QueueNames.NOTIFICATIONS, QueueNames.ANALYTICS]; + return [QueueNames.FEEDS, QueueNames.SEARCH, QueueNames.ANALYTICS]; case EventNames.CARD_ADDED_TO_COLLECTION: return [QueueNames.FEEDS]; - // Future: return [QueueNames.FEEDS, QueueNames.ANALYTICS]; - case EventNames.COLLECTION_CREATED: - return [QueueNames.FEEDS]; - // Future: return [QueueNames.FEEDS, QueueNames.ANALYTICS]; default: - return []; // Default to feeds queue + return [QueueNames.FEEDS]; } } diff --git a/src/shared/infrastructure/events/QueueConfig.ts b/src/shared/infrastructure/events/QueueConfig.ts index 426f6ea0..5145b116 100644 --- a/src/shared/infrastructure/events/QueueConfig.ts +++ b/src/shared/infrastructure/events/QueueConfig.ts @@ -1,8 +1,7 @@ export const QueueNames = { FEEDS: 'feeds', - // Future queues can be added here: - // NOTIFICATIONS: 'notifications', - // ANALYTICS: 'analytics', + SEARCH: 'search', + ANALYTICS: 'analytics', } as const; export type QueueName = (typeof QueueNames)[keyof typeof QueueNames]; @@ -15,19 +14,18 @@ export const QueueOptions = { removeOnFail: 25, concurrency: 15, }, - // Future queue configurations: - // [QueueNames.NOTIFICATIONS]: { - // attempts: 5, - // backoff: { type: 'exponential' as const, delay: 1000 }, - // removeOnComplete: 100, - // removeOnFail: 50, - // concurrency: 5, - // }, - // [QueueNames.ANALYTICS]: { - // attempts: 2, - // backoff: { type: 'exponential' as const, delay: 5000 }, - // removeOnComplete: 25, - // removeOnFail: 10, - // concurrency: 20, - // }, + [QueueNames.SEARCH]: { + attempts: 3, + backoff: { type: 'exponential' as const, delay: 2000 }, + removeOnComplete: 50, + removeOnFail: 25, + concurrency: 10, + }, + [QueueNames.ANALYTICS]: { + attempts: 2, + backoff: { type: 'exponential' as const, delay: 5000 }, + removeOnComplete: 25, + removeOnFail: 10, + concurrency: 20, + }, } as const; diff --git a/src/shared/infrastructure/http/factories/ServiceFactory.ts b/src/shared/infrastructure/http/factories/ServiceFactory.ts index e7b136a7..ecb76436 100644 --- a/src/shared/infrastructure/http/factories/ServiceFactory.ts +++ b/src/shared/infrastructure/http/factories/ServiceFactory.ts @@ -222,18 +222,12 @@ export class ServiceFactory { }; } - // Create saga for worker - const cardCollectionSaga = new CardCollectionSaga( - // We'll need to create this use case in the worker context - null as any, // Will be set properly in worker - ); - return { ...sharedServices, redisConnection: redisConnection, eventPublisher, createEventSubscriber, - cardCollectionSaga, + cardCollectionSaga: null as any, // Will be created in worker process }; } diff --git a/src/shared/infrastructure/processes/BaseWorkerProcess.ts b/src/shared/infrastructure/processes/BaseWorkerProcess.ts new file mode 100644 index 00000000..9daf7f06 --- /dev/null +++ b/src/shared/infrastructure/processes/BaseWorkerProcess.ts @@ -0,0 +1,54 @@ +import { IProcess } from '../../domain/IProcess'; +import { EnvironmentConfigService } from '../config/EnvironmentConfigService'; +import { QueueName } from '../events/QueueConfig'; +import { IEventSubscriber } from '../../application/events/IEventSubscriber'; +import { RepositoryFactory, Repositories } from '../http/factories/RepositoryFactory'; +import { WorkerServices } from '../http/factories/ServiceFactory'; + +export abstract class BaseWorkerProcess implements IProcess { + constructor( + protected configService: EnvironmentConfigService, + protected queueName: QueueName, + ) {} + + async start(): Promise { + console.log(`Starting ${this.queueName} worker...`); + + const repositories = RepositoryFactory.create(this.configService); + const services = this.createServices(repositories); + + await this.validateDependencies(services); + + const eventSubscriber = services.createEventSubscriber(this.queueName); + await this.registerHandlers(eventSubscriber, services, repositories); + await eventSubscriber.start(); + + console.log(`${this.queueName} worker started`); + + this.setupShutdownHandlers(eventSubscriber, services); + } + + protected abstract createServices(repositories: Repositories): WorkerServices; + protected abstract validateDependencies(services: WorkerServices): Promise; + protected abstract registerHandlers( + subscriber: IEventSubscriber, + services: WorkerServices, + repositories: Repositories, + ): Promise; + + private setupShutdownHandlers( + subscriber: IEventSubscriber, + services: WorkerServices, + ): void { + const shutdown = async () => { + console.log(`Shutting down ${this.queueName} worker...`); + await subscriber.stop(); + if (services.redisConnection) { + await services.redisConnection.quit(); + } + process.exit(0); + }; + process.on('SIGTERM', shutdown); + process.on('SIGINT', shutdown); + } +} diff --git a/src/shared/infrastructure/processes/FeedWorkerProcess.ts b/src/shared/infrastructure/processes/FeedWorkerProcess.ts index 7196ced0..27819900 100644 --- a/src/shared/infrastructure/processes/FeedWorkerProcess.ts +++ b/src/shared/infrastructure/processes/FeedWorkerProcess.ts @@ -1,50 +1,45 @@ -import { IProcess } from '../../domain/IProcess'; import { EnvironmentConfigService } from '../config/EnvironmentConfigService'; -import { RepositoryFactory } from '../http/factories/RepositoryFactory'; -import { ServiceFactory } from '../http/factories/ServiceFactory'; +import { ServiceFactory, WorkerServices } from '../http/factories/ServiceFactory'; import { UseCaseFactory } from '../http/factories/UseCaseFactory'; import { CardAddedToLibraryEventHandler } from '../../../modules/feeds/application/eventHandlers/CardAddedToLibraryEventHandler'; import { CardAddedToCollectionEventHandler } from '../../../modules/feeds/application/eventHandlers/CardAddedToCollectionEventHandler'; import { CardCollectionSaga } from '../../../modules/feeds/application/sagas/CardCollectionSaga'; +import { RedisSagaStateStore } from '../../../modules/feeds/infrastructure/RedisSagaStateStore'; import { QueueNames } from '../events/QueueConfig'; import { EventNames } from '../events/EventConfig'; +import { BaseWorkerProcess } from './BaseWorkerProcess'; +import { IEventSubscriber } from '../../application/events/IEventSubscriber'; +import { Repositories } from '../http/factories/RepositoryFactory'; -export class FeedWorkerProcess implements IProcess { - constructor(private configService: EnvironmentConfigService) {} - - async start(): Promise { - console.log('Starting feed worker...'); +export class FeedWorkerProcess extends BaseWorkerProcess { + constructor(configService: EnvironmentConfigService) { + super(configService, QueueNames.FEEDS); + } - // Create dependencies using factories - const repositories = RepositoryFactory.create(this.configService); - const services = ServiceFactory.createForWorker( - this.configService, - repositories, - ); - const useCases = UseCaseFactory.createForWorker(repositories, services); + protected createServices(repositories: Repositories): WorkerServices { + return ServiceFactory.createForWorker(this.configService, repositories); + } - // Test Redis connection (only if using Redis) - if (services.redisConnection) { - try { - await services.redisConnection.ping(); - console.log('Connected to Redis successfully'); - } catch (error) { - console.error('Failed to connect to Redis:', error); - process.exit(1); - } - } else { - console.log('Using in-memory event system'); + protected async validateDependencies(services: WorkerServices): Promise { + if (!services.redisConnection) { + throw new Error('Redis connection required for feed worker'); } + await services.redisConnection.ping(); + } - // Create subscriber for feeds queue - const eventSubscriber = services.createEventSubscriber(QueueNames.FEEDS); + protected async registerHandlers( + subscriber: IEventSubscriber, + services: WorkerServices, + repositories: Repositories, + ): Promise { + const useCases = UseCaseFactory.createForWorker(repositories, services); - // Create the saga with proper dependencies + const stateStore = new RedisSagaStateStore(services.redisConnection!); const cardCollectionSaga = new CardCollectionSaga( useCases.addActivityToFeedUseCase, + stateStore, ); - // Create event handlers with the saga const cardAddedToLibraryHandler = new CardAddedToLibraryEventHandler( cardCollectionSaga, ); @@ -52,33 +47,14 @@ export class FeedWorkerProcess implements IProcess { cardCollectionSaga, ); - // Register handlers - await eventSubscriber.subscribe( + await subscriber.subscribe( EventNames.CARD_ADDED_TO_LIBRARY, cardAddedToLibraryHandler, ); - await eventSubscriber.subscribe( + await subscriber.subscribe( EventNames.CARD_ADDED_TO_COLLECTION, cardAddedToCollectionHandler, ); - - // Start the worker - await eventSubscriber.start(); - - console.log('Feed worker started and listening for events...'); - - // Graceful shutdown - const shutdown = async () => { - console.log('Shutting down feed worker...'); - await eventSubscriber.stop(); - if (services.redisConnection) { - await services.redisConnection.quit(); - } - process.exit(0); - }; - - process.on('SIGTERM', shutdown); - process.on('SIGINT', shutdown); } }