This repository has no description
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}