import { changedColumns, extractCustomFields, pick } from "../../../shared/repository"; import { markInvariantsSatisfied, type PurchaseOrder, type PurchaseOrderFieldChange, type PurchaseOrderHeader, type PurchaseOrderLine, type PurchaseOrderRevision, } from "../domain/purchaseOrder"; import type { Insertable, Selectable, Transaction, Updateable } from "../generated/kysely-tailordb"; export interface PurchaseOrderRepository { findById(id: string, opts?: { forUpdate?: boolean }): Promise; /** The distinct parent orders of the given lines, loaded in sorted id order. */ findByLineIds(lineIds: string[], opts?: { forUpdate?: boolean }): Promise; save(order: PurchaseOrder): Promise; } // ===== Columns ===== /** The domain-owned columns of each table. Everything else round-trips as custom fields. */ const PURCHASE_ORDER_HEADER_COLUMNS = [ "companyId", "supplierAccountId", "currencyId", "receivingSiteId", "orderStatus", "receiptStatus", "billingStatus", "orderDate", "externalSupplierOrderReference", "supplierSnapshotName", "rejectionReason", "closeReason", ] as const satisfies readonly Exclude[]; const PURCHASE_ORDER_LINE_COLUMNS = [ "id", "itemId", "itemSnapshotName", "itemSnapshotSku", "quantity", "unitPrice", "unitId", "receivingSiteId", "requiresPhysicalReceipt", "billedQuantity", "receivedQuantity", ] as const satisfies readonly Exclude[]; const PURCHASE_ORDER_FIELD_CHANGE_COLUMNS = [ "recordType", "recordId", "fieldName", "changeKind", "oldValue", "newValue", ] as const satisfies readonly (keyof PurchaseOrderFieldChange)[]; const HEADER_ROW_ONLY_COLUMNS = ["id", "createdAt", "updatedAt"] as const; const LINE_ROW_ONLY_COLUMNS = ["purchaseOrderId", "createdAt", "updatedAt"] as const; // ===== Mapping ===== export function toPurchaseOrder( headerRow: Selectable<"PurchaseOrder">, lineRows: Selectable<"PurchaseOrderLine">[], revisionRows: Selectable<"PurchaseOrderRevision">[] = [], fieldChangeRows: Selectable<"PurchaseOrderFieldChange">[] = [], ): PurchaseOrder { const header: PurchaseOrderHeader = { ...pick(headerRow, PURCHASE_ORDER_HEADER_COLUMNS), customFields: extractCustomFields(headerRow, [ ...PURCHASE_ORDER_HEADER_COLUMNS, ...HEADER_ROW_ONLY_COLUMNS, ]), }; const lines = lineRows.map((lineRow): PurchaseOrderLine => ({ ...pick(lineRow, PURCHASE_ORDER_LINE_COLUMNS), customFields: extractCustomFields(lineRow, [ ...PURCHASE_ORDER_LINE_COLUMNS, ...LINE_ROW_ONLY_COLUMNS, ]), })); const revisions = revisionRows.map((revisionRow): PurchaseOrderRevision => ({ id: revisionRow.id, revisionNumber: revisionRow.revisionNumber, reason: revisionRow.reason, amendedByUserId: revisionRow.amendedByUserId, fieldChanges: fieldChangeRows .filter((fieldChangeRow) => fieldChangeRow.revisionId === revisionRow.id) .map((fieldChangeRow) => pick(fieldChangeRow, PURCHASE_ORDER_FIELD_CHANGE_COLUMNS)), })); return markInvariantsSatisfied({ id: headerRow.id, header, lines, revisions }); } /** Custom fields first so domain-owned columns always win. */ function toPurchaseOrderRow(order: PurchaseOrder) { const row: Record = { ...order.header.customFields, id: order.id }; for (const column of PURCHASE_ORDER_HEADER_COLUMNS) { row[column] = order.header[column]; } return row; } /** Custom fields first so domain-owned columns always win. */ function toPurchaseOrderLineRow(line: PurchaseOrderLine, purchaseOrderId: string) { const row: Record = { ...line.customFields, purchaseOrderId }; for (const column of PURCHASE_ORDER_LINE_COLUMNS) { row[column] = line[column]; } return row; } // ===== Repository ===== export function createPurchaseOrderRepository(db: Transaction): PurchaseOrderRepository { return { async findById(id, opts) { let headerQuery = db.selectFrom("PurchaseOrder").selectAll().where("id", "=", id); if (opts?.forUpdate) { headerQuery = headerQuery.forUpdate(); } const headerRow = await headerQuery.executeTakeFirst(); if (!headerRow) { return null; } const lineRows = await db .selectFrom("PurchaseOrderLine") .selectAll() .where("purchaseOrderId", "=", id) .execute(); const revisionRows = await db .selectFrom("PurchaseOrderRevision") .selectAll() .where("purchaseOrderId", "=", id) .orderBy("revisionNumber") .execute(); const fieldChangeRows = revisionRows.length === 0 ? [] : await db .selectFrom("PurchaseOrderFieldChange") .selectAll() .where( "revisionId", "in", revisionRows.map((revisionRow) => revisionRow.id), ) .execute(); return toPurchaseOrder(headerRow, lineRows, revisionRows, fieldChangeRows); }, // purchaseOrderId is immutable, so resolving the parents before locking is // safe; locking the headers in sorted id order keeps concurrent syncs // deadlock-free. async findByLineIds(lineIds, opts) { if (lineIds.length === 0) { return []; } const lineRefRows = await db .selectFrom("PurchaseOrderLine") .select("purchaseOrderId") .where("id", "in", lineIds) .execute(); const orderIds = [...new Set(lineRefRows.map((row) => row.purchaseOrderId))]; if (orderIds.length === 0) { return []; } let headerQuery = db .selectFrom("PurchaseOrder") .selectAll() .where("id", "in", orderIds) .orderBy("id"); if (opts?.forUpdate) { headerQuery = headerQuery.forUpdate(); } const headerRows = await headerQuery.execute(); if (headerRows.length === 0) { return []; } const foundOrderIds = headerRows.map((headerRow) => headerRow.id); const lineRows = await db .selectFrom("PurchaseOrderLine") .selectAll() .where("purchaseOrderId", "in", foundOrderIds) .execute(); const revisionRows = await db .selectFrom("PurchaseOrderRevision") .selectAll() .where("purchaseOrderId", "in", foundOrderIds) .orderBy("revisionNumber") .execute(); const fieldChangeRows = revisionRows.length === 0 ? [] : await db .selectFrom("PurchaseOrderFieldChange") .selectAll() .where( "revisionId", "in", revisionRows.map((revisionRow) => revisionRow.id), ) .execute(); return headerRows.map((headerRow) => toPurchaseOrder( headerRow, lineRows.filter((lineRow) => lineRow.purchaseOrderId === headerRow.id), revisionRows.filter((revisionRow) => revisionRow.purchaseOrderId === headerRow.id), fieldChangeRows, ), ); }, // Diffing against the stored rows keeps a no-op save from bumping updatedAt. async save(order) { const storedHeader = await db .selectFrom("PurchaseOrder") .selectAll() .where("id", "=", order.id) .forUpdate() .executeTakeFirst(); let storedRevisionIds: ReadonlySet = new Set(); if (!storedHeader) { await db .insertInto("PurchaseOrder") .values(toPurchaseOrderRow(order) as Insertable<"PurchaseOrder">) .execute(); await db .insertInto("PurchaseOrderLine") .values( order.lines.map( (line) => toPurchaseOrderLineRow(line, order.id) as Insertable<"PurchaseOrderLine">, ), ) .execute(); } else { const storedLines = await db .selectFrom("PurchaseOrderLine") .selectAll() .where("purchaseOrderId", "=", order.id) .execute(); const storedLineById = new Map(storedLines.map((lineRow) => [lineRow.id, lineRow])); for (const line of order.lines) { const storedLine = storedLineById.get(line.id); if (!storedLine) { continue; } const lineChanges = changedColumns(storedLine, toPurchaseOrderLineRow(line, order.id)); if (Object.keys(lineChanges).length > 0) { await db .updateTable("PurchaseOrderLine") .set(lineChanges as Updateable<"PurchaseOrderLine">) .where("id", "=", line.id) .execute(); } } const keptLineIds = new Set(order.lines.map((line) => line.id)); const removedLineIds = storedLines .filter((lineRow) => !keptLineIds.has(lineRow.id)) .map((lineRow) => lineRow.id); if (removedLineIds.length > 0) { await db.deleteFrom("PurchaseOrderLine").where("id", "in", removedLineIds).execute(); } const addedRows = order.lines .filter((line) => !storedLineById.has(line.id)) .map((line) => toPurchaseOrderLineRow(line, order.id) as Insertable<"PurchaseOrderLine">); if (addedRows.length > 0) { await db.insertInto("PurchaseOrderLine").values(addedRows).execute(); } const headerChanges = changedColumns(storedHeader, toPurchaseOrderRow(order)); if (Object.keys(headerChanges).length > 0) { await db .updateTable("PurchaseOrder") .set(headerChanges as Updateable<"PurchaseOrder">) .where("id", "=", order.id) .execute(); } const storedRevisionRows = await db .selectFrom("PurchaseOrderRevision") .select("id") .where("purchaseOrderId", "=", order.id) .execute(); storedRevisionIds = new Set(storedRevisionRows.map((row) => row.id)); } // Revisions are immutable audit facts, so save only ever appends new ones. for (const revision of order.revisions) { if (storedRevisionIds.has(revision.id)) { continue; } await db .insertInto("PurchaseOrderRevision") .values({ id: revision.id, purchaseOrderId: order.id, revisionNumber: revision.revisionNumber, reason: revision.reason, amendedByUserId: revision.amendedByUserId, }) .execute(); if (revision.fieldChanges.length > 0) { await db .insertInto("PurchaseOrderFieldChange") .values( revision.fieldChanges.map((fieldChange) => ({ ...fieldChange, revisionId: revision.id, })), ) .execute(); } } }, }; }