diff --git a/src/modules/atproto/application/useCases/ProcessFirehoseEventUseCase.ts b/src/modules/atproto/application/useCases/ProcessFirehoseEventUseCase.ts index defb9db5..2cdfa30f 100644 --- a/src/modules/atproto/application/useCases/ProcessFirehoseEventUseCase.ts +++ b/src/modules/atproto/application/useCases/ProcessFirehoseEventUseCase.ts @@ -11,6 +11,7 @@ import { ProcessCollectionLinkFirehoseEventUseCase } from './ProcessCollectionLi import { ProcessMarginBookmarkFirehoseEventUseCase } from './ProcessMarginBookmarkFirehoseEventUseCase'; import { ProcessMarginCollectionFirehoseEventUseCase } from './ProcessMarginCollectionFirehoseEventUseCase'; import { ProcessMarginCollectionItemFirehoseEventUseCase } from './ProcessMarginCollectionItemFirehoseEventUseCase'; +import { ProcessMarginNoteFirehoseEventUseCase } from './ProcessMarginNoteFirehoseEventUseCase'; import { ProcessCollectionLinkRemovalFirehoseEventUseCase } from './ProcessCollectionLinkRemovalFirehoseEventUseCase'; import { ProcessFollowFirehoseEventUseCase } from './ProcessFollowFirehoseEventUseCase'; import { ProcessConnectionFirehoseEventUseCase } from './ProcessConnectionFirehoseEventUseCase'; @@ -25,6 +26,22 @@ import { Record as CollectionLinkRemovalRecord } from '../../infrastructure/lexi import { Record as FollowRecord } from '../../infrastructure/lexicon/types/network/cosmik/follow'; import { Record as ConnectionRecord } from '../../infrastructure/lexicon/types/network/cosmik/connection'; +// Margin Note Record type definition +interface MarginNoteTarget { + title?: string; + source: string; + sourceHash?: string; + [k: string]: unknown; +} + +interface MarginNoteRecord { + $type: 'at.margin.note'; + target: MarginNoteTarget; + createdAt: string; + motivation?: string; + [k: string]: unknown; +} + export interface ProcessFirehoseEventDTO { atUri: string; cid: string | null; @@ -50,6 +67,7 @@ export class ProcessFirehoseEventUseCase private processMarginBookmarkFirehoseEventUseCase: ProcessMarginBookmarkFirehoseEventUseCase, private processMarginCollectionFirehoseEventUseCase: ProcessMarginCollectionFirehoseEventUseCase, private processMarginCollectionItemFirehoseEventUseCase: ProcessMarginCollectionItemFirehoseEventUseCase, + private processMarginNoteFirehoseEventUseCase: ProcessMarginNoteFirehoseEventUseCase, private processCollectionLinkRemovalFirehoseEventUseCase: ProcessCollectionLinkRemovalFirehoseEventUseCase, private processFollowFirehoseEventUseCase: ProcessFollowFirehoseEventUseCase, private processConnectionFirehoseEventUseCase: ProcessConnectionFirehoseEventUseCase, @@ -218,6 +236,23 @@ export class ProcessFirehoseEventUseCase ...request, record: request.record as MarginCollectionItemRecord | undefined, }); + case collections.marginNote: + // Validate MarginNoteRecord structure + if ( + request.record && + (request.eventType === 'create' || request.eventType === 'update') + ) { + const noteRecord = request.record as MarginNoteRecord; + if (!noteRecord.target || !noteRecord.target.source) { + return err( + new ValidationError('Invalid Margin note record structure'), + ); + } + } + return this.processMarginNoteFirehoseEventUseCase.execute({ + ...request, + record: request.record as MarginNoteRecord | undefined, + }); case collections.follow: // Validate FollowRecord structure if ( diff --git a/src/modules/atproto/application/useCases/ProcessMarginNoteFirehoseEventUseCase.ts b/src/modules/atproto/application/useCases/ProcessMarginNoteFirehoseEventUseCase.ts new file mode 100644 index 00000000..09c8e207 --- /dev/null +++ b/src/modules/atproto/application/useCases/ProcessMarginNoteFirehoseEventUseCase.ts @@ -0,0 +1,237 @@ +import { Result, ok, err } from 'src/shared/core/Result'; +import { UseCase } from 'src/shared/core/UseCase'; +import { AppError } from 'src/shared/core/AppError'; +import { IAtUriResolutionService } from '../../../cards/domain/services/IAtUriResolutionService'; +import { PublishedRecordId } from '../../../cards/domain/value-objects/PublishedRecordId'; +import { ATUri } from '../../domain/ATUri'; +import { AddUrlToLibraryUseCase } from '../../../cards/application/useCases/commands/AddUrlToLibraryUseCase'; +import { RemoveCardFromLibraryUseCase } from '../../../cards/application/useCases/commands/RemoveCardFromLibraryUseCase'; + +// Margin Note Record type definition +interface MarginNoteTarget { + title?: string; + source: string; + sourceHash?: string; + [k: string]: unknown; +} + +interface MarginNoteRecord { + $type: 'at.margin.note'; + target: MarginNoteTarget; + createdAt: string; + motivation?: string; + [k: string]: unknown; +} + +export interface ProcessMarginNoteFirehoseEventDTO { + atUri: string; + cid: string | null; + eventType: 'create' | 'update' | 'delete'; + record?: MarginNoteRecord; +} + +const ENABLE_FIREHOSE_LOGGING = true; + +export class ProcessMarginNoteFirehoseEventUseCase + implements UseCase> +{ + constructor( + private atUriResolutionService: IAtUriResolutionService, + private addUrlToLibraryUseCase: AddUrlToLibraryUseCase, + private removeCardFromLibraryUseCase: RemoveCardFromLibraryUseCase, + ) {} + + async execute( + request: ProcessMarginNoteFirehoseEventDTO, + ): Promise> { + try { + if (ENABLE_FIREHOSE_LOGGING) { + console.log( + `[FirehoseWorker] Processing Margin note event: ${request.atUri} (${request.eventType})`, + ); + } + + switch (request.eventType) { + case 'create': + return await this.handleNoteCreate(request); + case 'update': + // Margin notes don't support updates for now + if (ENABLE_FIREHOSE_LOGGING) { + console.log( + `[FirehoseWorker] Ignoring Margin note update: ${request.atUri}`, + ); + } + return ok(undefined); + case 'delete': + return await this.handleNoteDelete(request); + } + + return ok(undefined); + } catch (error) { + return err(AppError.UnexpectedError.create(error)); + } + } + + private async handleNoteCreate( + request: ProcessMarginNoteFirehoseEventDTO, + ): Promise> { + if (!request.record || !request.cid) { + if (ENABLE_FIREHOSE_LOGGING) { + console.warn( + `[FirehoseWorker] Margin note create event missing record or cid, skipping: ${request.atUri}`, + ); + } + return ok(undefined); + } + + // Only process notes with 'bookmarking' motivation + if (request.record.motivation !== 'bookmarking') { + if (ENABLE_FIREHOSE_LOGGING) { + console.log( + `[FirehoseWorker] Skipping Margin note with motivation '${request.record.motivation}' (only processing 'bookmarking'): ${request.atUri}`, + ); + } + return ok(undefined); + } + + try { + // Parse AT URI to extract curator DID + const atUriResult = ATUri.create(request.atUri); + if (atUriResult.isErr()) { + if (ENABLE_FIREHOSE_LOGGING) { + console.warn( + `[FirehoseWorker] Invalid AT URI format: ${request.atUri} - ${atUriResult.error.message}`, + ); + } + return ok(undefined); + } + const atUri = atUriResult.value; + const curatorDid = atUri.did.value; + + // Extract URL from Margin note's 'target.source' field + const url = request.record.target.source; + if (!url) { + if (ENABLE_FIREHOSE_LOGGING) { + console.warn( + `[FirehoseWorker] Margin note missing target source URL - user: ${curatorDid}, uri: ${request.atUri}`, + ); + } + return ok(undefined); + } + + // Extract timestamp from AT Protocol record (Margin note has required createdAt) + const timestamp = new Date(request.record.createdAt); + + const publishedRecordId = PublishedRecordId.create({ + uri: request.atUri, + cid: request.cid, + }); + + const result = await this.addUrlToLibraryUseCase.execute({ + url: url, + curatorId: curatorDid, + publishedRecordId: publishedRecordId, + viaCardId: undefined, // Margin notes don't have 'via' references + timestamp: timestamp, + }); + + if (result.isErr()) { + if (ENABLE_FIREHOSE_LOGGING) { + console.warn( + `[FirehoseWorker] Failed to add Margin note to library - user: ${curatorDid}, uri: ${request.atUri}, error: ${result.error.message}`, + ); + } + return ok(undefined); + } + + if (ENABLE_FIREHOSE_LOGGING) { + console.log( + `[FirehoseWorker] Successfully created Margin note - user: ${curatorDid}, cardId: ${result.value.urlCardId}, uri: ${request.atUri}`, + ); + } + + return ok(undefined); + } catch (error) { + if (ENABLE_FIREHOSE_LOGGING) { + console.error( + `[FirehoseWorker] Error processing Margin note create event - uri: ${request.atUri}, error: ${error}`, + ); + } + return ok(undefined); // Don't fail the firehose processing + } + } + + private async handleNoteDelete( + request: ProcessMarginNoteFirehoseEventDTO, + ): Promise> { + try { + // Parse AT URI to extract curator DID + const atUriResult = ATUri.create(request.atUri); + if (atUriResult.isErr()) { + if (ENABLE_FIREHOSE_LOGGING) { + console.warn( + `[FirehoseWorker] Invalid AT URI format: ${request.atUri} - ${atUriResult.error.message}`, + ); + } + return ok(undefined); + } + const curatorDid = atUriResult.value.did.value; + + const cardIdResult = await this.atUriResolutionService.resolveCardId( + request.atUri, + ); + if (cardIdResult.isErr()) { + if (ENABLE_FIREHOSE_LOGGING) { + console.warn( + `[FirehoseWorker] Failed to resolve Margin note card ID - user: ${curatorDid}, uri: ${request.atUri}, error: ${cardIdResult.error.message}`, + ); + } + return ok(undefined); + } + + if (cardIdResult.value) { + if (ENABLE_FIREHOSE_LOGGING) { + console.log( + `[FirehoseWorker] Margin note deleted externally - user: ${curatorDid}, cardId: ${cardIdResult.value.getStringValue()}, uri: ${request.atUri}`, + ); + } + + // For delete events, we don't have a record, so no timestamp available + const publishedRecordId = PublishedRecordId.create({ + uri: request.atUri, + cid: request.cid || 'deleted', + }); + + const result = await this.removeCardFromLibraryUseCase.execute({ + cardId: cardIdResult.value.getStringValue(), + curatorId: curatorDid, + publishedRecordId: publishedRecordId, + }); + + if (result.isErr()) { + if (ENABLE_FIREHOSE_LOGGING) { + console.warn( + `[FirehoseWorker] Failed to remove Margin note from library - user: ${curatorDid}, cardId: ${cardIdResult.value.getStringValue()}, uri: ${request.atUri}, error: ${result.error.message}`, + ); + } + return ok(undefined); + } + + if (ENABLE_FIREHOSE_LOGGING) { + console.log( + `[FirehoseWorker] Successfully removed Margin note from library - user: ${curatorDid}, cardId: ${result.value.cardId}, uri: ${request.atUri}`, + ); + } + } + + return ok(undefined); + } catch (error) { + if (ENABLE_FIREHOSE_LOGGING) { + console.error( + `[FirehoseWorker] Error processing Margin note delete event - uri: ${request.atUri}, error: ${error}`, + ); + } + return ok(undefined); + } + } +} diff --git a/src/modules/atproto/infrastructure/services/AtProtoJetstreamService.ts b/src/modules/atproto/infrastructure/services/AtProtoJetstreamService.ts index 78f1a7b9..d6cb8e05 100644 --- a/src/modules/atproto/infrastructure/services/AtProtoJetstreamService.ts +++ b/src/modules/atproto/infrastructure/services/AtProtoJetstreamService.ts @@ -263,6 +263,7 @@ export class AtProtoJetstreamService implements IFirehoseService { collections.card, collections.collection, collections.collectionLink, + collections.marginNote, collections.marginBookmark, collections.marginCollection, collections.marginCollectionItem, diff --git a/src/modules/atproto/tests/application/ProcessFirehoseEventUseCase.test.ts b/src/modules/atproto/tests/application/ProcessFirehoseEventUseCase.test.ts index 4104d37a..182baec5 100644 --- a/src/modules/atproto/tests/application/ProcessFirehoseEventUseCase.test.ts +++ b/src/modules/atproto/tests/application/ProcessFirehoseEventUseCase.test.ts @@ -7,6 +7,7 @@ import { ProcessCollectionLinkFirehoseEventUseCase } from '../../application/use import { ProcessMarginBookmarkFirehoseEventUseCase } from '../../application/useCases/ProcessMarginBookmarkFirehoseEventUseCase'; import { ProcessMarginCollectionFirehoseEventUseCase } from '../../application/useCases/ProcessMarginCollectionFirehoseEventUseCase'; import { ProcessMarginCollectionItemFirehoseEventUseCase } from '../../application/useCases/ProcessMarginCollectionItemFirehoseEventUseCase'; +import { ProcessMarginNoteFirehoseEventUseCase } from '../../application/useCases/ProcessMarginNoteFirehoseEventUseCase'; import { ProcessCollectionLinkRemovalFirehoseEventUseCase } from '../../application/useCases/ProcessCollectionLinkRemovalFirehoseEventUseCase'; import { ProcessFollowFirehoseEventUseCase } from '../../application/useCases/ProcessFollowFirehoseEventUseCase'; import { ProcessConnectionFirehoseEventUseCase } from '../../application/useCases/ProcessConnectionFirehoseEventUseCase'; @@ -52,6 +53,7 @@ describe('ProcessFirehoseEventUseCase', () => { let processMarginBookmarkFirehoseEventUseCase: ProcessMarginBookmarkFirehoseEventUseCase; let processMarginCollectionFirehoseEventUseCase: ProcessMarginCollectionFirehoseEventUseCase; let processMarginCollectionItemFirehoseEventUseCase: ProcessMarginCollectionItemFirehoseEventUseCase; + let processMarginNoteFirehoseEventUseCase: ProcessMarginNoteFirehoseEventUseCase; let processCollectionLinkRemovalFirehoseEventUseCase: ProcessCollectionLinkRemovalFirehoseEventUseCase; let processFollowFirehoseEventUseCase: ProcessFollowFirehoseEventUseCase; let processConnectionFirehoseEventUseCase: ProcessConnectionFirehoseEventUseCase; @@ -201,6 +203,12 @@ describe('ProcessFirehoseEventUseCase', () => { atUriResolutionService, updateUrlCardAssociationsUseCase, ); + processMarginNoteFirehoseEventUseCase = + new ProcessMarginNoteFirehoseEventUseCase( + atUriResolutionService, + addUrlToLibraryUseCase, + removeCardFromLibraryUseCase, + ); processCollectionLinkRemovalFirehoseEventUseCase = new ProcessCollectionLinkRemovalFirehoseEventUseCase( atUriResolutionService, @@ -272,6 +280,7 @@ describe('ProcessFirehoseEventUseCase', () => { processMarginBookmarkFirehoseEventUseCase, processMarginCollectionFirehoseEventUseCase, processMarginCollectionItemFirehoseEventUseCase, + processMarginNoteFirehoseEventUseCase, processCollectionLinkRemovalFirehoseEventUseCase, processFollowFirehoseEventUseCase, processConnectionFirehoseEventUseCase, diff --git a/src/shared/constants/atproto.ts b/src/shared/constants/atproto.ts index 463ddea7..b8bfef11 100644 --- a/src/shared/constants/atproto.ts +++ b/src/shared/constants/atproto.ts @@ -8,6 +8,7 @@ export const ATPROTO_NSID = { BOOKMARK: 'at.margin.bookmark', COLLECTION: 'at.margin.collection', COLLECTION_ITEM: 'at.margin.collectionItem', + NOTE: 'at.margin.note', }, COSMIK: { NAMESPACE: 'network.cosmik', diff --git a/src/shared/infrastructure/config/EnvironmentConfigService.ts b/src/shared/infrastructure/config/EnvironmentConfigService.ts index 899f8ac5..107049a8 100644 --- a/src/shared/infrastructure/config/EnvironmentConfigService.ts +++ b/src/shared/infrastructure/config/EnvironmentConfigService.ts @@ -33,6 +33,7 @@ export interface EnvironmentConfig { collection: string; collectionLink: string; marginBookmark: string; + marginNote: string; marginCollection: string; marginCollectionItem: string; collectionLinkRemoval: string; @@ -149,6 +150,7 @@ export class EnvironmentConfigService { marginBookmark: ATPROTO_NSID.MARGIN.BOOKMARK, marginCollection: ATPROTO_NSID.MARGIN.COLLECTION, marginCollectionItem: ATPROTO_NSID.MARGIN.COLLECTION_ITEM, + marginNote: ATPROTO_NSID.MARGIN.NOTE, }, serviceAccount: { identifier: process.env.BSKY_SERVICE_ACCOUNT_IDENTIFIER || '', diff --git a/src/shared/infrastructure/http/factories/UseCaseFactory.ts b/src/shared/infrastructure/http/factories/UseCaseFactory.ts index d763a77e..7d8c5f76 100644 --- a/src/shared/infrastructure/http/factories/UseCaseFactory.ts +++ b/src/shared/infrastructure/http/factories/UseCaseFactory.ts @@ -48,6 +48,7 @@ import { ProcessCollectionLinkFirehoseEventUseCase } from '../../../../modules/a import { ProcessMarginBookmarkFirehoseEventUseCase } from '../../../../modules/atproto/application/useCases/ProcessMarginBookmarkFirehoseEventUseCase'; import { ProcessMarginCollectionFirehoseEventUseCase } from '../../../../modules/atproto/application/useCases/ProcessMarginCollectionFirehoseEventUseCase'; import { ProcessMarginCollectionItemFirehoseEventUseCase } from '../../../../modules/atproto/application/useCases/ProcessMarginCollectionItemFirehoseEventUseCase'; +import { ProcessMarginNoteFirehoseEventUseCase } from '../../../../modules/atproto/application/useCases/ProcessMarginNoteFirehoseEventUseCase'; import { ProcessCollectionLinkRemovalFirehoseEventUseCase } from '../../../../modules/atproto/application/useCases/ProcessCollectionLinkRemovalFirehoseEventUseCase'; import { ProcessFollowFirehoseEventUseCase } from '../../../../modules/atproto/application/useCases/ProcessFollowFirehoseEventUseCase'; import { ProcessConnectionFirehoseEventUseCase } from '../../../../modules/atproto/application/useCases/ProcessConnectionFirehoseEventUseCase'; @@ -94,6 +95,7 @@ export interface WorkerUseCases { processMarginBookmarkFirehoseEventUseCase: ProcessMarginBookmarkFirehoseEventUseCase; processMarginCollectionFirehoseEventUseCase: ProcessMarginCollectionFirehoseEventUseCase; processMarginCollectionItemFirehoseEventUseCase: ProcessMarginCollectionItemFirehoseEventUseCase; + processMarginNoteFirehoseEventUseCase: ProcessMarginNoteFirehoseEventUseCase; processCollectionLinkRemovalFirehoseEventUseCase: ProcessCollectionLinkRemovalFirehoseEventUseCase; processFollowFirehoseEventUseCase: ProcessFollowFirehoseEventUseCase; processConnectionFirehoseEventUseCase: ProcessConnectionFirehoseEventUseCase; @@ -670,6 +672,13 @@ export class UseCaseFactory { updateUrlCardAssociationsUseCase, ); + const processMarginNoteFirehoseEventUseCase = + new ProcessMarginNoteFirehoseEventUseCase( + repositories.atUriResolutionService, + addUrlToLibraryUseCase, + removeCardFromLibraryUseCase, + ); + const processFollowFirehoseEventUseCase = new ProcessFollowFirehoseEventUseCase( repositories.atUriResolutionService, @@ -724,6 +733,7 @@ export class UseCaseFactory { processMarginBookmarkFirehoseEventUseCase, processMarginCollectionFirehoseEventUseCase, processMarginCollectionItemFirehoseEventUseCase, + processMarginNoteFirehoseEventUseCase, processFollowFirehoseEventUseCase, processConnectionFirehoseEventUseCase, // Level 3 diff --git a/src/workers/firehose-worker.ts b/src/workers/firehose-worker.ts index bc1f1615..9c2b8c71 100644 --- a/src/workers/firehose-worker.ts +++ b/src/workers/firehose-worker.ts @@ -41,6 +41,7 @@ async function main() { useCases.processMarginBookmarkFirehoseEventUseCase, useCases.processMarginCollectionFirehoseEventUseCase, useCases.processMarginCollectionItemFirehoseEventUseCase, + useCases.processMarginNoteFirehoseEventUseCase, useCases.processCollectionLinkRemovalFirehoseEventUseCase, useCases.processFollowFirehoseEventUseCase, useCases.processConnectionFirehoseEventUseCase,