This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

semble / src / modules / atproto / application / useCases / ProcessCollectionFirehoseEventUseCase.ts
9.8 kB 278 lines
1import { Result, ok, err } from 'src/shared/core/Result'; 2import { UseCase } from 'src/shared/core/UseCase'; 3import { AppError } from 'src/shared/core/AppError'; 4import { IAtUriResolutionService } from '../../../cards/domain/services/IAtUriResolutionService'; 5import { PublishedRecordId } from '../../../cards/domain/value-objects/PublishedRecordId'; 6import { ATUri } from '../../domain/ATUri'; 7import { Record as CollectionRecord } from '../../infrastructure/lexicon/types/network/cosmik/collection'; 8import { CreateCollectionUseCase } from '../../../cards/application/useCases/commands/CreateCollectionUseCase'; 9import { UpdateCollectionUseCase } from '../../../cards/application/useCases/commands/UpdateCollectionUseCase'; 10import { DeleteCollectionUseCase } from '../../../cards/application/useCases/commands/DeleteCollectionUseCase'; 11import { CollectionAccessType } from '../../../cards/domain/Collection'; 12export interface ProcessCollectionFirehoseEventDTO { 13 atUri: string; 14 cid: string | null; 15 eventType: 'create' | 'update' | 'delete'; 16 record?: CollectionRecord; 17} 18const ENABLE_FIREHOSE_LOGGING = true; 19 20export class ProcessCollectionFirehoseEventUseCase 21 implements UseCase<ProcessCollectionFirehoseEventDTO, Result<void>> 22{ 23 constructor( 24 private atUriResolutionService: IAtUriResolutionService, 25 private createCollectionUseCase: CreateCollectionUseCase, 26 private updateCollectionUseCase: UpdateCollectionUseCase, 27 private deleteCollectionUseCase: DeleteCollectionUseCase, 28 ) {} 29 30 async execute( 31 request: ProcessCollectionFirehoseEventDTO, 32 ): Promise<Result<void>> { 33 try { 34 if (ENABLE_FIREHOSE_LOGGING) { 35 console.log( 36 `[FirehoseWorker] Processing collection event: ${request.atUri} (${request.eventType})`, 37 ); 38 } 39 40 switch (request.eventType) { 41 case 'create': 42 return await this.handleCollectionCreate(request); 43 case 'update': 44 return await this.handleCollectionUpdate(request); 45 case 'delete': 46 return await this.handleCollectionDelete(request); 47 } 48 49 return ok(undefined); 50 } catch (error) { 51 return err(AppError.UnexpectedError.create(error)); 52 } 53 } 54 55 private async handleCollectionCreate( 56 request: ProcessCollectionFirehoseEventDTO, 57 ): Promise<Result<void>> { 58 if (!request.record || !request.cid) { 59 if (ENABLE_FIREHOSE_LOGGING) { 60 console.warn( 61 `[FirehoseWorker] Collection create event missing record or cid, skipping: ${request.atUri}`, 62 ); 63 } 64 return ok(undefined); 65 } 66 67 try { 68 // Parse AT URI to extract author DID 69 const atUriResult = ATUri.create(request.atUri); 70 if (atUriResult.isErr()) { 71 if (ENABLE_FIREHOSE_LOGGING) { 72 console.warn( 73 `[FirehoseWorker] Invalid AT URI format: ${request.atUri} - ${atUriResult.error.message}`, 74 ); 75 } 76 return ok(undefined); 77 } 78 const authorDid = atUriResult.value.did.value; 79 80 const publishedRecordId = PublishedRecordId.create({ 81 uri: request.atUri, 82 cid: request.cid, 83 }); 84 85 const result = await this.createCollectionUseCase.execute({ 86 name: request.record.name, 87 description: request.record.description, 88 accessType: request.record.accessType as 89 | CollectionAccessType 90 | undefined, 91 curatorId: authorDid, 92 publishedRecordId: publishedRecordId, 93 }); 94 95 if (result.isErr()) { 96 if (ENABLE_FIREHOSE_LOGGING) { 97 console.warn( 98 `[FirehoseWorker] Failed to create collection - user: ${authorDid}, uri: ${request.atUri}, error: ${result.error.message}`, 99 ); 100 } 101 return ok(undefined); 102 } 103 104 if (ENABLE_FIREHOSE_LOGGING) { 105 console.log( 106 `[FirehoseWorker] Successfully created collection - user: ${authorDid}, collectionId: ${result.value.collectionId}, uri: ${request.atUri}`, 107 ); 108 } 109 return ok(undefined); 110 } catch (error) { 111 if (ENABLE_FIREHOSE_LOGGING) { 112 console.error( 113 `[FirehoseWorker] Error processing collection create event - uri: ${request.atUri}, error: ${error}`, 114 ); 115 } 116 return ok(undefined); // Don't fail the firehose processing 117 } 118 } 119 120 private async handleCollectionUpdate( 121 request: ProcessCollectionFirehoseEventDTO, 122 ): Promise<Result<void>> { 123 if (!request.record || !request.cid) { 124 if (ENABLE_FIREHOSE_LOGGING) { 125 console.warn( 126 `[FirehoseWorker] Collection update event missing record or cid, skipping: ${request.atUri}`, 127 ); 128 } 129 return ok(undefined); 130 } 131 132 try { 133 // Parse AT URI to extract author DID 134 const atUriResult = ATUri.create(request.atUri); 135 if (atUriResult.isErr()) { 136 if (ENABLE_FIREHOSE_LOGGING) { 137 console.warn( 138 `[FirehoseWorker] Invalid AT URI format: ${request.atUri} - ${atUriResult.error.message}`, 139 ); 140 } 141 return ok(undefined); 142 } 143 const authorDid = atUriResult.value.did.value; 144 145 // Resolve existing collection 146 const collectionIdResult = 147 await this.atUriResolutionService.resolveCollectionId(request.atUri); 148 if (collectionIdResult.isErr()) { 149 if (ENABLE_FIREHOSE_LOGGING) { 150 console.warn( 151 `[FirehoseWorker] Failed to resolve collection ID - user: ${authorDid}, uri: ${request.atUri}, error: ${collectionIdResult.error.message}`, 152 ); 153 } 154 return ok(undefined); 155 } 156 157 if (!collectionIdResult.value) { 158 if (ENABLE_FIREHOSE_LOGGING) { 159 console.log( 160 `[FirehoseWorker] Collection not found in our system - user: ${authorDid}, uri: ${request.atUri}`, 161 ); 162 } 163 return ok(undefined); 164 } 165 166 const publishedRecordId = PublishedRecordId.create({ 167 uri: request.atUri, 168 cid: request.cid, 169 }); 170 171 const result = await this.updateCollectionUseCase.execute({ 172 collectionId: collectionIdResult.value.getStringValue(), 173 name: request.record.name, 174 description: request.record.description, 175 accessType: request.record.accessType as 176 | CollectionAccessType 177 | undefined, 178 curatorId: authorDid, 179 publishedRecordId: publishedRecordId, 180 }); 181 182 if (result.isErr()) { 183 if (ENABLE_FIREHOSE_LOGGING) { 184 console.warn( 185 `[FirehoseWorker] Failed to update collection - user: ${authorDid}, collectionId: ${collectionIdResult.value.getStringValue()}, uri: ${request.atUri}, error: ${result.error.message}`, 186 ); 187 } 188 return ok(undefined); 189 } 190 191 if (ENABLE_FIREHOSE_LOGGING) { 192 console.log( 193 `[FirehoseWorker] Successfully updated collection - user: ${authorDid}, collectionId: ${result.value.collectionId}, uri: ${request.atUri}`, 194 ); 195 } 196 return ok(undefined); 197 } catch (error) { 198 if (ENABLE_FIREHOSE_LOGGING) { 199 console.error( 200 `[FirehoseWorker] Error processing collection update event - uri: ${request.atUri}, error: ${error}`, 201 ); 202 } 203 return ok(undefined); // Don't fail the firehose processing 204 } 205 } 206 207 private async handleCollectionDelete( 208 request: ProcessCollectionFirehoseEventDTO, 209 ): Promise<Result<void>> { 210 try { 211 // Parse AT URI to extract author DID 212 const atUriResult = ATUri.create(request.atUri); 213 if (atUriResult.isErr()) { 214 if (ENABLE_FIREHOSE_LOGGING) { 215 console.warn( 216 `[FirehoseWorker] Invalid AT URI format: ${request.atUri} - ${atUriResult.error.message}`, 217 ); 218 } 219 return ok(undefined); 220 } 221 const authorDid = atUriResult.value.did.value; 222 223 const collectionIdResult = 224 await this.atUriResolutionService.resolveCollectionId(request.atUri); 225 if (collectionIdResult.isErr()) { 226 if (ENABLE_FIREHOSE_LOGGING) { 227 console.warn( 228 `[FirehoseWorker] Failed to resolve collection ID - user: ${authorDid}, uri: ${request.atUri}, error: ${collectionIdResult.error.message}`, 229 ); 230 } 231 return ok(undefined); 232 } 233 234 if (collectionIdResult.value) { 235 if (ENABLE_FIREHOSE_LOGGING) { 236 console.log( 237 `[FirehoseWorker] Collection deleted externally - user: ${authorDid}, collectionId: ${collectionIdResult.value.getStringValue()}, uri: ${request.atUri}`, 238 ); 239 } 240 241 const publishedRecordId = PublishedRecordId.create({ 242 uri: request.atUri, 243 cid: request.cid || 'deleted', 244 }); 245 246 const result = await this.deleteCollectionUseCase.execute({ 247 collectionId: collectionIdResult.value.getStringValue(), 248 curatorId: authorDid, 249 publishedRecordId: publishedRecordId, 250 }); 251 252 if (result.isErr()) { 253 if (ENABLE_FIREHOSE_LOGGING) { 254 console.warn( 255 `[FirehoseWorker] Failed to delete collection - user: ${authorDid}, collectionId: ${collectionIdResult.value.getStringValue()}, uri: ${request.atUri}, error: ${result.error.message}`, 256 ); 257 } 258 return ok(undefined); 259 } 260 261 if (ENABLE_FIREHOSE_LOGGING) { 262 console.log( 263 `[FirehoseWorker] Successfully deleted collection - user: ${authorDid}, collectionId: ${result.value.collectionId}, uri: ${request.atUri}`, 264 ); 265 } 266 } 267 268 return ok(undefined); 269 } catch (error) { 270 if (ENABLE_FIREHOSE_LOGGING) { 271 console.error( 272 `[FirehoseWorker] Error processing collection delete event - uri: ${request.atUri}, error: ${error}`, 273 ); 274 } 275 return ok(undefined); 276 } 277 } 278}