diff --git a/src/modules/atproto/application/useCases/ProcessCardFirehoseEventUseCase.ts b/src/modules/atproto/application/useCases/ProcessCardFirehoseEventUseCase.ts index 478b5fe5..310cb85f 100644 --- a/src/modules/atproto/application/useCases/ProcessCardFirehoseEventUseCase.ts +++ b/src/modules/atproto/application/useCases/ProcessCardFirehoseEventUseCase.ts @@ -86,6 +86,11 @@ export class ProcessCardFirehoseEventUseCase const atUri = atUriResult.value; const curatorDid = atUri.did.value; + // Extract timestamp from AT Protocol record + const timestamp = request.record.createdAt + ? new Date(request.record.createdAt) + : undefined; + const publishedRecordId = PublishedRecordId.create({ uri: request.atUri, cid: request.cid, @@ -118,6 +123,7 @@ export class ProcessCardFirehoseEventUseCase curatorId: curatorDid, publishedRecordId: publishedRecordId, viaCardId: viaCardId?.getStringValue(), + timestamp: timestamp, }); if (result.isErr()) { @@ -235,6 +241,7 @@ export class ProcessCardFirehoseEventUseCase publishedRecordIds: { noteCard: publishedRecordId, }, + timestamp: timestamp, }); if (result.isErr()) { @@ -332,6 +339,11 @@ export class ProcessCardFirehoseEventUseCase return ok(undefined); } + // Extract timestamp from AT Protocol record + const timestamp = request.record.createdAt + ? new Date(request.record.createdAt) + : undefined; + const publishedRecordId = PublishedRecordId.create({ uri: request.atUri, cid: request.cid, @@ -345,6 +357,7 @@ export class ProcessCardFirehoseEventUseCase publishedRecordIds: { noteCard: publishedRecordId, }, + timestamp: timestamp, }); if (result.isErr()) { @@ -407,6 +420,7 @@ export class ProcessCardFirehoseEventUseCase ); } + // For delete events, we don't have a record, so no timestamp available const publishedRecordId = PublishedRecordId.create({ uri: request.atUri, cid: request.cid || 'deleted', diff --git a/src/modules/atproto/application/useCases/ProcessCollectionFirehoseEventUseCase.ts b/src/modules/atproto/application/useCases/ProcessCollectionFirehoseEventUseCase.ts index 5cdd6de1..06658151 100644 --- a/src/modules/atproto/application/useCases/ProcessCollectionFirehoseEventUseCase.ts +++ b/src/modules/atproto/application/useCases/ProcessCollectionFirehoseEventUseCase.ts @@ -76,6 +76,11 @@ export class ProcessCollectionFirehoseEventUseCase } const authorDid = atUriResult.value.did.value; + // Extract timestamp from AT Protocol record + const timestamp = request.record.createdAt + ? new Date(request.record.createdAt) + : undefined; + const publishedRecordId = PublishedRecordId.create({ uri: request.atUri, cid: request.cid, @@ -86,6 +91,7 @@ export class ProcessCollectionFirehoseEventUseCase description: request.record.description, curatorId: authorDid, publishedRecordId: publishedRecordId, + createdAt: timestamp, }); if (result.isErr()) { @@ -159,6 +165,11 @@ export class ProcessCollectionFirehoseEventUseCase return ok(undefined); } + // Extract timestamp from AT Protocol record for the published record + const timestamp = request.record.createdAt + ? new Date(request.record.createdAt) + : undefined; + const publishedRecordId = PublishedRecordId.create({ uri: request.atUri, cid: request.cid, @@ -231,6 +242,7 @@ export class ProcessCollectionFirehoseEventUseCase ); } + // For delete events, we don't have a record, so no timestamp available const publishedRecordId = PublishedRecordId.create({ uri: request.atUri, cid: request.cid || 'deleted', diff --git a/src/modules/atproto/application/useCases/ProcessCollectionLinkFirehoseEventUseCase.ts b/src/modules/atproto/application/useCases/ProcessCollectionLinkFirehoseEventUseCase.ts index 74856ab0..b13bd971 100644 --- a/src/modules/atproto/application/useCases/ProcessCollectionLinkFirehoseEventUseCase.ts +++ b/src/modules/atproto/application/useCases/ProcessCollectionLinkFirehoseEventUseCase.ts @@ -109,6 +109,14 @@ export class ProcessCollectionLinkFirehoseEventUseCase return ok(undefined); } + // Extract timestamps from AT Protocol record + // Use createdAt for the published record timestamp + const recordedAt = request.record.createdAt + ? new Date(request.record.createdAt) + : undefined; + // Use addedAt for the collection link timestamp (when card was added to collection) + const timestamp = new Date(request.record.addedAt); + const publishedRecordId = PublishedRecordId.create({ uri: request.atUri, cid: request.cid, @@ -138,6 +146,7 @@ export class ProcessCollectionLinkFirehoseEventUseCase collectionLinks: collectionLinkMap, }, viaCardId: viaCardId?.getStringValue(), + timestamp: timestamp, }); if (result.isErr()) { @@ -202,6 +211,7 @@ export class ProcessCollectionLinkFirehoseEventUseCase ); } + // For delete events, we don't have a record, so no timestamp available const publishedRecordId = PublishedRecordId.create({ uri: request.atUri, cid: request.cid || 'deleted', diff --git a/src/modules/atproto/application/useCases/ProcessMarginBookmarkFirehoseEventUseCase.ts b/src/modules/atproto/application/useCases/ProcessMarginBookmarkFirehoseEventUseCase.ts index e60cf5ac..7ce4dd03 100644 --- a/src/modules/atproto/application/useCases/ProcessMarginBookmarkFirehoseEventUseCase.ts +++ b/src/modules/atproto/application/useCases/ProcessMarginBookmarkFirehoseEventUseCase.ts @@ -94,6 +94,9 @@ export class ProcessMarginBookmarkFirehoseEventUseCase return ok(undefined); } + // Extract timestamp from AT Protocol record (Margin bookmark has required createdAt) + const timestamp = new Date(request.record.createdAt); + const publishedRecordId = PublishedRecordId.create({ uri: request.atUri, cid: request.cid, @@ -104,6 +107,7 @@ export class ProcessMarginBookmarkFirehoseEventUseCase curatorId: curatorDid, publishedRecordId: publishedRecordId, viaCardId: undefined, // Margin bookmarks don't have 'via' references + timestamp: timestamp, }); if (result.isErr()) { @@ -167,6 +171,7 @@ export class ProcessMarginBookmarkFirehoseEventUseCase ); } + // For delete events, we don't have a record, so no timestamp available const publishedRecordId = PublishedRecordId.create({ uri: request.atUri, cid: request.cid || 'deleted', diff --git a/src/modules/atproto/application/useCases/ProcessMarginCollectionFirehoseEventUseCase.ts b/src/modules/atproto/application/useCases/ProcessMarginCollectionFirehoseEventUseCase.ts index f51f8660..49368ce2 100644 --- a/src/modules/atproto/application/useCases/ProcessMarginCollectionFirehoseEventUseCase.ts +++ b/src/modules/atproto/application/useCases/ProcessMarginCollectionFirehoseEventUseCase.ts @@ -78,6 +78,9 @@ export class ProcessMarginCollectionFirehoseEventUseCase } const authorDid = atUriResult.value.did.value; + // Extract timestamp from AT Protocol record (Margin collection has required createdAt) + const timestamp = new Date(request.record.createdAt); + const publishedRecordId = PublishedRecordId.create({ uri: request.atUri, cid: request.cid, @@ -89,6 +92,7 @@ export class ProcessMarginCollectionFirehoseEventUseCase description: request.record.description, curatorId: authorDid, publishedRecordId: publishedRecordId, + createdAt: timestamp, }); if (result.isErr()) { @@ -162,6 +166,9 @@ export class ProcessMarginCollectionFirehoseEventUseCase return ok(undefined); } + // Extract timestamp from AT Protocol record for the published record + const timestamp = new Date(request.record.createdAt); + const publishedRecordId = PublishedRecordId.create({ uri: request.atUri, cid: request.cid, @@ -234,6 +241,7 @@ export class ProcessMarginCollectionFirehoseEventUseCase ); } + // For delete events, we don't have a record, so no timestamp available const publishedRecordId = PublishedRecordId.create({ uri: request.atUri, cid: request.cid || 'deleted', diff --git a/src/modules/atproto/application/useCases/ProcessMarginCollectionItemFirehoseEventUseCase.ts b/src/modules/atproto/application/useCases/ProcessMarginCollectionItemFirehoseEventUseCase.ts index d9413ee8..93fe6602 100644 --- a/src/modules/atproto/application/useCases/ProcessMarginCollectionItemFirehoseEventUseCase.ts +++ b/src/modules/atproto/application/useCases/ProcessMarginCollectionItemFirehoseEventUseCase.ts @@ -134,6 +134,9 @@ export class ProcessMarginCollectionItemFirehoseEventUseCase return ok(undefined); } + // Extract timestamp from AT Protocol record (Margin collection item has required createdAt) + const timestamp = new Date(request.record.createdAt); + const publishedRecordId = PublishedRecordId.create({ uri: request.atUri, cid: request.cid, @@ -155,6 +158,7 @@ export class ProcessMarginCollectionItemFirehoseEventUseCase collectionLinks: collectionLinkMap, }, viaCardId: undefined, // Margin doesn't have 'via' provenance + timestamp: timestamp, }); if (result.isErr()) { @@ -219,6 +223,7 @@ export class ProcessMarginCollectionItemFirehoseEventUseCase ); } + // For delete events, we don't have a record, so no timestamp available const publishedRecordId = PublishedRecordId.create({ uri: request.atUri, cid: request.cid || 'deleted', diff --git a/src/modules/cards/application/useCases/commands/AddUrlToLibraryUseCase.ts b/src/modules/cards/application/useCases/commands/AddUrlToLibraryUseCase.ts index 979866c4..a2cefd41 100644 --- a/src/modules/cards/application/useCases/commands/AddUrlToLibraryUseCase.ts +++ b/src/modules/cards/application/useCases/commands/AddUrlToLibraryUseCase.ts @@ -29,6 +29,7 @@ export interface AddUrlToLibraryDTO { curatorId: string; viaCardId?: string; publishedRecordId?: PublishedRecordId; // For firehose events - skip publishing if provided + timestamp?: Date; // For firehose events - use historical timestamp from AT Protocol record } export interface AddUrlToLibraryResponseDTO { @@ -121,6 +122,7 @@ export class AddUrlToLibraryUseCase extends BaseUseCase< const urlCardResult = CardFactory.create({ curatorId: request.curatorId, cardInput: urlCardInput, + createdAt: request.timestamp, }); if (urlCardResult.isErr()) { @@ -135,8 +137,11 @@ export class AddUrlToLibraryUseCase extends BaseUseCase< ? { skipPublishing: true, publishedRecordId: request.publishedRecordId, + timestamp: request.timestamp, } - : undefined; + : request.timestamp + ? { timestamp: request.timestamp } + : undefined; const addUrlCardToLibraryResult = await this.cardLibraryService.addCardToLibrary( @@ -222,6 +227,7 @@ export class AddUrlToLibraryUseCase extends BaseUseCase< const noteCardResult = CardFactory.create({ curatorId: request.curatorId, cardInput: noteCardInput, + createdAt: request.timestamp, }); if (noteCardResult.isErr()) { @@ -233,8 +239,15 @@ export class AddUrlToLibraryUseCase extends BaseUseCase< // Add note card to library using domain service // Note: For note cards, we don't pass publishedRecordId here since it's handled // separately in UpdateUrlCardAssociationsUseCase for firehose events + const noteLibraryOptions = request.timestamp + ? { timestamp: request.timestamp } + : undefined; const addNoteCardToLibraryResult = - await this.cardLibraryService.addCardToLibrary(noteCard, curatorId); + await this.cardLibraryService.addCardToLibrary( + noteCard, + curatorId, + noteLibraryOptions, + ); if (addNoteCardToLibraryResult.isErr()) { // Propagate authentication errors if ( @@ -293,12 +306,16 @@ export class AddUrlToLibraryUseCase extends BaseUseCase< } // Add card to collections using domain service + const collectionOptions = request.timestamp + ? { timestamp: request.timestamp } + : undefined; const addToCollectionsResult = await this.cardCollectionService.addCardToCollections( cardToAdd, collectionIds, curatorId, viaCardId, + collectionOptions, ); if (addToCollectionsResult.isErr()) { // Propagate authentication errors diff --git a/src/modules/cards/application/useCases/commands/CreateCollectionUseCase.ts b/src/modules/cards/application/useCases/commands/CreateCollectionUseCase.ts index 72b15019..78a2b282 100644 --- a/src/modules/cards/application/useCases/commands/CreateCollectionUseCase.ts +++ b/src/modules/cards/application/useCases/commands/CreateCollectionUseCase.ts @@ -14,6 +14,7 @@ export interface CreateCollectionDTO { description?: string; curatorId: string; publishedRecordId?: PublishedRecordId; // For firehose events - skip publishing if provided + createdAt?: Date; // For firehose events - use historical timestamp from AT Protocol record } export interface CreateCollectionResponseDTO { @@ -62,14 +63,15 @@ export class CreateCollectionUseCase const curatorId = curatorIdResult.value; // Create collection + const timestamp = request.createdAt ?? new Date(); const collectionResult = Collection.create({ authorId: curatorId, name: request.name, description: request.description, accessType: CollectionAccessType.CLOSED, collaboratorIds: [], - createdAt: new Date(), - updatedAt: new Date(), + createdAt: timestamp, + updatedAt: timestamp, }); if (collectionResult.isErr()) { diff --git a/src/modules/cards/application/useCases/commands/UpdateUrlCardAssociationsUseCase.ts b/src/modules/cards/application/useCases/commands/UpdateUrlCardAssociationsUseCase.ts index aa0959eb..37d263ff 100644 --- a/src/modules/cards/application/useCases/commands/UpdateUrlCardAssociationsUseCase.ts +++ b/src/modules/cards/application/useCases/commands/UpdateUrlCardAssociationsUseCase.ts @@ -34,6 +34,7 @@ export interface UpdateUrlCardAssociationsDTO { noteCard?: PublishedRecordId; collectionLinks?: Map; }; + timestamp?: Date; // For firehose events - use historical timestamp from AT Protocol record } export interface UpdateUrlCardAssociationsResponseDTO { @@ -172,8 +173,11 @@ export class UpdateUrlCardAssociationsUseCase extends BaseUseCase< ? { skipPublishing: true, publishedRecordId: request.publishedRecordIds.noteCard, + timestamp: request.timestamp, } - : undefined; + : request.timestamp + ? { timestamp: request.timestamp } + : undefined; // Update note card in library (handles save and republish) const updateNoteResult = @@ -207,6 +211,7 @@ export class UpdateUrlCardAssociationsUseCase extends BaseUseCase< const noteCardResult = CardFactory.create({ curatorId: request.curatorId, cardInput: noteCardInput, + createdAt: request.timestamp, }); if (noteCardResult.isErr()) { @@ -223,8 +228,11 @@ export class UpdateUrlCardAssociationsUseCase extends BaseUseCase< ? { skipPublishing: true, publishedRecordId: request.publishedRecordIds.noteCard, + timestamp: request.timestamp, } - : undefined; + : request.timestamp + ? { timestamp: request.timestamp } + : undefined; // Add note card to library using domain service const addNoteCardToLibraryResult = @@ -283,8 +291,11 @@ export class UpdateUrlCardAssociationsUseCase extends BaseUseCase< ? { skipPublishing: true, publishedRecordIds: request.publishedRecordIds.collectionLinks, + timestamp: request.timestamp, } - : undefined; + : request.timestamp + ? { timestamp: request.timestamp } + : undefined; // Validate and create viaCardId if provided let viaCardId: CardId | undefined; diff --git a/src/modules/cards/domain/Card.ts b/src/modules/cards/domain/Card.ts index b28448ba..0ca90374 100644 --- a/src/modules/cards/domain/Card.ts +++ b/src/modules/cards/domain/Card.ts @@ -241,7 +241,10 @@ export class Card extends AggregateRoot { return ok(undefined); } - public addToLibrary(userId: CuratorId): Result { + public addToLibrary( + userId: CuratorId, + addedAt?: Date, + ): Result { if ( this.props.libraryMemberships.find((link) => link.curatorId.equals(userId), @@ -261,15 +264,20 @@ export class Card extends AggregateRoot { ); } + const membershipAddedAt = addedAt ?? new Date(); this.props.libraryMemberships.push({ curatorId: userId, - addedAt: new Date(), + addedAt: membershipAddedAt, }); this.props.libraryCount = this.props.libraryMemberships.length; this.props.updatedAt = new Date(); // Raise domain event - const domainEvent = CardAddedToLibraryEvent.create(this.cardId, userId); + const domainEvent = CardAddedToLibraryEvent.create( + this.cardId, + userId, + membershipAddedAt, + ); if (domainEvent.isErr()) { return err(new CardValidationError(domainEvent.error.message)); } diff --git a/src/modules/cards/domain/CardFactory.ts b/src/modules/cards/domain/CardFactory.ts index 8cec18ea..c87140e6 100644 --- a/src/modules/cards/domain/CardFactory.ts +++ b/src/modules/cards/domain/CardFactory.ts @@ -28,6 +28,7 @@ export type CardCreationInput = IUrlCardInput | INoteCardInput; interface CreateCardProps { curatorId: string; cardInput: CardCreationInput; + createdAt?: Date; } export class CardFactory { @@ -123,6 +124,7 @@ export class CardFactory { url, parentCardId, viaCardId, + createdAt: props.createdAt, }); } catch (error) { return err( diff --git a/src/modules/cards/domain/Collection.ts b/src/modules/cards/domain/Collection.ts index f0aa8f33..f93def12 100644 --- a/src/modules/cards/domain/Collection.ts +++ b/src/modules/cards/domain/Collection.ts @@ -203,6 +203,7 @@ export class Collection extends AggregateRoot { cardId: CardId, userId: CuratorId, viaCardId?: CardId, + addedAt?: Date, ): Result { if (!this.canAddCard(userId)) { return err( @@ -220,10 +221,11 @@ export class Collection extends AggregateRoot { return ok(existingLink); // Return existing link } + const linkAddedAt = addedAt ?? new Date(); const newLink: CardLink = { cardId, addedBy: userId, - addedAt: new Date(), + addedAt: linkAddedAt, viaCardId, publishedRecordId: undefined, // Will be set when published }; @@ -238,6 +240,7 @@ export class Collection extends AggregateRoot { cardId, this.collectionId, userId, + linkAddedAt, ).unwrap(), ); diff --git a/src/modules/cards/domain/events/CardAddedToCollectionEvent.ts b/src/modules/cards/domain/events/CardAddedToCollectionEvent.ts index 91200ea6..993836f3 100644 --- a/src/modules/cards/domain/events/CardAddedToCollectionEvent.ts +++ b/src/modules/cards/domain/events/CardAddedToCollectionEvent.ts @@ -14,6 +14,7 @@ export class CardAddedToCollectionEvent implements IDomainEvent { public readonly cardId: CardId, public readonly collectionId: CollectionId, public readonly addedBy: CuratorId, + public readonly addedAt: Date, dateTimeOccurred?: Date, ) { this.dateTimeOccurred = dateTimeOccurred || new Date(); @@ -23,14 +24,18 @@ export class CardAddedToCollectionEvent implements IDomainEvent { cardId: CardId, collectionId: CollectionId, addedBy: CuratorId, + addedAt: Date, ): Result { - return ok(new CardAddedToCollectionEvent(cardId, collectionId, addedBy)); + return ok( + new CardAddedToCollectionEvent(cardId, collectionId, addedBy, addedAt), + ); } public static reconstruct( cardId: CardId, collectionId: CollectionId, addedBy: CuratorId, + addedAt: Date, dateTimeOccurred: Date, ): Result { return ok( @@ -38,6 +43,7 @@ export class CardAddedToCollectionEvent implements IDomainEvent { cardId, collectionId, addedBy, + addedAt, dateTimeOccurred, ), ); diff --git a/src/modules/cards/domain/events/CardAddedToLibraryEvent.ts b/src/modules/cards/domain/events/CardAddedToLibraryEvent.ts index 3e30695b..1215259c 100644 --- a/src/modules/cards/domain/events/CardAddedToLibraryEvent.ts +++ b/src/modules/cards/domain/events/CardAddedToLibraryEvent.ts @@ -12,6 +12,7 @@ export class CardAddedToLibraryEvent implements IDomainEvent { private constructor( public readonly cardId: CardId, public readonly curatorId: CuratorId, + public readonly addedAt: Date, dateTimeOccurred?: Date, ) { this.dateTimeOccurred = dateTimeOccurred || new Date(); @@ -20,16 +21,20 @@ export class CardAddedToLibraryEvent implements IDomainEvent { public static create( cardId: CardId, curatorId: CuratorId, + addedAt: Date, ): Result { - return ok(new CardAddedToLibraryEvent(cardId, curatorId)); + return ok(new CardAddedToLibraryEvent(cardId, curatorId, addedAt)); } public static reconstruct( cardId: CardId, curatorId: CuratorId, + addedAt: Date, dateTimeOccurred: Date, ): Result { - return ok(new CardAddedToLibraryEvent(cardId, curatorId, dateTimeOccurred)); + return ok( + new CardAddedToLibraryEvent(cardId, curatorId, addedAt, dateTimeOccurred), + ); } getAggregateId(): UniqueEntityID { diff --git a/src/modules/cards/domain/services/CardCollectionService.ts b/src/modules/cards/domain/services/CardCollectionService.ts index f7d2146c..b3979b62 100644 --- a/src/modules/cards/domain/services/CardCollectionService.ts +++ b/src/modules/cards/domain/services/CardCollectionService.ts @@ -18,6 +18,7 @@ import { AuthenticationError } from '../../../../shared/core/AuthenticationError export interface CardCollectionServiceOptions { skipPublishing?: boolean; publishedRecordIds?: Map; // collectionId -> publishedRecordId + timestamp?: Date; } export class CardCollectionValidationError extends Error { @@ -70,6 +71,7 @@ export class CardCollectionService implements DomainService { card.cardId, curatorId, viaCardId, + options?.timestamp, ); if (addCardResult.isErr()) { return err( diff --git a/src/modules/cards/domain/services/CardLibraryService.ts b/src/modules/cards/domain/services/CardLibraryService.ts index 59f05d87..8e5ffa52 100644 --- a/src/modules/cards/domain/services/CardLibraryService.ts +++ b/src/modules/cards/domain/services/CardLibraryService.ts @@ -15,6 +15,7 @@ export interface CardLibraryServiceOptions { publishedRecordId?: PublishedRecordId; skipUnpublishing?: boolean; skipCollectionUnpublishing?: boolean; + timestamp?: Date; } export class CardLibraryValidationError extends Error { @@ -268,7 +269,7 @@ export class CardLibraryService implements DomainService { // Card is already in library but not published, nothing to do return ok(card); } - const addToLibResult = card.addToLibrary(curatorId); + const addToLibResult = card.addToLibrary(curatorId, options?.timestamp); if (addToLibResult.isErr()) { return err( new CardLibraryValidationError( diff --git a/src/modules/cards/tests/integration/BullMQEventSystem.integration.test.ts b/src/modules/cards/tests/integration/BullMQEventSystem.integration.test.ts index 3b04dfeb..9c8c60b4 100644 --- a/src/modules/cards/tests/integration/BullMQEventSystem.integration.test.ts +++ b/src/modules/cards/tests/integration/BullMQEventSystem.integration.test.ts @@ -79,7 +79,11 @@ describe('BullMQ Event System Integration', () => { // Create test event const cardId = CardId.createFromString('test-card-123').unwrap(); const curatorId = CuratorId.create('did:plc:testuser123').unwrap(); - const event = CardAddedToLibraryEvent.create(cardId, curatorId).unwrap(); + const event = CardAddedToLibraryEvent.create( + cardId, + curatorId, + new Date(), + ).unwrap(); // Act - Publish event const publishResult = await publisher.publishEvents([event]); @@ -123,14 +127,17 @@ describe('BullMQ Event System Integration', () => { CardAddedToLibraryEvent.create( CardId.createFromString('card-1').unwrap(), CuratorId.create('did:plc:user1').unwrap(), + new Date(), ).unwrap(), CardAddedToLibraryEvent.create( CardId.createFromString('card-2').unwrap(), CuratorId.create('did:plc:user2').unwrap(), + new Date(), ).unwrap(), CardAddedToLibraryEvent.create( CardId.createFromString('card-3').unwrap(), CuratorId.create('did:plc:user3').unwrap(), + new Date(), ).unwrap(), ]; @@ -175,6 +182,7 @@ describe('BullMQ Event System Integration', () => { const event = CardAddedToLibraryEvent.create( CardId.createFromString('failing-card').unwrap(), CuratorId.create('did:plc:failuser').unwrap(), + new Date(), ).unwrap(); // Act - Publish event that will initially fail @@ -197,6 +205,7 @@ describe('BullMQ Event System Integration', () => { const event = CardAddedToLibraryEvent.create( CardId.createFromString('unhandled-card').unwrap(), CuratorId.create('did:plc:unhandleduser').unwrap(), + new Date(), ).unwrap(); // Act - Publish event @@ -236,6 +245,7 @@ describe('BullMQ Event System Integration', () => { const originalEvent = CardAddedToLibraryEvent.create( originalCardId, originalCuratorId, + new Date(), ).unwrap(); const originalTimestamp = originalEvent.dateTimeOccurred; @@ -264,6 +274,7 @@ describe('BullMQ Event System Integration', () => { const event = CardAddedToLibraryEvent.create( CardId.createFromString('queue-test-card').unwrap(), CuratorId.create('did:plc:queueuser').unwrap(), + new Date(), ).unwrap(); await publisher.publishEvents([event]); @@ -307,11 +318,13 @@ describe('BullMQ Event System Integration', () => { const libraryEvent = CardAddedToLibraryEvent.create( cardId, curatorId, + new Date(), ).unwrap(); const collectionEvent = CardAddedToCollectionEvent.create( cardId, CollectionId.createFromString('test-collection').unwrap(), curatorId, + new Date(), ).unwrap(); // Act - Process events with different saga instances @@ -356,16 +369,18 @@ describe('BullMQ Event System Integration', () => { CollectionId.createFromString('collection-2').unwrap(); const events = [ - CardAddedToLibraryEvent.create(cardId, curatorId).unwrap(), + CardAddedToLibraryEvent.create(cardId, curatorId, new Date()).unwrap(), CardAddedToCollectionEvent.create( cardId, collectionId1, curatorId, + new Date(), ).unwrap(), CardAddedToCollectionEvent.create( cardId, collectionId2, curatorId, + new Date(), ).unwrap(), ]; @@ -414,6 +429,7 @@ describe('BullMQ Event System Integration', () => { cardId, collectionId, curatorId, + new Date(), ).unwrap(); return { saga, event }; }); @@ -423,6 +439,7 @@ describe('BullMQ Event System Integration', () => { const libraryEvent = CardAddedToLibraryEvent.create( cardId, curatorId, + new Date(), ).unwrap(); // Act - Process all events concurrently @@ -470,7 +487,11 @@ describe('BullMQ Event System Integration', () => { // Create saga and event const saga = new CardCollectionSaga(mockUseCase, stateStore); - const event = CardAddedToLibraryEvent.create(cardId, curatorId).unwrap(); + const event = CardAddedToLibraryEvent.create( + cardId, + curatorId, + new Date(), + ).unwrap(); // Act - Try to process event (should initially be blocked by lock) // But should succeed after lock expires and retry mechanism kicks in @@ -502,11 +523,16 @@ describe('BullMQ Event System Integration', () => { const lockHoldingSaga = new CardCollectionSaga(mockUseCase, stateStore); const retryingSaga = new CardCollectionSaga(mockUseCase, stateStore); - const event1 = CardAddedToLibraryEvent.create(cardId, curatorId).unwrap(); + const event1 = CardAddedToLibraryEvent.create( + cardId, + curatorId, + new Date(), + ).unwrap(); const event2 = CardAddedToCollectionEvent.create( cardId, CollectionId.createFromString('retry-collection').unwrap(), curatorId, + new Date(), ).unwrap(); // Act - Start first saga (will acquire lock) @@ -573,6 +599,7 @@ describe('BullMQ Event System Integration', () => { const event = CardAddedToLibraryEvent.create( CardId.createFromString('multi-queue-card').unwrap(), CuratorId.create('did:plc:multiuser').unwrap(), + new Date(), ).unwrap(); await publisher.publishEvents([event]); diff --git a/src/modules/feeds/application/sagas/CardCollectionSaga.ts b/src/modules/feeds/application/sagas/CardCollectionSaga.ts index 5ee00db5..fd82fba5 100644 --- a/src/modules/feeds/application/sagas/CardCollectionSaga.ts +++ b/src/modules/feeds/application/sagas/CardCollectionSaga.ts @@ -12,7 +12,8 @@ interface PendingCardActivity { cardId: string; actorId: string; collectionIds: string[]; - timestamp: Date; + timestamp: Date; // Timestamp for the aggregation window + eventTimestamps: Date[]; // Track all addedAt timestamps from events hasLibraryEvent: boolean; hasCollectionEvents: boolean; } @@ -96,8 +97,11 @@ export class CardCollectionSaga { if (!data) return null; const parsed = JSON.parse(data); - // Convert timestamp string back to Date object + // Convert timestamp strings back to Date objects parsed.timestamp = new Date(parsed.timestamp); + parsed.eventTimestamps = parsed.eventTimestamps.map( + (ts: string) => new Date(ts), + ); return parsed; } @@ -162,6 +166,7 @@ export class CardCollectionSaga { actorId, collectionIds: [], timestamp: new Date(), + eventTimestamps: [], hasLibraryEvent: false, hasCollectionEvents: false, }; @@ -174,6 +179,9 @@ export class CardCollectionSaga { existing: PendingCardActivity, event: CardAddedToLibraryEvent | CardAddedToCollectionEvent, ): void { + // Track the addedAt timestamp from the event + existing.eventTimestamps.push(event.addedAt); + if ('curatorId' in event) { // CardAddedToLibraryEvent existing.hasLibraryEvent = true; @@ -201,12 +209,22 @@ export class CardCollectionSaga { const pending = await this.getPendingActivity(aggregationKey); if (!pending) return; + // Calculate the earliest timestamp from all events + // If no timestamps (shouldn't happen), use current time as fallback + const earliestTimestamp = + pending.eventTimestamps.length > 0 + ? pending.eventTimestamps.reduce((earliest, current) => + current < earliest ? current : earliest, + ) + : undefined; + const request: AddCardCollectedActivityDTO = { type: ActivityTypeEnum.CARD_COLLECTED, actorId: pending.actorId, cardId: pending.cardId, collectionIds: pending.collectionIds.length > 0 ? pending.collectionIds : undefined, + createdAt: earliestTimestamp, }; await this.addActivityToFeedUseCase.execute(request); diff --git a/src/modules/feeds/application/useCases/commands/AddActivityToFeedUseCase.ts b/src/modules/feeds/application/useCases/commands/AddActivityToFeedUseCase.ts index 3655cbab..e96cf97e 100644 --- a/src/modules/feeds/application/useCases/commands/AddActivityToFeedUseCase.ts +++ b/src/modules/feeds/application/useCases/commands/AddActivityToFeedUseCase.ts @@ -16,6 +16,7 @@ export interface AddCardCollectedActivityDTO { actorId: string; cardId: string; collectionIds?: string[]; + createdAt?: Date; // Timestamp from earliest event (for historical data) } export type AddActivityToFeedDTO = AddCardCollectedActivityDTO; @@ -126,12 +127,26 @@ export class AddActivityToFeedUseCase } } + // Determine createdAt timestamp for the activity + // For collections: use the provided timestamp (earliest addedAt from saga) + // For library-only: use the card's createdAt timestamp + let createdAt = request.createdAt; + if ( + !createdAt && + cardResult.value && + (!collectionIds || collectionIds.length === 0) + ) { + // Library-only scenario: use card's creation timestamp + createdAt = cardResult.value.createdAt; + } + const activityResult = await this.feedService.addCardCollectedActivity( actorId, cardId, collectionIds, urlType, source, + createdAt, ); if (activityResult.isErr()) { diff --git a/src/modules/feeds/domain/services/FeedService.ts b/src/modules/feeds/domain/services/FeedService.ts index d14dd0e4..13a0fa83 100644 --- a/src/modules/feeds/domain/services/FeedService.ts +++ b/src/modules/feeds/domain/services/FeedService.ts @@ -23,6 +23,7 @@ export class FeedService implements DomainService { collectionIds?: CollectionId[], urlType?: UrlType, source?: string, + createdAt?: Date, ): Promise> { try { // Check for recent duplicate activity (within 2 minutes) @@ -71,6 +72,7 @@ export class FeedService implements DomainService { collectionIds, urlType, source, + createdAt, ); if (activityResult.isErr()) { diff --git a/src/shared/infrastructure/events/EventMapper.ts b/src/shared/infrastructure/events/EventMapper.ts index 4b3e4a56..1a44c9f7 100644 --- a/src/shared/infrastructure/events/EventMapper.ts +++ b/src/shared/infrastructure/events/EventMapper.ts @@ -18,6 +18,7 @@ export interface SerializedCardAddedToLibraryEvent extends SerializedEvent { eventType: typeof EventNames.CARD_ADDED_TO_LIBRARY; cardId: string; curatorId: string; + addedAt: string; } export interface SerializedCardAddedToCollectionEvent extends SerializedEvent { @@ -25,6 +26,7 @@ export interface SerializedCardAddedToCollectionEvent extends SerializedEvent { cardId: string; collectionId: string; addedBy: string; + addedAt: string; } export interface SerializedCardRemovedFromLibraryEvent extends SerializedEvent { @@ -55,6 +57,7 @@ export class EventMapper { dateTimeOccurred: event.dateTimeOccurred.toISOString(), cardId: event.cardId.getValue().toString(), curatorId: event.curatorId.value, + addedAt: event.addedAt.toISOString(), }; } @@ -66,6 +69,7 @@ export class EventMapper { cardId: event.cardId.getValue().toString(), collectionId: event.collectionId.getValue().toString(), addedBy: event.addedBy.value, + addedAt: event.addedAt.toISOString(), }; } @@ -100,11 +104,13 @@ export class EventMapper { case EventNames.CARD_ADDED_TO_LIBRARY: { const cardId = CardId.createFromString(eventData.cardId).unwrap(); const curatorId = CuratorId.create(eventData.curatorId).unwrap(); + const addedAt = new Date(eventData.addedAt); const dateTimeOccurred = new Date(eventData.dateTimeOccurred); return CardAddedToLibraryEvent.reconstruct( cardId, curatorId, + addedAt, dateTimeOccurred, ).unwrap(); } @@ -114,12 +120,14 @@ export class EventMapper { eventData.collectionId, ).unwrap(); const addedBy = CuratorId.create(eventData.addedBy).unwrap(); + const addedAt = new Date(eventData.addedAt); const dateTimeOccurred = new Date(eventData.dateTimeOccurred); return CardAddedToCollectionEvent.reconstruct( cardId, collectionId, addedBy, + addedAt, dateTimeOccurred, ).unwrap(); }