diff --git a/src/modules/cards/domain/Connection.ts b/src/modules/cards/domain/Connection.ts new file mode 100644 index 00000000..3e342ff3 --- /dev/null +++ b/src/modules/cards/domain/Connection.ts @@ -0,0 +1,155 @@ +import { AggregateRoot } from '../../../shared/domain/AggregateRoot'; +import { UniqueEntityID } from '../../../shared/domain/UniqueEntityID'; +import { ok, err, Result } from '../../../shared/core/Result'; +import { ConnectionId } from './value-objects/ConnectionId'; +import { ConnectionType } from './value-objects/ConnectionType'; +import { UrlOrCardId } from './value-objects/UrlOrCardId'; +import { ConnectionNote } from './value-objects/ConnectionNote'; +import { CuratorId } from './value-objects/CuratorId'; +import { PublishedRecordId } from './value-objects/PublishedRecordId'; +import { ConnectionCreatedEvent } from './events/ConnectionCreatedEvent'; +import { ConnectionRemovedEvent } from './events/ConnectionRemovedEvent'; + +export class ConnectionValidationError extends Error { + constructor(message: string) { + super(message); + this.name = 'ConnectionValidationError'; + } +} + +interface ConnectionProps { + source: UrlOrCardId; + target: UrlOrCardId; + type?: ConnectionType; + note?: ConnectionNote; + curatorId: CuratorId; + publishedRecordId?: PublishedRecordId; + createdAt: Date; + updatedAt: Date; +} + +export class Connection extends AggregateRoot { + get connectionId(): ConnectionId { + return ConnectionId.create(this._id).unwrap(); + } + + get source(): UrlOrCardId { + return this.props.source; + } + + get target(): UrlOrCardId { + return this.props.target; + } + + get type(): ConnectionType | undefined { + return this.props.type; + } + + get note(): ConnectionNote | undefined { + return this.props.note; + } + + get curatorId(): CuratorId { + return this.props.curatorId; + } + + get publishedRecordId(): PublishedRecordId | undefined { + return this.props.publishedRecordId; + } + + get createdAt(): Date { + return this.props.createdAt; + } + + get updatedAt(): Date { + return this.props.updatedAt; + } + + get isPublished(): boolean { + return this.props.publishedRecordId !== undefined; + } + + private constructor(props: ConnectionProps, id?: UniqueEntityID) { + super(props, id); + } + + public static create( + props: Omit & { + createdAt?: Date; + updatedAt?: Date; + }, + id?: UniqueEntityID, + ): Result { + // Validate that source and target are not the same + if (props.source.equals(props.target)) { + return err( + new ConnectionValidationError( + 'Connection source and target cannot be the same', + ), + ); + } + + const now = new Date(); + const connectionProps: ConnectionProps = { + ...props, + createdAt: props.createdAt || now, + updatedAt: props.updatedAt || now, + }; + + const connection = new Connection(connectionProps, id); + connection.raiseCreatedEvent(); + return ok(connection); + } + + public updateNote( + note: ConnectionNote, + ): Result { + this.props.note = note; + this.props.updatedAt = new Date(); + return ok(undefined); + } + + public removeNote(): Result { + this.props.note = undefined; + this.props.updatedAt = new Date(); + return ok(undefined); + } + + private raiseCreatedEvent(): Result { + const event = ConnectionCreatedEvent.create( + this.connectionId, + this.curatorId, + ); + + if (event.isErr()) { + return err(new Error(event.error.message)); + } + + this.addDomainEvent(event.value); + return ok(undefined); + } + + public markForRemoval(): Result { + const event = ConnectionRemovedEvent.create( + this.connectionId, + this.curatorId, + ); + + if (event.isErr()) { + return err(new Error(event.error.message)); + } + + this.addDomainEvent(event.value); + return ok(undefined); + } + + public markAsPublished(publishedRecordId: PublishedRecordId): void { + this.props.publishedRecordId = publishedRecordId; + this.props.updatedAt = new Date(); + } + + public markAsUnpublished(): void { + this.props.publishedRecordId = undefined; + this.props.updatedAt = new Date(); + } +} diff --git a/src/modules/cards/domain/IConnectionRepository.ts b/src/modules/cards/domain/IConnectionRepository.ts new file mode 100644 index 00000000..9ef27e20 --- /dev/null +++ b/src/modules/cards/domain/IConnectionRepository.ts @@ -0,0 +1,19 @@ +import { Result } from '../../../shared/core/Result'; +import { Connection } from './Connection'; +import { ConnectionId } from './value-objects/ConnectionId'; +import { UrlOrCardId } from './value-objects/UrlOrCardId'; +import { CuratorId } from './value-objects/CuratorId'; + +export interface IConnectionRepository { + findById(id: ConnectionId): Promise>; + findByIds(ids: ConnectionId[]): Promise>; + findByCuratorId(curatorId: CuratorId): Promise>; + findBySource(source: UrlOrCardId): Promise>; + findByTarget(target: UrlOrCardId): Promise>; + findBetween( + source: UrlOrCardId, + target: UrlOrCardId, + ): Promise>; + save(connection: Connection): Promise>; + delete(connectionId: ConnectionId): Promise>; +} diff --git a/src/modules/cards/domain/events/ConnectionCreatedEvent.ts b/src/modules/cards/domain/events/ConnectionCreatedEvent.ts new file mode 100644 index 00000000..3806bb5a --- /dev/null +++ b/src/modules/cards/domain/events/ConnectionCreatedEvent.ts @@ -0,0 +1,40 @@ +import { IDomainEvent } from '../../../../shared/domain/events/IDomainEvent'; +import { UniqueEntityID } from '../../../../shared/domain/UniqueEntityID'; +import { ConnectionId } from '../value-objects/ConnectionId'; +import { CuratorId } from '../value-objects/CuratorId'; +import { EventNames } from '../../../../shared/infrastructure/events/EventConfig'; +import { Result, ok } from '../../../../shared/core/Result'; + +export class ConnectionCreatedEvent implements IDomainEvent { + public readonly eventName = EventNames.CONNECTION_CREATED; + public readonly dateTimeOccurred: Date; + + private constructor( + public readonly connectionId: ConnectionId, + public readonly curatorId: CuratorId, + dateTimeOccurred?: Date, + ) { + this.dateTimeOccurred = dateTimeOccurred || new Date(); + } + + public static create( + connectionId: ConnectionId, + curatorId: CuratorId, + ): Result { + return ok(new ConnectionCreatedEvent(connectionId, curatorId)); + } + + public static reconstruct( + connectionId: ConnectionId, + curatorId: CuratorId, + dateTimeOccurred: Date, + ): Result { + return ok( + new ConnectionCreatedEvent(connectionId, curatorId, dateTimeOccurred), + ); + } + + getAggregateId(): UniqueEntityID { + return this.connectionId.getValue(); + } +} diff --git a/src/modules/cards/domain/events/ConnectionRemovedEvent.ts b/src/modules/cards/domain/events/ConnectionRemovedEvent.ts new file mode 100644 index 00000000..28ecfb47 --- /dev/null +++ b/src/modules/cards/domain/events/ConnectionRemovedEvent.ts @@ -0,0 +1,40 @@ +import { IDomainEvent } from '../../../../shared/domain/events/IDomainEvent'; +import { UniqueEntityID } from '../../../../shared/domain/UniqueEntityID'; +import { ConnectionId } from '../value-objects/ConnectionId'; +import { CuratorId } from '../value-objects/CuratorId'; +import { EventNames } from '../../../../shared/infrastructure/events/EventConfig'; +import { Result, ok } from '../../../../shared/core/Result'; + +export class ConnectionRemovedEvent implements IDomainEvent { + public readonly eventName = EventNames.CONNECTION_REMOVED; + public readonly dateTimeOccurred: Date; + + private constructor( + public readonly connectionId: ConnectionId, + public readonly curatorId: CuratorId, + dateTimeOccurred?: Date, + ) { + this.dateTimeOccurred = dateTimeOccurred || new Date(); + } + + public static create( + connectionId: ConnectionId, + curatorId: CuratorId, + ): Result { + return ok(new ConnectionRemovedEvent(connectionId, curatorId)); + } + + public static reconstruct( + connectionId: ConnectionId, + curatorId: CuratorId, + dateTimeOccurred: Date, + ): Result { + return ok( + new ConnectionRemovedEvent(connectionId, curatorId, dateTimeOccurred), + ); + } + + getAggregateId(): UniqueEntityID { + return this.connectionId.getValue(); + } +} diff --git a/src/modules/cards/domain/value-objects/ConnectionId.ts b/src/modules/cards/domain/value-objects/ConnectionId.ts new file mode 100644 index 00000000..ba8bcea0 --- /dev/null +++ b/src/modules/cards/domain/value-objects/ConnectionId.ts @@ -0,0 +1,39 @@ +import { ok, Result, err } from 'src/shared/core/Result'; +import { Guard } from 'src/shared/core/Guard'; +import { UniqueEntityID } from 'src/shared/domain/UniqueEntityID'; +import { ValueObject } from 'src/shared/domain/ValueObject'; + +interface ConnectionIdProps { + value: UniqueEntityID; +} + +export class ConnectionId extends ValueObject { + getStringValue(): string { + return this.props.value.toString(); + } + + getValue(): UniqueEntityID { + return this.props.value; + } + + private constructor(value: UniqueEntityID) { + super({ value }); + } + + public static create(id: UniqueEntityID): Result { + const guardResult = Guard.againstNullOrUndefined(id, 'id'); + if (guardResult.isErr()) { + return err(new Error(guardResult.error)); + } + return ok(new ConnectionId(id)); + } + + public static createFromString(value: string): Result { + const guardResult = Guard.againstNullOrUndefined(value, 'value'); + if (guardResult.isErr()) { + return err(new Error(guardResult.error)); + } + const uniqueEntityID = new UniqueEntityID(value); + return ok(new ConnectionId(uniqueEntityID)); + } +} diff --git a/src/modules/cards/domain/value-objects/ConnectionNote.ts b/src/modules/cards/domain/value-objects/ConnectionNote.ts new file mode 100644 index 00000000..fffc7e30 --- /dev/null +++ b/src/modules/cards/domain/value-objects/ConnectionNote.ts @@ -0,0 +1,51 @@ +import { ValueObject } from '../../../../shared/domain/ValueObject'; +import { Result, ok, err } from '../../../../shared/core/Result'; + +export class InvalidConnectionNoteError extends Error { + constructor(message: string) { + super(message); + this.name = 'InvalidConnectionNoteError'; + } +} + +interface ConnectionNoteProps { + value: string; +} + +export class ConnectionNote extends ValueObject { + public static readonly MAX_LENGTH = 1000; + + get value(): string { + return this.props.value; + } + + private constructor(props: ConnectionNoteProps) { + super(props); + } + + public static create( + note: string, + ): Result { + const trimmedNote = note.trim(); + + if (trimmedNote.length === 0) { + return err( + new InvalidConnectionNoteError('Connection note cannot be empty'), + ); + } + + if (trimmedNote.length > this.MAX_LENGTH) { + return err( + new InvalidConnectionNoteError( + `Connection note cannot exceed ${this.MAX_LENGTH} characters`, + ), + ); + } + + return ok(new ConnectionNote({ value: trimmedNote })); + } + + public toString(): string { + return this.value; + } +} diff --git a/src/modules/cards/domain/value-objects/ConnectionType.ts b/src/modules/cards/domain/value-objects/ConnectionType.ts new file mode 100644 index 00000000..3f252bb4 --- /dev/null +++ b/src/modules/cards/domain/value-objects/ConnectionType.ts @@ -0,0 +1,94 @@ +import { ok, Result, err } from '../../../../shared/core/Result'; +import { ValueObject } from '../../../../shared/domain/ValueObject'; + +export enum ConnectionTypeEnum { + SUPPORTS = 'SUPPORTS', + OPPOSES = 'OPPOSES', + ADDRESSES = 'ADDRESSES', + HELPFUL = 'HELPFUL', + LEADS_TO = 'LEADS_TO', + RELATED = 'RELATED', + SUPPLEMENT = 'SUPPLEMENT', + EXPLAINER = 'EXPLAINER', +} + +// Metadata about each connection type +interface ConnectionTypeMetadata { + isDirectional: boolean; + displayName: string; +} + +const CONNECTION_TYPE_METADATA: Record< + ConnectionTypeEnum, + ConnectionTypeMetadata +> = { + [ConnectionTypeEnum.SUPPORTS]: { + isDirectional: true, + displayName: 'Supports', + }, + [ConnectionTypeEnum.OPPOSES]: { + isDirectional: true, + displayName: 'Opposes', + }, + [ConnectionTypeEnum.ADDRESSES]: { + isDirectional: true, + displayName: 'Addresses', + }, + [ConnectionTypeEnum.HELPFUL]: { + isDirectional: false, + displayName: 'Helpful', + }, + [ConnectionTypeEnum.LEADS_TO]: { + isDirectional: true, + displayName: 'Leads to', + }, + [ConnectionTypeEnum.RELATED]: { + isDirectional: false, + displayName: 'Related', + }, + [ConnectionTypeEnum.SUPPLEMENT]: { + isDirectional: true, + displayName: 'Supplement', + }, + [ConnectionTypeEnum.EXPLAINER]: { + isDirectional: true, + displayName: 'Explainer', + }, +}; + +interface ConnectionTypeProps { + value: ConnectionTypeEnum; +} + +export class ConnectionType extends ValueObject { + get value(): ConnectionTypeEnum { + return this.props.value; + } + + get isDirectional(): boolean { + return CONNECTION_TYPE_METADATA[this.props.value].isDirectional; + } + + get displayName(): string { + return CONNECTION_TYPE_METADATA[this.props.value].displayName; + } + + private constructor(props: ConnectionTypeProps) { + super(props); + } + + public static create(type: ConnectionTypeEnum): Result { + if (!Object.values(ConnectionTypeEnum).includes(type)) { + return err(new Error(`Invalid connection type: ${type}`)); + } + return ok(new ConnectionType({ value: type })); + } + + public static createFromString(type: string): Result { + const upperType = type.toUpperCase(); + if (!Object.values(ConnectionTypeEnum).includes(upperType as any)) { + return err(new Error(`Invalid connection type: ${type}`)); + } + return ok(new ConnectionType({ value: upperType as ConnectionTypeEnum })); + } +} diff --git a/src/modules/cards/domain/value-objects/UrlOrCardId.ts b/src/modules/cards/domain/value-objects/UrlOrCardId.ts new file mode 100644 index 00000000..0cf7815b --- /dev/null +++ b/src/modules/cards/domain/value-objects/UrlOrCardId.ts @@ -0,0 +1,112 @@ +import { ok, Result, err } from '../../../../shared/core/Result'; +import { ValueObject } from '../../../../shared/domain/ValueObject'; +import { URL } from './URL'; +import { CardId } from './CardId'; + +export enum UrlOrCardIdType { + URL = 'URL', + CARD = 'CARD', +} + +// Discriminated union type +type UrlOrCardIdValue = + | { type: UrlOrCardIdType.URL; url: URL } + | { type: UrlOrCardIdType.CARD; cardId: CardId }; + +interface UrlOrCardIdProps { + value: UrlOrCardIdValue; +} + +export class UrlOrCardId extends ValueObject { + get type(): UrlOrCardIdType { + return this.props.value.type; + } + + get isUrl(): boolean { + return this.props.value.type === UrlOrCardIdType.URL; + } + + get isCard(): boolean { + return this.props.value.type === UrlOrCardIdType.CARD; + } + + get url(): URL | null { + if (this.props.value.type === UrlOrCardIdType.URL) { + return this.props.value.url; + } + return null; + } + + get cardId(): CardId | null { + if (this.props.value.type === UrlOrCardIdType.CARD) { + return this.props.value.cardId; + } + return null; + } + + // Get the string representation for persistence + get stringValue(): string { + if (this.props.value.type === UrlOrCardIdType.URL) { + return this.props.value.url.value; + } else { + return this.props.value.cardId.getStringValue(); + } + } + + private constructor(props: UrlOrCardIdProps) { + super(props); + } + + public static createFromUrl(url: URL): Result { + return ok( + new UrlOrCardId({ + value: { type: UrlOrCardIdType.URL, url }, + }), + ); + } + + public static createFromCard(cardId: CardId): Result { + return ok( + new UrlOrCardId({ + value: { type: UrlOrCardIdType.CARD, cardId }, + }), + ); + } + + // Factory method for reconstruction from persistence + public static reconstruct( + type: UrlOrCardIdType, + value: string, + ): Result { + if (type === UrlOrCardIdType.URL) { + const urlResult = URL.create(value); + if (urlResult.isErr()) { + return err(urlResult.error); + } + return UrlOrCardId.createFromUrl(urlResult.value); + } else if (type === UrlOrCardIdType.CARD) { + const cardIdResult = CardId.createFromString(value); + if (cardIdResult.isErr()) { + return err(cardIdResult.error); + } + return UrlOrCardId.createFromCard(cardIdResult.value); + } else { + return err(new Error(`Invalid UrlOrCardId type: ${type}`)); + } + } + + public equals(vo?: ValueObject): boolean { + if (vo === null || vo === undefined) return false; + if (!(vo instanceof UrlOrCardId)) return false; + + if (this.type !== vo.type) return false; + + if (this.isUrl && vo.isUrl) { + return this.url!.equals(vo.url!); + } else if (this.isCard && vo.isCard) { + return this.cardId!.equals(vo.cardId!); + } + + return false; + } +} diff --git a/src/modules/cards/infrastructure/repositories/DrizzleConnectionRepository.ts b/src/modules/cards/infrastructure/repositories/DrizzleConnectionRepository.ts new file mode 100644 index 00000000..ba0a34e8 --- /dev/null +++ b/src/modules/cards/infrastructure/repositories/DrizzleConnectionRepository.ts @@ -0,0 +1,433 @@ +import { eq, inArray, and } from 'drizzle-orm'; +import { PostgresJsDatabase } from 'drizzle-orm/postgres-js'; +import { IConnectionRepository } from '../../domain/IConnectionRepository'; +import { Connection } from '../../domain/Connection'; +import { ConnectionId } from '../../domain/value-objects/ConnectionId'; +import { UrlOrCardId } from '../../domain/value-objects/UrlOrCardId'; +import { CuratorId } from '../../domain/value-objects/CuratorId'; +import { connections } from './schema/connection.sql'; +import { publishedRecords } from './schema/publishedRecord.sql'; +import { ConnectionDTO, ConnectionMapper } from './mappers/ConnectionMapper'; +import { Result, ok, err } from '../../../../shared/core/Result'; + +export class DrizzleConnectionRepository implements IConnectionRepository { + constructor(private db: PostgresJsDatabase) {} + + async findById(id: ConnectionId): Promise> { + try { + const connectionId = id.getStringValue(); + + const connectionResult = await this.db + .select({ + connection: connections, + publishedRecord: publishedRecords, + }) + .from(connections) + .leftJoin( + publishedRecords, + eq(connections.publishedRecordId, publishedRecords.id), + ) + .where(eq(connections.id, connectionId)) + .limit(1); + + if (connectionResult.length === 0) { + return ok(null); + } + + const result = connectionResult[0]; + if (!result || !result.connection) { + return ok(null); + } + + const connectionDTO: ConnectionDTO = { + id: result.connection.id, + curatorId: result.connection.curatorId, + sourceType: result.connection.sourceType, + sourceValue: result.connection.sourceValue, + targetType: result.connection.targetType, + targetValue: result.connection.targetValue, + connectionType: result.connection.connectionType, + note: result.connection.note || undefined, + createdAt: result.connection.createdAt, + updatedAt: result.connection.updatedAt, + publishedRecordId: result.publishedRecord?.id || null, + publishedRecord: result.publishedRecord || undefined, + }; + + const domainResult = ConnectionMapper.toDomain(connectionDTO); + if (domainResult.isErr()) { + return err(domainResult.error); + } + + return ok(domainResult.value); + } catch (error) { + return err(error as Error); + } + } + + async findByIds(ids: ConnectionId[]): Promise> { + try { + if (ids.length === 0) { + return ok([]); + } + + const connectionIds = ids.map((id) => id.getStringValue()); + + const connectionResults = await this.db + .select({ + connection: connections, + publishedRecord: publishedRecords, + }) + .from(connections) + .leftJoin( + publishedRecords, + eq(connections.publishedRecordId, publishedRecords.id), + ) + .where(inArray(connections.id, connectionIds)); + + const domainConnections: Connection[] = []; + for (const result of connectionResults) { + if (!result.connection) continue; + + const connectionDTO: ConnectionDTO = { + id: result.connection.id, + curatorId: result.connection.curatorId, + sourceType: result.connection.sourceType, + sourceValue: result.connection.sourceValue, + targetType: result.connection.targetType, + targetValue: result.connection.targetValue, + connectionType: result.connection.connectionType, + note: result.connection.note || undefined, + createdAt: result.connection.createdAt, + updatedAt: result.connection.updatedAt, + publishedRecordId: result.publishedRecord?.id || null, + publishedRecord: result.publishedRecord || undefined, + }; + + const domainResult = ConnectionMapper.toDomain(connectionDTO); + if (domainResult.isErr()) { + console.error( + 'Error mapping connection to domain:', + domainResult.error, + ); + continue; + } + domainConnections.push(domainResult.value); + } + + return ok(domainConnections); + } catch (error) { + return err(error as Error); + } + } + + async findByCuratorId(curatorId: CuratorId): Promise> { + try { + const curatorIdString = curatorId.value; + + const connectionResults = await this.db + .select({ + connection: connections, + publishedRecord: publishedRecords, + }) + .from(connections) + .leftJoin( + publishedRecords, + eq(connections.publishedRecordId, publishedRecords.id), + ) + .where(eq(connections.curatorId, curatorIdString)) + .orderBy(connections.createdAt); + + const domainConnections: Connection[] = []; + for (const result of connectionResults) { + if (!result.connection) continue; + + const connectionDTO: ConnectionDTO = { + id: result.connection.id, + curatorId: result.connection.curatorId, + sourceType: result.connection.sourceType, + sourceValue: result.connection.sourceValue, + targetType: result.connection.targetType, + targetValue: result.connection.targetValue, + connectionType: result.connection.connectionType, + note: result.connection.note || undefined, + createdAt: result.connection.createdAt, + updatedAt: result.connection.updatedAt, + publishedRecordId: result.publishedRecord?.id || null, + publishedRecord: result.publishedRecord || undefined, + }; + + const domainResult = ConnectionMapper.toDomain(connectionDTO); + if (domainResult.isErr()) { + console.error( + 'Error mapping connection to domain:', + domainResult.error, + ); + continue; + } + domainConnections.push(domainResult.value); + } + + return ok(domainConnections); + } catch (error) { + return err(error as Error); + } + } + + async findBySource(source: UrlOrCardId): Promise> { + try { + const connectionResults = await this.db + .select({ + connection: connections, + publishedRecord: publishedRecords, + }) + .from(connections) + .leftJoin( + publishedRecords, + eq(connections.publishedRecordId, publishedRecords.id), + ) + .where( + and( + eq(connections.sourceType, source.type), + eq(connections.sourceValue, source.stringValue), + ), + ) + .orderBy(connections.createdAt); + + const domainConnections: Connection[] = []; + for (const result of connectionResults) { + if (!result.connection) continue; + + const connectionDTO: ConnectionDTO = { + id: result.connection.id, + curatorId: result.connection.curatorId, + sourceType: result.connection.sourceType, + sourceValue: result.connection.sourceValue, + targetType: result.connection.targetType, + targetValue: result.connection.targetValue, + connectionType: result.connection.connectionType, + note: result.connection.note || undefined, + createdAt: result.connection.createdAt, + updatedAt: result.connection.updatedAt, + publishedRecordId: result.publishedRecord?.id || null, + publishedRecord: result.publishedRecord || undefined, + }; + + const domainResult = ConnectionMapper.toDomain(connectionDTO); + if (domainResult.isErr()) { + console.error( + 'Error mapping connection to domain:', + domainResult.error, + ); + continue; + } + domainConnections.push(domainResult.value); + } + + return ok(domainConnections); + } catch (error) { + return err(error as Error); + } + } + + async findByTarget(target: UrlOrCardId): Promise> { + try { + const connectionResults = await this.db + .select({ + connection: connections, + publishedRecord: publishedRecords, + }) + .from(connections) + .leftJoin( + publishedRecords, + eq(connections.publishedRecordId, publishedRecords.id), + ) + .where( + and( + eq(connections.targetType, target.type), + eq(connections.targetValue, target.stringValue), + ), + ) + .orderBy(connections.createdAt); + + const domainConnections: Connection[] = []; + for (const result of connectionResults) { + if (!result.connection) continue; + + const connectionDTO: ConnectionDTO = { + id: result.connection.id, + curatorId: result.connection.curatorId, + sourceType: result.connection.sourceType, + sourceValue: result.connection.sourceValue, + targetType: result.connection.targetType, + targetValue: result.connection.targetValue, + connectionType: result.connection.connectionType, + note: result.connection.note || undefined, + createdAt: result.connection.createdAt, + updatedAt: result.connection.updatedAt, + publishedRecordId: result.publishedRecord?.id || null, + publishedRecord: result.publishedRecord || undefined, + }; + + const domainResult = ConnectionMapper.toDomain(connectionDTO); + if (domainResult.isErr()) { + console.error( + 'Error mapping connection to domain:', + domainResult.error, + ); + continue; + } + domainConnections.push(domainResult.value); + } + + return ok(domainConnections); + } catch (error) { + return err(error as Error); + } + } + + async findBetween( + source: UrlOrCardId, + target: UrlOrCardId, + ): Promise> { + try { + const connectionResults = await this.db + .select({ + connection: connections, + publishedRecord: publishedRecords, + }) + .from(connections) + .leftJoin( + publishedRecords, + eq(connections.publishedRecordId, publishedRecords.id), + ) + .where( + and( + eq(connections.sourceType, source.type), + eq(connections.sourceValue, source.stringValue), + eq(connections.targetType, target.type), + eq(connections.targetValue, target.stringValue), + ), + ) + .orderBy(connections.createdAt); + + const domainConnections: Connection[] = []; + for (const result of connectionResults) { + if (!result.connection) continue; + + const connectionDTO: ConnectionDTO = { + id: result.connection.id, + curatorId: result.connection.curatorId, + sourceType: result.connection.sourceType, + sourceValue: result.connection.sourceValue, + targetType: result.connection.targetType, + targetValue: result.connection.targetValue, + connectionType: result.connection.connectionType, + note: result.connection.note || undefined, + createdAt: result.connection.createdAt, + updatedAt: result.connection.updatedAt, + publishedRecordId: result.publishedRecord?.id || null, + publishedRecord: result.publishedRecord || undefined, + }; + + const domainResult = ConnectionMapper.toDomain(connectionDTO); + if (domainResult.isErr()) { + console.error( + 'Error mapping connection to domain:', + domainResult.error, + ); + continue; + } + domainConnections.push(domainResult.value); + } + + return ok(domainConnections); + } catch (error) { + return err(error as Error); + } + } + + async save(connection: Connection): Promise> { + try { + const { connection: connectionData, publishedRecord } = + ConnectionMapper.toPersistence(connection); + + await this.db.transaction(async (tx) => { + // Handle published record if it exists + let publishedRecordId: string | undefined = undefined; + + if (publishedRecord) { + const publishedRecordResult = await tx + .insert(publishedRecords) + .values({ + id: publishedRecord.id, + uri: publishedRecord.uri, + cid: publishedRecord.cid, + recordedAt: publishedRecord.recordedAt || new Date(), + }) + .onConflictDoNothing({ + target: [publishedRecords.uri, publishedRecords.cid], + }) + .returning({ id: publishedRecords.id }); + + if (publishedRecordResult.length === 0) { + const existingRecord = await tx + .select() + .from(publishedRecords) + .where( + and( + eq(publishedRecords.uri, publishedRecord.uri), + eq(publishedRecords.cid, publishedRecord.cid), + ), + ) + .limit(1); + + if (existingRecord.length > 0) { + publishedRecordId = existingRecord[0]!.id; + } + } else { + publishedRecordId = publishedRecordResult[0]!.id; + } + } + + // Upsert the connection + await tx + .insert(connections) + .values({ + ...connectionData, + publishedRecordId: publishedRecordId, + }) + .onConflictDoUpdate({ + target: connections.id, + set: { + curatorId: connectionData.curatorId, + sourceType: connectionData.sourceType, + sourceValue: connectionData.sourceValue, + targetType: connectionData.targetType, + targetValue: connectionData.targetValue, + connectionType: connectionData.connectionType, + note: connectionData.note, + updatedAt: connectionData.updatedAt, + publishedRecordId: publishedRecordId, + }, + }); + }); + + return ok(undefined); + } catch (error) { + return err(error as Error); + } + } + + async delete(connectionId: ConnectionId): Promise> { + try { + const id = connectionId.getStringValue(); + + await this.db.delete(connections).where(eq(connections.id, id)); + + return ok(undefined); + } catch (error) { + return err(error as Error); + } + } +} diff --git a/src/modules/cards/infrastructure/repositories/mappers/ConnectionMapper.ts b/src/modules/cards/infrastructure/repositories/mappers/ConnectionMapper.ts new file mode 100644 index 00000000..31fb4459 --- /dev/null +++ b/src/modules/cards/infrastructure/repositories/mappers/ConnectionMapper.ts @@ -0,0 +1,146 @@ +import { UniqueEntityID } from '../../../../../shared/domain/UniqueEntityID'; +import { Connection } from '../../../domain/Connection'; +import { ConnectionId } from '../../../domain/value-objects/ConnectionId'; +import { ConnectionType } from '../../../domain/value-objects/ConnectionType'; +import { + UrlOrCardId, + UrlOrCardIdType, +} from '../../../domain/value-objects/UrlOrCardId'; +import { ConnectionNote } from '../../../domain/value-objects/ConnectionNote'; +import { CuratorId } from '../../../domain/value-objects/CuratorId'; +import { PublishedRecordId } from '../../../domain/value-objects/PublishedRecordId'; +import { PublishedRecordDTO, PublishedRecordRefDTO } from './DTOTypes'; +import { err, ok, Result } from '../../../../../shared/core/Result'; + +// Database representation of a connection +export interface ConnectionDTO extends PublishedRecordRefDTO { + id: string; + curatorId: string; + sourceType: string; + sourceValue: string; + targetType: string; + targetValue: string; + connectionType?: string; + note?: string; + createdAt: Date; + updatedAt: Date; +} + +export class ConnectionMapper { + public static toDomain(dto: ConnectionDTO): Result { + try { + // Create curator ID + const curatorIdOrError = CuratorId.create(dto.curatorId); + if (curatorIdOrError.isErr()) return err(curatorIdOrError.error); + + // Create source + const sourceOrError = UrlOrCardId.reconstruct( + dto.sourceType as UrlOrCardIdType, + dto.sourceValue, + ); + if (sourceOrError.isErr()) return err(sourceOrError.error); + + // Create target + const targetOrError = UrlOrCardId.reconstruct( + dto.targetType as UrlOrCardIdType, + dto.targetValue, + ); + if (targetOrError.isErr()) return err(targetOrError.error); + + // Create optional connection type + let type: ConnectionType | undefined; + if (dto.connectionType) { + const typeOrError = ConnectionType.createFromString(dto.connectionType); + if (typeOrError.isErr()) return err(typeOrError.error); + type = typeOrError.value; + } + + // Create optional note + let note: ConnectionNote | undefined; + if (dto.note) { + const noteOrError = ConnectionNote.create(dto.note); + if (noteOrError.isErr()) return err(noteOrError.error); + note = noteOrError.value; + } + + // Create optional published record ID + let publishedRecordId: PublishedRecordId | undefined; + if (dto.publishedRecord) { + publishedRecordId = PublishedRecordId.create({ + uri: dto.publishedRecord.uri, + cid: dto.publishedRecord.cid, + }); + } + + // Create the connection + const connectionOrError = Connection.create( + { + source: sourceOrError.value, + target: targetOrError.value, + type, + note, + curatorId: curatorIdOrError.value, + publishedRecordId, + createdAt: dto.createdAt, + updatedAt: dto.updatedAt, + }, + new UniqueEntityID(dto.id), + ); + + if (connectionOrError.isErr()) return err(connectionOrError.error); + + return ok(connectionOrError.value); + } catch (error) { + return err(error as Error); + } + } + + public static toPersistence(connection: Connection): { + connection: { + id: string; + curatorId: string; + sourceType: string; + sourceValue: string; + targetType: string; + targetValue: string; + connectionType?: string; + note?: string; + createdAt: Date; + updatedAt: Date; + publishedRecordId?: string; + }; + publishedRecord?: PublishedRecordDTO; + } { + // Create published record data if it exists + let publishedRecord: PublishedRecordDTO | undefined; + let publishedRecordId: string | undefined; + + if (connection.publishedRecordId) { + const recordId = new UniqueEntityID().toString(); + publishedRecord = { + id: recordId, + uri: connection.publishedRecordId.uri, + cid: connection.publishedRecordId.cid, + recordedAt: new Date(), + }; + publishedRecordId = recordId; + } + + return { + connection: { + id: connection.connectionId.getStringValue(), + curatorId: connection.curatorId.value, + sourceType: connection.source.type, + sourceValue: connection.source.stringValue, + targetType: connection.target.type, + targetValue: connection.target.stringValue, + connectionType: connection.type?.value, + note: connection.note?.value, + createdAt: connection.createdAt, + updatedAt: connection.updatedAt, + publishedRecordId, + }, + publishedRecord, + }; + } +} diff --git a/src/modules/cards/infrastructure/repositories/schema/connection.sql.ts b/src/modules/cards/infrastructure/repositories/schema/connection.sql.ts new file mode 100644 index 00000000..1c0c42a6 --- /dev/null +++ b/src/modules/cards/infrastructure/repositories/schema/connection.sql.ts @@ -0,0 +1,51 @@ +import { pgTable, text, timestamp, uuid, index } from 'drizzle-orm/pg-core'; +import { publishedRecords } from './publishedRecord.sql'; +import { cards } from './card.sql'; + +export const connections = pgTable( + 'connections', + { + id: uuid('id').primaryKey(), + curatorId: text('curator_id').notNull(), + sourceType: text('source_type').notNull(), // 'URL' or 'CARD' + sourceValue: text('source_value').notNull(), // URL string or Card UUID + targetType: text('target_type').notNull(), // 'URL' or 'CARD' + targetValue: text('target_value').notNull(), // URL string or Card UUID + connectionType: text('connection_type'), // SUPPORTS, OPPOSES, etc. + note: text('note'), + publishedRecordId: uuid('published_record_id').references( + () => publishedRecords.id, + ), + createdAt: timestamp('created_at').notNull().defaultNow(), + updatedAt: timestamp('updated_at').notNull().defaultNow(), + }, + (table) => { + return { + // Critical for querying connections by curator + curatorIdIdx: index('connections_curator_id_idx').on(table.curatorId), + + // For finding connections from a source + sourceIdx: index('connections_source_idx').on( + table.sourceType, + table.sourceValue, + ), + + // For finding connections to a target + targetIdx: index('connections_target_idx').on( + table.targetType, + table.targetValue, + ), + + // For paginated listings + createdAtIdx: index('connections_created_at_idx').on( + table.createdAt.desc(), + ), + + // Composite index for curator's connections sorted by time + curatorCreatedAtIdx: index('connections_curator_created_at_idx').on( + table.curatorId, + table.createdAt.desc(), + ), + }; + }, +); diff --git a/src/modules/cards/tests/test-utils/createTestSchema.ts b/src/modules/cards/tests/test-utils/createTestSchema.ts index 9f489518..bc482293 100644 --- a/src/modules/cards/tests/test-utils/createTestSchema.ts +++ b/src/modules/cards/tests/test-utils/createTestSchema.ts @@ -120,6 +120,21 @@ export async function createTestSchema(db: PostgresJsDatabase) { PRIMARY KEY (follower_id, target_id, target_type) )`, + // Connections table (references published_records) + sql`CREATE TABLE IF NOT EXISTS connections ( + id UUID PRIMARY KEY DEFAULT uuid_generate_v4(), + curator_id TEXT NOT NULL, + source_type TEXT NOT NULL, + source_value TEXT NOT NULL, + target_type TEXT NOT NULL, + target_value TEXT NOT NULL, + connection_type TEXT, + note TEXT, + published_record_id UUID REFERENCES published_records(id), + created_at TIMESTAMP WITH TIME ZONE NOT NULL DEFAULT NOW(), + updated_at TIMESTAMP WITH TIME ZONE NOT NULL DEFAULT NOW() + )`, + // Following feed items table (references feed_activities) sql`CREATE TABLE IF NOT EXISTS following_feed_items ( user_id TEXT NOT NULL, @@ -266,6 +281,23 @@ export async function createTestSchema(db: PostgresJsDatabase) { CREATE INDEX IF NOT EXISTS idx_follows_target ON follows(target_id, target_type); `); + // Connections table indexes + await db.execute(sql` + CREATE INDEX IF NOT EXISTS connections_curator_id_idx ON connections(curator_id); + `); + await db.execute(sql` + CREATE INDEX IF NOT EXISTS connections_source_idx ON connections(source_type, source_value); + `); + await db.execute(sql` + CREATE INDEX IF NOT EXISTS connections_target_idx ON connections(target_type, target_value); + `); + await db.execute(sql` + CREATE INDEX IF NOT EXISTS connections_created_at_idx ON connections(created_at DESC); + `); + await db.execute(sql` + CREATE INDEX IF NOT EXISTS connections_curator_created_at_idx ON connections(curator_id, created_at DESC); + `); + // Following feed items indexes await db.execute(sql` CREATE INDEX IF NOT EXISTS idx_following_feed_user_time ON following_feed_items(user_id, created_at DESC); diff --git a/src/shared/infrastructure/events/EventConfig.ts b/src/shared/infrastructure/events/EventConfig.ts index 24370fcb..c3eda9d6 100644 --- a/src/shared/infrastructure/events/EventConfig.ts +++ b/src/shared/infrastructure/events/EventConfig.ts @@ -6,6 +6,8 @@ export const EventNames = { CARD_REMOVED_FROM_COLLECTION: 'CardRemovedFromCollectionEvent', USER_FOLLOWED_TARGET: 'USER_FOLLOWED_TARGET', USER_UNFOLLOWED_TARGET: 'USER_UNFOLLOWED_TARGET', + CONNECTION_CREATED: 'ConnectionCreatedEvent', + CONNECTION_REMOVED: 'ConnectionRemovedEvent', } as const; export type EventName = (typeof EventNames)[keyof typeof EventNames];