From 5b0e19adb28a9fd6101f063b4fa27ffec4f3fb25 Mon Sep 17 00:00:00 2001 From: Wesley Finck Date: Thu, 26 Feb 2026 19:01:27 -0800 Subject: [PATCH] main command use cases for connections --- .../application/ports/IConnectionPublisher.ts | 12 ++ .../commands/CreateConnectionUseCase.ts | 203 ++++++++++++++++++ .../commands/DeleteConnectionUseCase.ts | 153 +++++++++++++ .../commands/UpdateConnectionUseCase.ts | 166 ++++++++++++++ .../DrizzleConnectionRepository.ts | 12 +- 5 files changed, 540 insertions(+), 6 deletions(-) create mode 100644 src/modules/cards/application/ports/IConnectionPublisher.ts create mode 100644 src/modules/cards/application/useCases/commands/CreateConnectionUseCase.ts create mode 100644 src/modules/cards/application/useCases/commands/DeleteConnectionUseCase.ts create mode 100644 src/modules/cards/application/useCases/commands/UpdateConnectionUseCase.ts diff --git a/src/modules/cards/application/ports/IConnectionPublisher.ts b/src/modules/cards/application/ports/IConnectionPublisher.ts new file mode 100644 index 00000000..dfe5a777 --- /dev/null +++ b/src/modules/cards/application/ports/IConnectionPublisher.ts @@ -0,0 +1,12 @@ +import { Connection } from '../../domain/Connection'; +import { Result } from '../../../../shared/core/Result'; +import { UseCaseError } from '../../../../shared/core/UseCaseError'; +import { PublishedRecordId } from '../../domain/value-objects/PublishedRecordId'; + +export interface IConnectionPublisher { + publish( + connection: Connection, + ): Promise>; + + unpublish(recordId: PublishedRecordId): Promise>; +} diff --git a/src/modules/cards/application/useCases/commands/CreateConnectionUseCase.ts b/src/modules/cards/application/useCases/commands/CreateConnectionUseCase.ts new file mode 100644 index 00000000..f544c93a --- /dev/null +++ b/src/modules/cards/application/useCases/commands/CreateConnectionUseCase.ts @@ -0,0 +1,203 @@ +import { Result, ok, err } from '../../../../../shared/core/Result'; +import { BaseUseCase } from '../../../../../shared/core/UseCase'; +import { UseCaseError } from '../../../../../shared/core/UseCaseError'; +import { AppError } from '../../../../../shared/core/AppError'; +import { IEventPublisher } from '../../../../../shared/application/events/IEventPublisher'; +import { IConnectionRepository } from '../../../domain/IConnectionRepository'; +import { Connection } from '../../../domain/Connection'; +import { CuratorId } from '../../../domain/value-objects/CuratorId'; +import { PublishedRecordId } from '../../../domain/value-objects/PublishedRecordId'; +import { IConnectionPublisher } from '../../ports/IConnectionPublisher'; +import { AuthenticationError } from '../../../../../shared/core/AuthenticationError'; +import { + UrlOrCardId, + UrlOrCardIdType, +} from '../../../domain/value-objects/UrlOrCardId'; +import { + ConnectionType, + ConnectionTypeEnum, +} from '../../../domain/value-objects/ConnectionType'; +import { ConnectionNote } from '../../../domain/value-objects/ConnectionNote'; + +export interface CreateConnectionDTO { + sourceType: UrlOrCardIdType; + sourceValue: string; + targetType: UrlOrCardIdType; + targetValue: string; + connectionType?: ConnectionTypeEnum; + note?: 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 CreateConnectionResponseDTO { + connectionId: string; +} + +export class ValidationError extends UseCaseError { + constructor(message: string) { + super(message); + } +} + +export class CreateConnectionUseCase extends BaseUseCase< + CreateConnectionDTO, + Result< + CreateConnectionResponseDTO, + ValidationError | AuthenticationError | AppError.UnexpectedError + > +> { + constructor( + private connectionRepository: IConnectionRepository, + private connectionPublisher: IConnectionPublisher, + eventPublisher: IEventPublisher, + ) { + super(eventPublisher); + } + + async execute( + request: CreateConnectionDTO, + ): Promise< + Result< + CreateConnectionResponseDTO, + ValidationError | AuthenticationError | AppError.UnexpectedError + > + > { + try { + // Validate and create CuratorId + const curatorIdResult = CuratorId.create(request.curatorId); + if (curatorIdResult.isErr()) { + return err( + new ValidationError( + `Invalid curator ID: ${curatorIdResult.error.message}`, + ), + ); + } + const curatorId = curatorIdResult.value; + + // Create source UrlOrCardId + const sourceResult = UrlOrCardId.reconstruct( + request.sourceType, + request.sourceValue, + ); + if (sourceResult.isErr()) { + return err( + new ValidationError( + `Invalid source: ${sourceResult.error.message}`, + ), + ); + } + const source = sourceResult.value; + + // Create target UrlOrCardId + const targetResult = UrlOrCardId.reconstruct( + request.targetType, + request.targetValue, + ); + if (targetResult.isErr()) { + return err( + new ValidationError( + `Invalid target: ${targetResult.error.message}`, + ), + ); + } + const target = targetResult.value; + + // Create optional connection type + let connectionType: ConnectionType | undefined; + if (request.connectionType) { + const typeResult = ConnectionType.createFromString( + request.connectionType, + ); + if (typeResult.isErr()) { + return err( + new ValidationError( + `Invalid connection type: ${typeResult.error.message}`, + ), + ); + } + connectionType = typeResult.value; + } + + // Create optional note + let note: ConnectionNote | undefined; + if (request.note) { + const noteResult = ConnectionNote.create(request.note); + if (noteResult.isErr()) { + return err( + new ValidationError(`Invalid note: ${noteResult.error.message}`), + ); + } + note = noteResult.value; + } + + // Create connection + const timestamp = request.createdAt ?? new Date(); + const connectionResult = Connection.create({ + source, + target, + type: connectionType, + note, + curatorId, + createdAt: timestamp, + updatedAt: timestamp, + }); + + if (connectionResult.isErr()) { + return err(new ValidationError(connectionResult.error.message)); + } + + const connection = connectionResult.value; + + // Handle publishing - skip if publishedRecordId provided (firehose event) + if (request.publishedRecordId) { + // Mark connection as published with provided record ID + connection.markAsPublished(request.publishedRecordId); + } else { + // Publish connection normally + const publishResult = + await this.connectionPublisher.publish(connection); + if (publishResult.isErr()) { + // Propagate authentication errors + if (publishResult.error instanceof AuthenticationError) { + return err(publishResult.error); + } + return err( + new ValidationError( + `Failed to publish connection: ${publishResult.error.message}`, + ), + ); + } + + // Mark connection as published + connection.markAsPublished(publishResult.value); + } + + // Save updated connection with published record ID + const saveUpdatedResult = + await this.connectionRepository.save(connection); + if (saveUpdatedResult.isErr()) { + return err(AppError.UnexpectedError.create(saveUpdatedResult.error)); + } + + // Publish domain events (ConnectionCreatedEvent) + const publishEventsResult = await this.publishEventsForAggregate( + connection, + ); + if (publishEventsResult.isErr()) { + console.error( + 'Failed to publish domain events:', + publishEventsResult.error, + ); + // Don't fail the operation + } + + return ok({ + connectionId: connection.connectionId.getStringValue(), + }); + } catch (error) { + return err(AppError.UnexpectedError.create(error)); + } + } +} diff --git a/src/modules/cards/application/useCases/commands/DeleteConnectionUseCase.ts b/src/modules/cards/application/useCases/commands/DeleteConnectionUseCase.ts new file mode 100644 index 00000000..7eef8429 --- /dev/null +++ b/src/modules/cards/application/useCases/commands/DeleteConnectionUseCase.ts @@ -0,0 +1,153 @@ +import { Result, ok, err } from '../../../../../shared/core/Result'; +import { BaseUseCase } from '../../../../../shared/core/UseCase'; +import { UseCaseError } from '../../../../../shared/core/UseCaseError'; +import { AppError } from '../../../../../shared/core/AppError'; +import { IEventPublisher } from '../../../../../shared/application/events/IEventPublisher'; +import { IConnectionRepository } from '../../../domain/IConnectionRepository'; +import { ConnectionId } from '../../../domain/value-objects/ConnectionId'; +import { CuratorId } from '../../../domain/value-objects/CuratorId'; +import { PublishedRecordId } from '../../../domain/value-objects/PublishedRecordId'; +import { IConnectionPublisher } from '../../ports/IConnectionPublisher'; +import { AuthenticationError } from '../../../../../shared/core/AuthenticationError'; + +export interface DeleteConnectionDTO { + connectionId: string; + curatorId: string; + publishedRecordId?: PublishedRecordId; // For firehose events - skip unpublishing if provided +} + +export interface DeleteConnectionResponseDTO { + connectionId: string; +} + +export class ValidationError extends UseCaseError { + constructor(message: string) { + super(message); + } +} + +export class DeleteConnectionUseCase extends BaseUseCase< + DeleteConnectionDTO, + Result< + DeleteConnectionResponseDTO, + ValidationError | AuthenticationError | AppError.UnexpectedError + > +> { + constructor( + private connectionRepository: IConnectionRepository, + private connectionPublisher: IConnectionPublisher, + eventPublisher: IEventPublisher, + ) { + super(eventPublisher); + } + + async execute( + request: DeleteConnectionDTO, + ): Promise< + Result< + DeleteConnectionResponseDTO, + ValidationError | AuthenticationError | AppError.UnexpectedError + > + > { + try { + // Validate and create CuratorId + const curatorIdResult = CuratorId.create(request.curatorId); + if (curatorIdResult.isErr()) { + return err( + new ValidationError( + `Invalid curator ID: ${curatorIdResult.error.message}`, + ), + ); + } + const curatorId = curatorIdResult.value; + + // Validate and create ConnectionId + const connectionIdResult = ConnectionId.createFromString( + request.connectionId, + ); + if (connectionIdResult.isErr()) { + return err( + new ValidationError( + `Invalid connection ID: ${connectionIdResult.error.message}`, + ), + ); + } + const connectionId = connectionIdResult.value; + + // Find the connection + const connectionResult = + await this.connectionRepository.findById(connectionId); + if (connectionResult.isErr()) { + return err(AppError.UnexpectedError.create(connectionResult.error)); + } + + const connection = connectionResult.value; + if (!connection) { + return err( + new ValidationError(`Connection not found: ${request.connectionId}`), + ); + } + + // Check if user is the curator + if (!connection.curatorId.equals(curatorId)) { + return err( + new ValidationError( + 'Only the connection curator can delete the connection', + ), + ); + } + + // Mark connection for removal (raises ConnectionRemovedEvent) + const removalResult = connection.markForRemoval(); + if (removalResult.isErr()) { + return err(new ValidationError(removalResult.error.message)); + } + + // Handle unpublishing - skip if publishedRecordId provided (firehose event) + if ( + !request.publishedRecordId && + connection.isPublished && + connection.publishedRecordId + ) { + const unpublishResult = await this.connectionPublisher.unpublish( + connection.publishedRecordId, + ); + if (unpublishResult.isErr()) { + // Propagate authentication errors + if (unpublishResult.error instanceof AuthenticationError) { + return err(unpublishResult.error); + } + return err( + new ValidationError( + `Failed to unpublish connection: ${unpublishResult.error.message}`, + ), + ); + } + } + + // Publish domain events (ConnectionRemovedEvent) + const publishEventsResult = await this.publishEventsForAggregate( + connection, + ); + if (publishEventsResult.isErr()) { + console.error( + 'Failed to publish domain events:', + publishEventsResult.error, + ); + // Don't fail the operation + } + + // Delete connection from repository + const deleteResult = await this.connectionRepository.delete(connectionId); + if (deleteResult.isErr()) { + return err(AppError.UnexpectedError.create(deleteResult.error)); + } + + return ok({ + connectionId: connection.connectionId.getStringValue(), + }); + } catch (error) { + return err(AppError.UnexpectedError.create(error)); + } + } +} diff --git a/src/modules/cards/application/useCases/commands/UpdateConnectionUseCase.ts b/src/modules/cards/application/useCases/commands/UpdateConnectionUseCase.ts new file mode 100644 index 00000000..bf8816ad --- /dev/null +++ b/src/modules/cards/application/useCases/commands/UpdateConnectionUseCase.ts @@ -0,0 +1,166 @@ +import { Result, ok, err } from '../../../../../shared/core/Result'; +import { UseCase } from '../../../../../shared/core/UseCase'; +import { UseCaseError } from '../../../../../shared/core/UseCaseError'; +import { AppError } from '../../../../../shared/core/AppError'; +import { IConnectionRepository } from '../../../domain/IConnectionRepository'; +import { ConnectionId } from '../../../domain/value-objects/ConnectionId'; +import { CuratorId } from '../../../domain/value-objects/CuratorId'; +import { PublishedRecordId } from '../../../domain/value-objects/PublishedRecordId'; +import { IConnectionPublisher } from '../../ports/IConnectionPublisher'; +import { AuthenticationError } from '../../../../../shared/core/AuthenticationError'; +import { ConnectionNote } from '../../../domain/value-objects/ConnectionNote'; + +export interface UpdateConnectionDTO { + connectionId: string; + note?: string; + removeNote?: boolean; + curatorId: string; + publishedRecordId?: PublishedRecordId; // For firehose events - skip republishing if provided +} + +export interface UpdateConnectionResponseDTO { + connectionId: string; +} + +export class ValidationError extends UseCaseError { + constructor(message: string) { + super(message); + } +} + +export class UpdateConnectionUseCase + implements + UseCase< + UpdateConnectionDTO, + Result< + UpdateConnectionResponseDTO, + ValidationError | AuthenticationError | AppError.UnexpectedError + > + > +{ + constructor( + private connectionRepository: IConnectionRepository, + private connectionPublisher: IConnectionPublisher, + ) {} + + async execute( + request: UpdateConnectionDTO, + ): Promise< + Result< + UpdateConnectionResponseDTO, + ValidationError | AuthenticationError | AppError.UnexpectedError + > + > { + try { + // Validate and create CuratorId + const curatorIdResult = CuratorId.create(request.curatorId); + if (curatorIdResult.isErr()) { + return err( + new ValidationError( + `Invalid curator ID: ${curatorIdResult.error.message}`, + ), + ); + } + const curatorId = curatorIdResult.value; + + // Validate and create ConnectionId + const connectionIdResult = ConnectionId.createFromString( + request.connectionId, + ); + if (connectionIdResult.isErr()) { + return err( + new ValidationError( + `Invalid connection ID: ${connectionIdResult.error.message}`, + ), + ); + } + const connectionId = connectionIdResult.value; + + // Find the connection + const connectionResult = + await this.connectionRepository.findById(connectionId); + if (connectionResult.isErr()) { + return err(AppError.UnexpectedError.create(connectionResult.error)); + } + + const connection = connectionResult.value; + if (!connection) { + return err( + new ValidationError(`Connection not found: ${request.connectionId}`), + ); + } + + // Check if user is the curator + if (!connection.curatorId.equals(curatorId)) { + return err( + new ValidationError( + 'Only the connection curator can update the connection', + ), + ); + } + + // Handle note update/removal + if (request.removeNote) { + const removeResult = connection.removeNote(); + if (removeResult.isErr()) { + return err(new ValidationError(removeResult.error.message)); + } + } else if (request.note !== undefined) { + const noteResult = ConnectionNote.create(request.note); + if (noteResult.isErr()) { + return err( + new ValidationError(`Invalid note: ${noteResult.error.message}`), + ); + } + const updateResult = connection.updateNote(noteResult.value); + if (updateResult.isErr()) { + return err(new ValidationError(updateResult.error.message)); + } + } + + // Handle republishing - skip if publishedRecordId provided (firehose event) + if (request.publishedRecordId) { + // Update published record ID with provided value + connection.markAsPublished(request.publishedRecordId); + + // Save connection with updated published record ID + const saveUpdatedResult = + await this.connectionRepository.save(connection); + if (saveUpdatedResult.isErr()) { + return err(AppError.UnexpectedError.create(saveUpdatedResult.error)); + } + } else if (connection.isPublished) { + // Republish connection normally + const republishResult = + await this.connectionPublisher.publish(connection); + if (republishResult.isErr()) { + // Propagate authentication errors + if (republishResult.error instanceof AuthenticationError) { + return err(republishResult.error); + } + return err( + new ValidationError( + `Failed to republish connection: ${republishResult.error.message}`, + ), + ); + } + + // Update published record ID + connection.markAsPublished(republishResult.value); + + // Save connection with updated published record ID + const saveUpdatedResult = + await this.connectionRepository.save(connection); + if (saveUpdatedResult.isErr()) { + return err(AppError.UnexpectedError.create(saveUpdatedResult.error)); + } + } + + return ok({ + connectionId: connection.connectionId.getStringValue(), + }); + } catch (error) { + return err(AppError.UnexpectedError.create(error)); + } + } +} diff --git a/src/modules/cards/infrastructure/repositories/DrizzleConnectionRepository.ts b/src/modules/cards/infrastructure/repositories/DrizzleConnectionRepository.ts index ba0a34e8..add3a34d 100644 --- a/src/modules/cards/infrastructure/repositories/DrizzleConnectionRepository.ts +++ b/src/modules/cards/infrastructure/repositories/DrizzleConnectionRepository.ts @@ -46,7 +46,7 @@ export class DrizzleConnectionRepository implements IConnectionRepository { sourceValue: result.connection.sourceValue, targetType: result.connection.targetType, targetValue: result.connection.targetValue, - connectionType: result.connection.connectionType, + connectionType: result.connection.connectionType || undefined, note: result.connection.note || undefined, createdAt: result.connection.createdAt, updatedAt: result.connection.updatedAt, @@ -96,7 +96,7 @@ export class DrizzleConnectionRepository implements IConnectionRepository { sourceValue: result.connection.sourceValue, targetType: result.connection.targetType, targetValue: result.connection.targetValue, - connectionType: result.connection.connectionType, + connectionType: result.connection.connectionType || undefined, note: result.connection.note || undefined, createdAt: result.connection.createdAt, updatedAt: result.connection.updatedAt, @@ -149,7 +149,7 @@ export class DrizzleConnectionRepository implements IConnectionRepository { sourceValue: result.connection.sourceValue, targetType: result.connection.targetType, targetValue: result.connection.targetValue, - connectionType: result.connection.connectionType, + connectionType: result.connection.connectionType || undefined, note: result.connection.note || undefined, createdAt: result.connection.createdAt, updatedAt: result.connection.updatedAt, @@ -205,7 +205,7 @@ export class DrizzleConnectionRepository implements IConnectionRepository { sourceValue: result.connection.sourceValue, targetType: result.connection.targetType, targetValue: result.connection.targetValue, - connectionType: result.connection.connectionType, + connectionType: result.connection.connectionType || undefined, note: result.connection.note || undefined, createdAt: result.connection.createdAt, updatedAt: result.connection.updatedAt, @@ -261,7 +261,7 @@ export class DrizzleConnectionRepository implements IConnectionRepository { sourceValue: result.connection.sourceValue, targetType: result.connection.targetType, targetValue: result.connection.targetValue, - connectionType: result.connection.connectionType, + connectionType: result.connection.connectionType || undefined, note: result.connection.note || undefined, createdAt: result.connection.createdAt, updatedAt: result.connection.updatedAt, @@ -322,7 +322,7 @@ export class DrizzleConnectionRepository implements IConnectionRepository { sourceValue: result.connection.sourceValue, targetType: result.connection.targetType, targetValue: result.connection.targetValue, - connectionType: result.connection.connectionType, + connectionType: result.connection.connectionType || undefined, note: result.connection.note || undefined, createdAt: result.connection.createdAt, updatedAt: result.connection.updatedAt, -- 2.51.2