From e939e6b07e1bbfa79e5d39347ee6a9c96ee3b16e Mon Sep 17 00:00:00 2001 From: Wesley Finck Date: Thu, 22 Jan 2026 21:52:47 -0800 Subject: [PATCH] update to firehose handlers as well as new handler for collectionLinkRemoval --- .../ProcessCollectionFirehoseEventUseCase.ts | 3 + ...llectionLinkRemovalFirehoseEventUseCase.ts | 151 ++++++++++++++++++ .../useCases/ProcessFirehoseEventUseCase.ts | 24 +++ .../services/AtProtoJetstreamService.ts | 1 + .../DrizzleFirehoseEventDuplicationService.ts | 5 + .../ProcessFirehoseEventUseCase.test.ts | 9 ++ .../UpdateUrlCardAssociationsUseCase.test.ts | 7 +- .../http/factories/UseCaseFactory.ts | 12 ++ src/workers/firehose-worker.ts | 1 + 9 files changed, 210 insertions(+), 3 deletions(-) create mode 100644 src/modules/atproto/application/useCases/ProcessCollectionLinkRemovalFirehoseEventUseCase.ts diff --git a/src/modules/atproto/application/useCases/ProcessCollectionFirehoseEventUseCase.ts b/src/modules/atproto/application/useCases/ProcessCollectionFirehoseEventUseCase.ts index 5cdd6de1..6cb7ced3 100644 --- a/src/modules/atproto/application/useCases/ProcessCollectionFirehoseEventUseCase.ts +++ b/src/modules/atproto/application/useCases/ProcessCollectionFirehoseEventUseCase.ts @@ -8,6 +8,7 @@ import { Record as CollectionRecord } from '../../infrastructure/lexicon/types/n import { CreateCollectionUseCase } from '../../../cards/application/useCases/commands/CreateCollectionUseCase'; import { UpdateCollectionUseCase } from '../../../cards/application/useCases/commands/UpdateCollectionUseCase'; import { DeleteCollectionUseCase } from '../../../cards/application/useCases/commands/DeleteCollectionUseCase'; +import { CollectionAccessType } from '../../../cards/domain/Collection'; export interface ProcessCollectionFirehoseEventDTO { atUri: string; cid: string | null; @@ -84,6 +85,7 @@ export class ProcessCollectionFirehoseEventUseCase const result = await this.createCollectionUseCase.execute({ name: request.record.name, description: request.record.description, + accessType: request.record.accessType as CollectionAccessType | undefined, curatorId: authorDid, publishedRecordId: publishedRecordId, }); @@ -168,6 +170,7 @@ export class ProcessCollectionFirehoseEventUseCase collectionId: collectionIdResult.value.getStringValue(), name: request.record.name, description: request.record.description, + accessType: request.record.accessType as CollectionAccessType | undefined, curatorId: authorDid, publishedRecordId: publishedRecordId, }); diff --git a/src/modules/atproto/application/useCases/ProcessCollectionLinkRemovalFirehoseEventUseCase.ts b/src/modules/atproto/application/useCases/ProcessCollectionLinkRemovalFirehoseEventUseCase.ts new file mode 100644 index 00000000..8db549fb --- /dev/null +++ b/src/modules/atproto/application/useCases/ProcessCollectionLinkRemovalFirehoseEventUseCase.ts @@ -0,0 +1,151 @@ +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 { ATUri } from '../../domain/ATUri'; +import { Record as CollectionLinkRemovalRecord } from '../../infrastructure/lexicon/types/network/cosmik/collectionLinkRemoval'; +import { + UpdateUrlCardAssociationsUseCase, + OperationContext, +} from '../../../cards/application/useCases/commands/UpdateUrlCardAssociationsUseCase'; + +export interface ProcessCollectionLinkRemovalFirehoseEventDTO { + atUri: string; + cid: string | null; + eventType: 'create' | 'update' | 'delete'; + record?: CollectionLinkRemovalRecord; +} + +const ENABLE_FIREHOSE_LOGGING = true; + +export class ProcessCollectionLinkRemovalFirehoseEventUseCase + implements + UseCase> +{ + constructor( + private atUriResolutionService: IAtUriResolutionService, + private updateUrlCardAssociationsUseCase: UpdateUrlCardAssociationsUseCase, + ) {} + + async execute( + request: ProcessCollectionLinkRemovalFirehoseEventDTO, + ): Promise> { + try { + if (ENABLE_FIREHOSE_LOGGING) { + console.log( + `[FirehoseWorker] Processing collection link removal event: ${request.atUri} (${request.eventType})`, + ); + } + + switch (request.eventType) { + case 'create': + return await this.handleCollectionLinkRemovalCreate(request); + case 'delete': + // For now, we don't handle deletion events for collectionLinkRemoval + if (ENABLE_FIREHOSE_LOGGING) { + console.log( + `[FirehoseWorker] Collection link removal delete event (not handled): ${request.atUri}`, + ); + } + break; + case 'update': + // Collection link removals don't typically have update operations + if (ENABLE_FIREHOSE_LOGGING) { + console.log( + `[FirehoseWorker] Collection link removal update event (unusual): ${request.atUri}`, + ); + } + break; + } + + return ok(undefined); + } catch (error) { + return err(AppError.UnexpectedError.create(error)); + } + } + + private async handleCollectionLinkRemovalCreate( + request: ProcessCollectionLinkRemovalFirehoseEventDTO, + ): Promise> { + if (!request.record || !request.cid) { + if (ENABLE_FIREHOSE_LOGGING) { + console.warn( + `[FirehoseWorker] Collection link removal create event missing record or cid, skipping: ${request.atUri}`, + ); + } + return ok(undefined); + } + + try { + // Parse AT URI to extract curator DID (the person who published the removal) + 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 removerDid = atUriResult.value.did.value; + + // Resolve the collection link that's being removed + const collectionLinkUri = request.record.removedLink.uri; + const linkInfoResult = + await this.atUriResolutionService.resolveCollectionLinkId( + collectionLinkUri, + ); + + if (linkInfoResult.isErr()) { + if (ENABLE_FIREHOSE_LOGGING) { + console.warn( + `[FirehoseWorker] Failed to resolve collection link - remover: ${removerDid}, collectionLinkUri: ${collectionLinkUri}, removalUri: ${request.atUri}`, + ); + } + return ok(undefined); + } + + if (!linkInfoResult.value) { + if (ENABLE_FIREHOSE_LOGGING) { + console.log( + `[FirehoseWorker] Collection link not found in our system (may have already been removed) - remover: ${removerDid}, collectionLinkUri: ${collectionLinkUri}, removalUri: ${request.atUri}`, + ); + } + return ok(undefined); + } + + const { cardId, collectionId } = linkInfoResult.value; + + // Remove the card from the collection using the context flag to skip publishing + const result = await this.updateUrlCardAssociationsUseCase.execute({ + cardId: cardId.getStringValue(), + curatorId: removerDid, + removeFromCollections: [collectionId.getStringValue()], + context: OperationContext.FIREHOSE_EVENT, // This tells the use case to skip publishing + }); + + if (result.isErr()) { + if (ENABLE_FIREHOSE_LOGGING) { + console.warn( + `[FirehoseWorker] Failed to remove card from collection - remover: ${removerDid}, cardId: ${cardId.getStringValue()}, collectionId: ${collectionId.getStringValue()}, removalUri: ${request.atUri}, error: ${result.error.message}`, + ); + } + return ok(undefined); + } + + if (ENABLE_FIREHOSE_LOGGING) { + console.log( + `[FirehoseWorker] Successfully removed card from collection via removal record - remover: ${removerDid}, cardId: ${cardId.getStringValue()}, collectionId: ${collectionId.getStringValue()}, removalUri: ${request.atUri}`, + ); + } + return ok(undefined); + } catch (error) { + if (ENABLE_FIREHOSE_LOGGING) { + console.error( + `[FirehoseWorker] Error processing collection link removal create event - uri: ${request.atUri}, error: ${error}`, + ); + } + return ok(undefined); // Don't fail the firehose processing + } + } +} diff --git a/src/modules/atproto/application/useCases/ProcessFirehoseEventUseCase.ts b/src/modules/atproto/application/useCases/ProcessFirehoseEventUseCase.ts index 7b68bf64..6215255b 100644 --- a/src/modules/atproto/application/useCases/ProcessFirehoseEventUseCase.ts +++ b/src/modules/atproto/application/useCases/ProcessFirehoseEventUseCase.ts @@ -8,10 +8,12 @@ import { EnvironmentConfigService } from 'src/shared/infrastructure/config/Envir import { ProcessCardFirehoseEventUseCase } from './ProcessCardFirehoseEventUseCase'; import { ProcessCollectionFirehoseEventUseCase } from './ProcessCollectionFirehoseEventUseCase'; import { ProcessCollectionLinkFirehoseEventUseCase } from './ProcessCollectionLinkFirehoseEventUseCase'; +import { ProcessCollectionLinkRemovalFirehoseEventUseCase } from './ProcessCollectionLinkRemovalFirehoseEventUseCase'; import type { RepoRecord } from '@atproto/lexicon'; import { Record as CardRecord } from '../../infrastructure/lexicon/types/network/cosmik/card'; import { Record as CollectionRecord } from '../../infrastructure/lexicon/types/network/cosmik/collection'; import { Record as CollectionLinkRecord } from '../../infrastructure/lexicon/types/network/cosmik/collectionLink'; +import { Record as CollectionLinkRemovalRecord } from '../../infrastructure/lexicon/types/network/cosmik/collectionLinkRemoval'; export interface ProcessFirehoseEventDTO { atUri: string; @@ -35,6 +37,7 @@ export class ProcessFirehoseEventUseCase private processCardFirehoseEventUseCase: ProcessCardFirehoseEventUseCase, private processCollectionFirehoseEventUseCase: ProcessCollectionFirehoseEventUseCase, private processCollectionLinkFirehoseEventUseCase: ProcessCollectionLinkFirehoseEventUseCase, + private processCollectionLinkRemovalFirehoseEventUseCase: ProcessCollectionLinkRemovalFirehoseEventUseCase, ) {} async execute(request: ProcessFirehoseEventDTO): Promise> { @@ -121,6 +124,27 @@ export class ProcessFirehoseEventUseCase ...request, record: request.record as CollectionLinkRecord | undefined, }); + case collections.collectionLinkRemoval: + // Validate CollectionLinkRemovalRecord structure + if ( + request.record && + (request.eventType === 'create' || request.eventType === 'update') + ) { + const removalRecord = request.record as CollectionLinkRemovalRecord; + if ( + !removalRecord.collection || + !removalRecord.removedLink || + !removalRecord.removedAt + ) { + return err( + new ValidationError('Invalid collection link removal record structure'), + ); + } + } + return this.processCollectionLinkRemovalFirehoseEventUseCase.execute({ + ...request, + record: request.record as CollectionLinkRemovalRecord | undefined, + }); default: return err( new ValidationError(`Unknown collection type: ${collection}`), diff --git a/src/modules/atproto/infrastructure/services/AtProtoJetstreamService.ts b/src/modules/atproto/infrastructure/services/AtProtoJetstreamService.ts index f7e42255..9a1b4148 100644 --- a/src/modules/atproto/infrastructure/services/AtProtoJetstreamService.ts +++ b/src/modules/atproto/infrastructure/services/AtProtoJetstreamService.ts @@ -262,6 +262,7 @@ export class AtProtoJetstreamService implements IFirehoseService { collections.card, collections.collection, collections.collectionLink, + collections.collectionLinkRemoval, ]; } diff --git a/src/modules/atproto/infrastructure/services/DrizzleFirehoseEventDuplicationService.ts b/src/modules/atproto/infrastructure/services/DrizzleFirehoseEventDuplicationService.ts index a806740e..4b246aeb 100644 --- a/src/modules/atproto/infrastructure/services/DrizzleFirehoseEventDuplicationService.ts +++ b/src/modules/atproto/infrastructure/services/DrizzleFirehoseEventDuplicationService.ts @@ -93,6 +93,11 @@ export class DrizzleFirehoseEventDuplicationService } return ok(linkInfoResult.value === null); } + case collections.collectionLinkRemoval: { + // For removal records, we track them in published_records table + // If a record exists, it hasn't been deleted + return ok(records.length === 0); + } default: return err(new Error(`Unknown collection type: ${collection}`)); } diff --git a/src/modules/atproto/tests/application/ProcessFirehoseEventUseCase.test.ts b/src/modules/atproto/tests/application/ProcessFirehoseEventUseCase.test.ts index 5a757276..b6f1c583 100644 --- a/src/modules/atproto/tests/application/ProcessFirehoseEventUseCase.test.ts +++ b/src/modules/atproto/tests/application/ProcessFirehoseEventUseCase.test.ts @@ -3,6 +3,7 @@ import { InMemoryFirehoseEventDuplicationService } from '../utils/InMemoryFireho import { ProcessCardFirehoseEventUseCase } from '../../application/useCases/ProcessCardFirehoseEventUseCase'; import { ProcessCollectionFirehoseEventUseCase } from '../../application/useCases/ProcessCollectionFirehoseEventUseCase'; import { ProcessCollectionLinkFirehoseEventUseCase } from '../../application/useCases/ProcessCollectionLinkFirehoseEventUseCase'; +import { ProcessCollectionLinkRemovalFirehoseEventUseCase } from '../../application/useCases/ProcessCollectionLinkRemovalFirehoseEventUseCase'; import { EnvironmentConfigService } from '../../../../shared/infrastructure/config/EnvironmentConfigService'; import { InMemoryAtUriResolutionService } from '../../../cards/tests/utils/InMemoryAtUriResolutionService'; import { AddUrlToLibraryUseCase } from '../../../cards/application/useCases/commands/AddUrlToLibraryUseCase'; @@ -30,6 +31,7 @@ describe('ProcessFirehoseEventUseCase', () => { let processCardFirehoseEventUseCase: ProcessCardFirehoseEventUseCase; let processCollectionFirehoseEventUseCase: ProcessCollectionFirehoseEventUseCase; let processCollectionLinkFirehoseEventUseCase: ProcessCollectionLinkFirehoseEventUseCase; + let processCollectionLinkRemovalFirehoseEventUseCase: ProcessCollectionLinkRemovalFirehoseEventUseCase; // Dependencies for real use cases let atUriResolutionService: InMemoryAtUriResolutionService; @@ -138,12 +140,19 @@ describe('ProcessFirehoseEventUseCase', () => { updateUrlCardAssociationsUseCase, ); + processCollectionLinkRemovalFirehoseEventUseCase = + new ProcessCollectionLinkRemovalFirehoseEventUseCase( + atUriResolutionService, + updateUrlCardAssociationsUseCase, + ); + useCase = new ProcessFirehoseEventUseCase( duplicationService, configService, processCardFirehoseEventUseCase, processCollectionFirehoseEventUseCase, processCollectionLinkFirehoseEventUseCase, + processCollectionLinkRemovalFirehoseEventUseCase, ); }); diff --git a/src/modules/cards/tests/application/UpdateUrlCardAssociationsUseCase.test.ts b/src/modules/cards/tests/application/UpdateUrlCardAssociationsUseCase.test.ts index bc9dfae5..0913e955 100644 --- a/src/modules/cards/tests/application/UpdateUrlCardAssociationsUseCase.test.ts +++ b/src/modules/cards/tests/application/UpdateUrlCardAssociationsUseCase.test.ts @@ -555,18 +555,19 @@ describe('UpdateUrlCardAssociationsUseCase', () => { expect(addResult.isOk()).toBe(true); const urlCardId = addResult.unwrap().urlCardId; - // Manually add card to collection (as collection owner) to simulate existing state + // Add card to collection as the current user (curatorId) + // In an OPEN collection, users can only remove cards they added themselves const cardResult = await cardRepository.findById( CardId.createFromString(urlCardId).unwrap(), ); expect(cardResult.isOk()).toBe(true); const card = cardResult.unwrap()!; - const addCardResult = collection.addCard(card.cardId, otherCuratorId); + const addCardResult = collection.addCard(card.cardId, curatorId); expect(addCardResult.isOk()).toBe(true); await collectionRepository.save(collection); - // Remove from collection + // Remove from collection - should succeed because curatorId added the card const updateRequest = { cardId: urlCardId, curatorId: curatorId.value, diff --git a/src/shared/infrastructure/http/factories/UseCaseFactory.ts b/src/shared/infrastructure/http/factories/UseCaseFactory.ts index 306c7122..97ea868f 100644 --- a/src/shared/infrastructure/http/factories/UseCaseFactory.ts +++ b/src/shared/infrastructure/http/factories/UseCaseFactory.ts @@ -41,6 +41,7 @@ import { SearchLeafletDocsForUrlUseCase } from '../../../../modules/search/appli import { ProcessCardFirehoseEventUseCase } from '../../../../modules/atproto/application/useCases/ProcessCardFirehoseEventUseCase'; import { ProcessCollectionFirehoseEventUseCase } from '../../../../modules/atproto/application/useCases/ProcessCollectionFirehoseEventUseCase'; import { ProcessCollectionLinkFirehoseEventUseCase } from '../../../../modules/atproto/application/useCases/ProcessCollectionLinkFirehoseEventUseCase'; +import { ProcessCollectionLinkRemovalFirehoseEventUseCase } from '../../../../modules/atproto/application/useCases/ProcessCollectionLinkRemovalFirehoseEventUseCase'; import { GetMyNotificationsUseCase } from '../../../../modules/notifications/application/useCases/queries/GetMyNotificationsUseCase'; import { GetUnreadNotificationCountUseCase } from '../../../../modules/notifications/application/useCases/queries/GetUnreadNotificationCountUseCase'; import { MarkNotificationsAsReadUseCase } from '../../../../modules/notifications/application/useCases/commands/MarkNotificationsAsReadUseCase'; @@ -61,6 +62,7 @@ export interface WorkerUseCases { processCardFirehoseEventUseCase: ProcessCardFirehoseEventUseCase; processCollectionFirehoseEventUseCase: ProcessCollectionFirehoseEventUseCase; processCollectionLinkFirehoseEventUseCase: ProcessCollectionLinkFirehoseEventUseCase; + processCollectionLinkRemovalFirehoseEventUseCase: ProcessCollectionLinkRemovalFirehoseEventUseCase; } export interface UseCases { @@ -414,6 +416,16 @@ export class UseCaseFactory { services.eventPublisher, ), ), + processCollectionLinkRemovalFirehoseEventUseCase: + new ProcessCollectionLinkRemovalFirehoseEventUseCase( + repositories.atUriResolutionService, + new UpdateUrlCardAssociationsUseCase( + repositories.cardRepository, + services.cardLibraryService, + services.cardCollectionService, + services.eventPublisher, + ), + ), }; } } diff --git a/src/workers/firehose-worker.ts b/src/workers/firehose-worker.ts index 1c94c503..82aa8448 100644 --- a/src/workers/firehose-worker.ts +++ b/src/workers/firehose-worker.ts @@ -35,6 +35,7 @@ async function main() { useCases.processCardFirehoseEventUseCase, useCases.processCollectionFirehoseEventUseCase, useCases.processCollectionLinkFirehoseEventUseCase, + useCases.processCollectionLinkRemovalFirehoseEventUseCase, ); const firehoseEventHandler = new FirehoseEventHandler( -- 2.51.2