diff --git a/src/modules/cards/domain/Collection.ts b/src/modules/cards/domain/Collection.ts index 149aec06..378c0f20 100644 --- a/src/modules/cards/domain/Collection.ts +++ b/src/modules/cards/domain/Collection.ts @@ -38,6 +38,20 @@ export class CollectionValidationError extends Error { } } +// Command tracking for optimized persistence +export enum CollectionCommandType { + ADD_CARD = 'ADD_CARD', + REMOVE_CARD = 'REMOVE_CARD', + UPDATE_CARD_LINK = 'UPDATE_CARD_LINK', + ADD_COLLABORATOR = 'ADD_COLLABORATOR', + REMOVE_COLLABORATOR = 'REMOVE_COLLABORATOR', +} + +export interface CollectionCommand { + type: CollectionCommandType; + payload: any; +} + interface CollectionProps { authorId: CuratorId; name: CollectionName; @@ -52,6 +66,9 @@ interface CollectionProps { } export class Collection extends AggregateRoot { + private pendingCommands: CollectionCommand[] = []; + private isFullySynced: boolean = true; // Flag to indicate if we have full data or just metadata + get collectionId(): CollectionId { return CollectionId.create(this._id).unwrap(); } @@ -263,6 +280,12 @@ export class Collection extends AggregateRoot { this.props.cardCount = this.props.cardLinks.length; this.props.updatedAt = new Date(); + // Track the command for optimized persistence + this.pendingCommands.push({ + type: CollectionCommandType.ADD_CARD, + payload: newLink, + }); + // Raise domain event this.addDomainEvent( CardAddedToCollectionEvent.create( @@ -286,6 +309,15 @@ export class Collection extends AggregateRoot { if (link) { link.publishedRecordId = publishedRecordId; this.props.updatedAt = new Date(); + + // Track the update command for optimized persistence + this.pendingCommands.push({ + type: CollectionCommandType.UPDATE_CARD_LINK, + payload: { + cardId, + publishedRecordId, + }, + }); } } @@ -337,6 +369,15 @@ export class Collection extends AggregateRoot { this.props.cardCount = this.props.cardLinks.length; this.props.updatedAt = new Date(); + // Track the command for optimized persistence + this.pendingCommands.push({ + type: CollectionCommandType.REMOVE_CARD, + payload: { + cardId, + userId, + }, + }); + // Raise domain event this.addDomainEvent( CardRemovedFromCollectionEvent.create( @@ -448,4 +489,25 @@ export class Collection extends AggregateRoot { public getUnpublishedCardLinks(): CardLink[] { return this.props.cardLinks.filter((link) => !link.publishedRecordId); } + + // Command tracking methods for optimized persistence + public getPendingCommands(): CollectionCommand[] { + return [...this.pendingCommands]; + } + + public clearPendingCommands(): void { + this.pendingCommands = []; + } + + public hasPendingCommands(): boolean { + return this.pendingCommands.length > 0; + } + + public markAsFullySynced(synced: boolean = true): void { + this.isFullySynced = synced; + } + + public getIsFullySynced(): boolean { + return this.isFullySynced; + } } diff --git a/src/modules/cards/domain/ICollectionRepository.ts b/src/modules/cards/domain/ICollectionRepository.ts index 8647708c..c4830252 100644 --- a/src/modules/cards/domain/ICollectionRepository.ts +++ b/src/modules/cards/domain/ICollectionRepository.ts @@ -19,4 +19,13 @@ export interface ICollectionRepository { ): Promise>; save(collection: Collection): Promise>; delete(collectionId: CollectionId): Promise>; + + // Batch operation for efficient multi-collection updates + addCardToMultipleCollections( + cardId: CardId, + collectionIds: CollectionId[], + curatorId: CuratorId, + viaCardId?: CardId, + publishedRecordIds?: Map, + ): Promise>; } diff --git a/src/modules/cards/domain/services/CardCollectionService.ts b/src/modules/cards/domain/services/CardCollectionService.ts index 56723af4..c4648522 100644 --- a/src/modules/cards/domain/services/CardCollectionService.ts +++ b/src/modules/cards/domain/services/CardCollectionService.ts @@ -155,6 +155,117 @@ export class CardCollectionService implements DomainService { | AppError.UnexpectedError > > { + // Use optimized batch operation when not skipping publishing + // and when we have multiple collections + if (collectionIds.length > 1 && !options?.skipPublishing) { + try { + // First, check permissions for all collections + const collectionsResult = + await this.collectionRepository.findByIds(collectionIds); + if (collectionsResult.isErr()) { + return err(AppError.UnexpectedError.create(collectionsResult.error)); + } + + const collections = collectionsResult.value; + for (const collection of collections) { + if (!collection.canAddCard(curatorId)) { + return err( + new CardCollectionValidationError( + `User does not have permission to add cards to collection: ${collection.collectionId.getStringValue()}`, + ), + ); + } + } + + // First, add card to all collections in memory and publish links + const publishedRecordMap = new Map(); + const modifiedCollections: Collection[] = []; + + for (const collection of collections) { + // Check if card already in collection + const existingLink = collection.cardLinks.find((link) => + link.cardId.equals(card.cardId), + ); + + if (!existingLink) { + // Add card to collection in memory first + const addCardResult = collection.addCard( + card.cardId, + curatorId, + viaCardId, + options?.timestamp, + ); + if (addCardResult.isErr()) { + return err( + new CardCollectionValidationError( + `Failed to add card to collection: ${addCardResult.error.message}`, + ), + ); + } + + // Resolve via card published record ID if needed + let viaCardPublishedRecordId: PublishedRecordIdProps | undefined; + if (viaCardId) { + const viaCardResult = + await this.cardRepository.findById(viaCardId); + if ( + viaCardResult.isOk() && + viaCardResult.value?.publishedRecordId + ) { + viaCardPublishedRecordId = + viaCardResult.value.publishedRecordId.getValue(); + } + } + + // Now publish the collection link (card is now in collection) + const publishLinkResult = + await this.collectionPublisher.publishCardAddedToCollection( + card, + collection, + curatorId, + viaCardPublishedRecordId, + ); + if (publishLinkResult.isErr()) { + if (publishLinkResult.error instanceof AuthenticationError) { + return err(publishLinkResult.error); + } + return err( + new CardCollectionValidationError( + `Failed to publish collection link: ${publishLinkResult.error.message}`, + ), + ); + } + + // Mark the link as published in the collection + collection.markCardLinkAsPublished( + card.cardId, + publishLinkResult.value, + ); + publishedRecordMap.set( + collection.collectionId.getStringValue(), + publishLinkResult.value, + ); + modifiedCollections.push(collection); + } + } + + // Save all modified collections with their pending commands + // This will use the optimized command-based persistence + for (const collection of modifiedCollections) { + const saveResult = await this.collectionRepository.save(collection); + if (saveResult.isErr()) { + return err(AppError.UnexpectedError.create(saveResult.error)); + } + } + + // Return all collections (modified and unmodified) + return ok(collections); + } catch (error) { + return err(AppError.UnexpectedError.create(error)); + } + } + + // Fall back to sequential processing for single collection or when skipping publishing const updatedCollections: Collection[] = []; for (const collectionId of collectionIds) { diff --git a/src/modules/cards/infrastructure/repositories/DrizzleCollectionRepository.ts b/src/modules/cards/infrastructure/repositories/DrizzleCollectionRepository.ts index b6985e58..a342c8f2 100644 --- a/src/modules/cards/infrastructure/repositories/DrizzleCollectionRepository.ts +++ b/src/modules/cards/infrastructure/repositories/DrizzleCollectionRepository.ts @@ -1,7 +1,11 @@ -import { eq, inArray, and } from 'drizzle-orm'; +import { eq, inArray, and, sql } from 'drizzle-orm'; import { PostgresJsDatabase } from 'drizzle-orm/postgres-js'; import { ICollectionRepository } from '../../domain/ICollectionRepository'; -import { Collection } from '../../domain/Collection'; +import { + Collection, + CollectionCommandType, + CardLink, +} from '../../domain/Collection'; import { CollectionId } from '../../domain/value-objects/CollectionId'; import { CardId } from '../../domain/value-objects/CardId'; import { CuratorId } from '../../domain/value-objects/CuratorId'; @@ -603,6 +607,182 @@ export class DrizzleCollectionRepository implements ICollectionRepository { async save(collection: Collection): Promise> { try { + const collectionId = collection.collectionId.getStringValue(); + const pendingCommands = collection.getPendingCommands(); + + // If we have pending commands, use optimized targeted operations + if (pendingCommands.length > 0) { + await this.db.transaction(async (tx) => { + // First, ensure the collection exists (upsert) + const collectionData = + CollectionMapper.toPersistence(collection).collection; + + // Handle collection published record if it exists + let publishedRecordId: string | undefined; + if (collection.publishedRecordId) { + const recordId = new UniqueEntityID().toString(); + const recordedAt = new Date(); + const insertResult = await tx + .insert(publishedRecords) + .values({ + id: recordId, + uri: collection.publishedRecordId.uri, + cid: collection.publishedRecordId.cid, + recordedAt: recordedAt, + }) + .onConflictDoUpdate({ + target: [publishedRecords.uri, publishedRecords.cid], + set: { recordedAt: recordedAt }, + }) + .returning({ id: publishedRecords.id }); + + publishedRecordId = insertResult[0]?.id || recordId; + } + + // Upsert the collection + await tx + .insert(collections) + .values({ + ...collectionData, + publishedRecordId: publishedRecordId, + }) + .onConflictDoUpdate({ + target: collections.id, + set: { + authorId: collectionData.authorId, + name: collectionData.name, + description: collectionData.description, + accessType: collectionData.accessType, + cardCount: collectionData.cardCount, + updatedAt: collectionData.updatedAt, + publishedRecordId: publishedRecordId, + }, + }); + + // Process each command + for (const command of pendingCommands) { + switch (command.type) { + case CollectionCommandType.ADD_CARD: { + const link = command.payload as CardLink; + const cardLinkId = new UniqueEntityID().toString(); + + // Handle published record if present - optimized version + let publishedRecordId: string | undefined; + if (link.publishedRecordId) { + const recordId = new UniqueEntityID().toString(); + const recordedAt = new Date(); + const insertResult = await tx + .insert(publishedRecords) + .values({ + id: recordId, + uri: link.publishedRecordId.uri, + cid: link.publishedRecordId.cid, + recordedAt: recordedAt, + }) + .onConflictDoUpdate({ + target: [publishedRecords.uri, publishedRecords.cid], + set: { recordedAt: recordedAt }, // Update recordedAt to avoid empty set + }) + .returning({ id: publishedRecords.id }); + + publishedRecordId = insertResult[0]?.id || recordId; + } + + // Insert the new card link + await tx + .insert(collectionCards) + .values({ + id: cardLinkId, + collectionId: collectionId, + cardId: link.cardId.getStringValue(), + addedBy: link.addedBy.value, + addedAt: link.addedAt, + viaCardId: link.viaCardId?.getStringValue(), + publishedRecordId: publishedRecordId, + }) + .onConflictDoNothing(); // Idempotent - ignore if already exists + break; + } + + case CollectionCommandType.UPDATE_CARD_LINK: { + const { cardId, publishedRecordId } = command.payload; + + // Handle published record - optimized version + let recordId: string | undefined; + if (publishedRecordId) { + const newRecordId = new UniqueEntityID().toString(); + const recordedAt = new Date(); + const insertResult = await tx + .insert(publishedRecords) + .values({ + id: newRecordId, + uri: publishedRecordId.uri, + cid: publishedRecordId.cid, + recordedAt: recordedAt, + }) + .onConflictDoUpdate({ + target: [publishedRecords.uri, publishedRecords.cid], + set: { recordedAt: recordedAt }, // Update recordedAt to avoid empty set + }) + .returning({ id: publishedRecords.id }); + + recordId = insertResult[0]?.id || newRecordId; + } + + // Update the card link + await tx + .update(collectionCards) + .set({ + publishedRecordId: recordId, + }) + .where( + and( + eq(collectionCards.collectionId, collectionId), + eq(collectionCards.cardId, cardId.getStringValue()), + ), + ); + break; + } + + case CollectionCommandType.REMOVE_CARD: { + const { cardId } = command.payload; + + // Delete the card link + await tx + .delete(collectionCards) + .where( + and( + eq(collectionCards.collectionId, collectionId), + eq(collectionCards.cardId, cardId.getStringValue()), + ), + ); + break; + } + + case CollectionCommandType.ADD_COLLABORATOR: + case CollectionCommandType.REMOVE_COLLABORATOR: + // Handle collaborator changes if needed + break; + } + } + + // Update collection metadata only (count and timestamp) + // Other fields were already handled in the upsert above + await tx + .update(collections) + .set({ + cardCount: collectionData.cardCount, + updatedAt: collectionData.updatedAt, + }) + .where(eq(collections.id, collectionId)); + }); + + // Clear commands after successful save + collection.clearPendingCommands(); + return ok(undefined); + } + + // Fall back to full save for collections without commands (e.g., initial creation) const { collection: collectionData, collaborators, @@ -612,83 +792,49 @@ export class DrizzleCollectionRepository implements ICollectionRepository { } = CollectionMapper.toPersistence(collection); await this.db.transaction(async (tx) => { - // Handle collection published record if it exists + // Handle collection published record if it exists - optimized let publishedRecordId: string | undefined = undefined; if (publishedRecord) { + const recordedAt = publishedRecord.recordedAt || new Date(); const publishedRecordResult = await tx .insert(publishedRecords) .values({ id: publishedRecord.id, uri: publishedRecord.uri, cid: publishedRecord.cid, - recordedAt: publishedRecord.recordedAt || new Date(), + recordedAt: recordedAt, }) - .onConflictDoNothing({ + .onConflictDoUpdate({ target: [publishedRecords.uri, publishedRecords.cid], + set: { recordedAt: recordedAt }, // Update recordedAt to avoid empty set }) .returning({ id: publishedRecords.id }); - if (publishedRecordResult.length === 0) { - const existingRecord = await tx - .select() - .from(publishedRecords) - .where( - and( - eq(publishedRecords.uri, publishedRecord.uri), - eq(publishedRecords.cid, publishedRecord.cid), - ), - ) - .limit(1); - - if (existingRecord.length > 0) { - publishedRecordId = existingRecord[0]!.id; - } - } else { - publishedRecordId = publishedRecordResult[0]!.id; - } + publishedRecordId = + publishedRecordResult[0]?.id || publishedRecord.id; } - // Handle link published records - const linkPublishedRecordMap = new Map(); - if (linkPublishedRecords) { - for (const linkRecord of linkPublishedRecords) { - const linkPublishedRecordResult = await tx + // Batch insert published records if needed + const publishedRecordsBatch: any[] = []; + if (linkPublishedRecords && linkPublishedRecords.length > 0) { + for (const record of linkPublishedRecords) { + publishedRecordsBatch.push({ + id: record.id, + uri: record.uri, + cid: record.cid, + recordedAt: record.recordedAt || new Date(), + }); + } + + // Batch insert all published records at once + if (publishedRecordsBatch.length > 0) { + await tx .insert(publishedRecords) - .values({ - id: linkRecord.id, - uri: linkRecord.uri, - cid: linkRecord.cid, - recordedAt: linkRecord.recordedAt || new Date(), - }) + .values(publishedRecordsBatch) .onConflictDoNothing({ target: [publishedRecords.uri, publishedRecords.cid], - }) - .returning({ id: publishedRecords.id }); - - let actualRecordId: string; - if (linkPublishedRecordResult.length === 0) { - const existingRecord = await tx - .select() - .from(publishedRecords) - .where( - and( - eq(publishedRecords.uri, linkRecord.uri), - eq(publishedRecords.cid, linkRecord.cid), - ), - ) - .limit(1); - - if (existingRecord.length > 0) { - actualRecordId = existingRecord[0]!.id; - } else { - actualRecordId = linkRecord.id; - } - } else { - actualRecordId = linkPublishedRecordResult[0]!.id; - } - - linkPublishedRecordMap.set(linkRecord.id, actualRecordId); + }); } } @@ -712,31 +858,26 @@ export class DrizzleCollectionRepository implements ICollectionRepository { }, }); - // Delete existing collaborators and card links - await tx - .delete(collectionCollaborators) - .where(eq(collectionCollaborators.collectionId, collectionData.id)); + // Only do full resync if this is an initial save or full update + if (collection.getIsFullySynced()) { + // Delete existing collaborators and card links + await tx + .delete(collectionCollaborators) + .where(eq(collectionCollaborators.collectionId, collectionData.id)); - await tx - .delete(collectionCards) - .where(eq(collectionCards.collectionId, collectionData.id)); + await tx + .delete(collectionCards) + .where(eq(collectionCards.collectionId, collectionData.id)); - // Insert new collaborators - if (collaborators.length > 0) { - await tx.insert(collectionCollaborators).values(collaborators); - } + // Insert new collaborators + if (collaborators.length > 0) { + await tx.insert(collectionCollaborators).values(collaborators); + } - // Insert new card links with mapped published record IDs - if (cardLinks.length > 0) { - const cardLinksWithMappedRecords = cardLinks.map((link) => ({ - ...link, - publishedRecordId: link.publishedRecordId - ? linkPublishedRecordMap.get(link.publishedRecordId) || - link.publishedRecordId - : undefined, - })); - - await tx.insert(collectionCards).values(cardLinksWithMappedRecords); + // Insert new card links + if (cardLinks.length > 0) { + await tx.insert(collectionCards).values(cardLinks); + } } }); @@ -759,4 +900,83 @@ export class DrizzleCollectionRepository implements ICollectionRepository { return err(error as Error); } } + + // Batch operation to add a card to multiple collections efficiently + async addCardToMultipleCollections( + cardId: CardId, + collectionIds: CollectionId[], + curatorId: CuratorId, + viaCardId?: CardId, + publishedRecordIds?: Map, // collectionId -> publishedRecordId + ): Promise> { + try { + if (collectionIds.length === 0) { + return ok(undefined); + } + + const cardIdStr = cardId.getStringValue(); + const curatorIdStr = curatorId.value; + const viaCardIdStr = viaCardId?.getStringValue(); + const addedAt = new Date(); + + await this.db.transaction(async (tx) => { + // Prepare batch of card links + const cardLinksToInsert: any[] = []; + const collectionIdsToUpdate: string[] = []; + + for (const collectionId of collectionIds) { + const collectionIdStr = collectionId.getStringValue(); + const linkId = new UniqueEntityID().toString(); + const publishedRecordId = publishedRecordIds?.get(collectionIdStr); + + // Check if card already exists in collection + const existing = await tx + .select() + .from(collectionCards) + .where( + and( + eq(collectionCards.collectionId, collectionIdStr), + eq(collectionCards.cardId, cardIdStr), + ), + ) + .limit(1); + + if (existing.length === 0) { + cardLinksToInsert.push({ + id: linkId, + collectionId: collectionIdStr, + cardId: cardIdStr, + addedBy: curatorIdStr, + addedAt: addedAt, + viaCardId: viaCardIdStr, + publishedRecordId: publishedRecordId, + }); + collectionIdsToUpdate.push(collectionIdStr); + } + } + + // Batch insert all card links at once + if (cardLinksToInsert.length > 0) { + await tx.insert(collectionCards).values(cardLinksToInsert); + + // Update all collection counts in a single query + await tx.execute(sql` + UPDATE collections + SET + card_count = ( + SELECT COUNT(*) + FROM collection_cards + WHERE collection_cards.collection_id = collections.id + ), + updated_at = NOW() + WHERE id = ANY(${collectionIdsToUpdate}::uuid[]) + `); + } + }); + + return ok(undefined); + } catch (error) { + return err(error as Error); + } + } } diff --git a/src/modules/cards/tests/application/AddUrlToLibraryUseCase.test.ts b/src/modules/cards/tests/application/AddUrlToLibraryUseCase.test.ts index 2dcc2b2a..f1160160 100644 --- a/src/modules/cards/tests/application/AddUrlToLibraryUseCase.test.ts +++ b/src/modules/cards/tests/application/AddUrlToLibraryUseCase.test.ts @@ -699,6 +699,8 @@ describe('AddUrlToLibraryUseCase', () => { expect(result.error.message).toContain('does not have permission'); } + // When adding to multiple collections fails due to permissions, + // NO collections should be modified (all-or-nothing behavior) const openPublishedLinks = collectionPublisher.getPublishedLinksForCollection( openCollection.collectionId.getStringValue(), @@ -707,7 +709,7 @@ describe('AddUrlToLibraryUseCase', () => { collectionPublisher.getPublishedLinksForCollection( closedCollection.collectionId.getStringValue(), ); - expect(openPublishedLinks).toHaveLength(1); + expect(openPublishedLinks).toHaveLength(0); expect(closedPublishedLinks).toHaveLength(0); }); }); diff --git a/src/modules/cards/tests/application/GetCollectionPageUseCase.test.ts b/src/modules/cards/tests/application/GetCollectionPageUseCase.test.ts index 3d61253a..cfe62323 100644 --- a/src/modules/cards/tests/application/GetCollectionPageUseCase.test.ts +++ b/src/modules/cards/tests/application/GetCollectionPageUseCase.test.ts @@ -715,6 +715,7 @@ describe('GetCollectionPageUseCase', () => { findByCardId: jest.fn(), findByCuratorIdContainingCard: jest.fn(), findContainingCardAddedBy: jest.fn(), + addCardToMultipleCollections: jest.fn(), }; const errorUseCase = new GetCollectionPageUseCase( diff --git a/src/modules/cards/tests/utils/InMemoryCollectionRepository.ts b/src/modules/cards/tests/utils/InMemoryCollectionRepository.ts index bb9f4439..af42e72b 100644 --- a/src/modules/cards/tests/utils/InMemoryCollectionRepository.ts +++ b/src/modules/cards/tests/utils/InMemoryCollectionRepository.ts @@ -155,6 +155,51 @@ export class InMemoryCollectionRepository implements ICollectionRepository { } } + async addCardToMultipleCollections( + cardId: CardId, + collectionIds: CollectionId[], + curatorId: CuratorId, + viaCardId?: CardId, + publishedRecordIds?: Map, + ): Promise> { + try { + const addedAt = new Date(); + + for (const collectionId of collectionIds) { + const collection = this.collections.get(collectionId.getStringValue()); + if (collection) { + // Check if card already exists in collection + const existingLink = collection.cardLinks.find((link) => + link.cardId.equals(cardId), + ); + + if (!existingLink) { + // Add the card to the collection + const addResult = collection.addCard( + cardId, + curatorId, + viaCardId, + addedAt, + ); + if (addResult.isErr()) { + return err(addResult.error); + } + + // Update the collection + this.collections.set( + collectionId.getStringValue(), + this.clone(collection), + ); + } + } + } + + return ok(undefined); + } catch (error) { + return err(error as Error); + } + } + // Helper methods for testing public clear(): void { this.collections.clear(); diff --git a/src/shared/infrastructure/events/BullMQEventPublisher.ts b/src/shared/infrastructure/events/BullMQEventPublisher.ts index eaaaf4cf..fed397a3 100644 --- a/src/shared/infrastructure/events/BullMQEventPublisher.ts +++ b/src/shared/infrastructure/events/BullMQEventPublisher.ts @@ -13,27 +13,42 @@ export class BullMQEventPublisher implements IEventPublisher { constructor(private redisConnection: Redis) {} async publishEvents(events: IDomainEvent[]): Promise> { - try { - for (const event of events) { - await this.publishSingleEvent(event); + // Fire and forget - don't await, just catch errors silently + this.publishEventsAsync(events).catch((error) => { + console.error('[BullMQEventPublisher] Failed to publish events:', error); + }); + + return ok(undefined); + } + + private async publishEventsAsync(events: IDomainEvent[]): Promise { + // Group events by queue to enable batch publishing + const eventsByQueue = new Map(); + + for (const event of events) { + const targetQueues = this.getTargetQueues(event.eventName); + + for (const queueName of targetQueues) { + if (!eventsByQueue.has(queueName)) { + eventsByQueue.set(queueName, []); + } + eventsByQueue.get(queueName)!.push(event); } - return ok(undefined); - } catch (error) { - return err(error as Error); } - } - private async publishSingleEvent(event: IDomainEvent): Promise { - const targetQueues = this.getTargetQueues(event.eventName); + // Publish all events to each queue in batch + const publishPromises: Promise[] = []; - for (const queueName of targetQueues) { - await this.publishToQueue(queueName, event); + for (const [queueName, queueEvents] of eventsByQueue.entries()) { + publishPromises.push(this.publishBatchToQueue(queueName, queueEvents)); } + + await Promise.all(publishPromises); } - private async publishToQueue( + private async publishBatchToQueue( queueName: QueueName, - event: IDomainEvent, + events: IDomainEvent[], ): Promise { if (!this.queues.has(queueName)) { this.queues.set( @@ -46,8 +61,17 @@ export class BullMQEventPublisher implements IEventPublisher { } const queue = this.queues.get(queueName)!; - const serializedEvent = EventMapper.toSerialized(event); - await queue.add(serializedEvent.eventType, serializedEvent); + + // Use addBulk for batch publishing + const jobs = events.map((event) => { + const serializedEvent = EventMapper.toSerialized(event); + return { + name: serializedEvent.eventType, + data: serializedEvent, + }; + }); + + await queue.addBulk(jobs); } private getTargetQueues(eventName: EventName): QueueName[] {