import { markInvariantsSatisfied, type BusinessPartner, type BusinessPartnerHeader, type PartnerAddress, type PartnerBankAccount, } from "../domain/businessPartner"; import type { Insertable, Selectable, Transaction, Updateable } from "../generated/kysely-tailordb"; export interface BusinessPartnerRepository { findById(id: string, opts?: { forUpdate?: boolean }): Promise; findByAddressId( addressId: string, opts?: { forUpdate?: boolean }, ): Promise; findByBankAccountId( bankAccountId: string, opts?: { forUpdate?: boolean }, ): Promise; save(partner: BusinessPartner): Promise; } const HEADER_COLUMNS = ["name", "type", "status"] as const satisfies readonly Exclude< keyof BusinessPartnerHeader, "customFields" >[]; const ADDRESS_COLUMNS = [ "id", "active", "line1", "line2", "city", "state", "postalCode", "country", "createdAt", "updatedAt", ] as const satisfies readonly Exclude[]; const BANK_ACCOUNT_COLUMNS = [ "id", "active", "bankName", "accountHolderName", "accountNumber", "routingNumber", "currencyId", "createdAt", "updatedAt", ] as const satisfies readonly Exclude[]; const HEADER_ROW_ONLY_COLUMNS = ["id", "createdAt", "updatedAt"] as const; const CHILD_ROW_ONLY_COLUMNS = ["partnerId"] as const; function pick(row: T, keys: readonly K[]): Pick { const picked = {} as Pick; for (const key of keys) picked[key] = row[key]; return picked; } function extractCustomFields(row: object, knownColumns: readonly string[]) { const known = new Set(knownColumns); return Object.fromEntries(Object.entries(row).filter(([key]) => !known.has(key))); } function isSameValue(a: unknown, b: unknown) { if (a instanceof Date || b instanceof Date) { return a instanceof Date && b instanceof Date && a.getTime() === b.getTime(); } return a === b; } function changedColumns(before: Record, after: Record) { return Object.fromEntries( Object.entries(after).filter(([column, value]) => !isSameValue(before[column], value)), ); } function toPartnerAddress(row: Selectable<"PartnerAddress">): PartnerAddress { return { ...pick(row, ADDRESS_COLUMNS), customFields: extractCustomFields(row, [...ADDRESS_COLUMNS, ...CHILD_ROW_ONLY_COLUMNS]), }; } function toPartnerBankAccount(row: Selectable<"PartnerBankAccount">): PartnerBankAccount { return { ...pick(row, BANK_ACCOUNT_COLUMNS), customFields: extractCustomFields(row, [...BANK_ACCOUNT_COLUMNS, ...CHILD_ROW_ONLY_COLUMNS]), }; } export function toBusinessPartner( row: Selectable<"BusinessPartner">, addressRows: Selectable<"PartnerAddress">[] = [], bankAccountRows: Selectable<"PartnerBankAccount">[] = [], ): BusinessPartner { return markInvariantsSatisfied({ id: row.id, header: { ...pick(row, HEADER_COLUMNS), customFields: extractCustomFields(row, [...HEADER_COLUMNS, ...HEADER_ROW_ONLY_COLUMNS]), }, addresses: addressRows.map(toPartnerAddress), bankAccounts: bankAccountRows.map(toPartnerBankAccount), }); } function toHeaderRow(partner: BusinessPartner) { return { ...partner.header.customFields, id: partner.id, ...pick(partner.header, HEADER_COLUMNS), }; } function toAddressRow(address: PartnerAddress, partnerId: string) { return { ...address.customFields, partnerId, ...pick(address, ADDRESS_COLUMNS), }; } function toBankAccountRow(bankAccount: PartnerBankAccount, partnerId: string) { return { ...bankAccount.customFields, partnerId, ...pick(bankAccount, BANK_ACCOUNT_COLUMNS), }; } async function syncAddresses( db: Transaction, partner: BusinessPartner, storedRows: Selectable<"PartnerAddress">[], ) { const storedById = new Map(storedRows.map((row) => [row.id, row])); for (const address of partner.addresses) { const stored = storedById.get(address.id); if (!stored) continue; const changes = changedColumns(stored, toAddressRow(address, partner.id)); if (Object.keys(changes).length > 0) { await db .updateTable("PartnerAddress") .set(changes as Updateable<"PartnerAddress">) .where("id", "=", address.id) .execute(); } } const addedRows = partner.addresses .filter((address) => !storedById.has(address.id)) .map((address) => toAddressRow(address, partner.id) as Insertable<"PartnerAddress">); if (addedRows.length > 0) { await db.insertInto("PartnerAddress").values(addedRows).execute(); } } async function syncBankAccounts( db: Transaction, partner: BusinessPartner, storedRows: Selectable<"PartnerBankAccount">[], ) { const storedById = new Map(storedRows.map((row) => [row.id, row])); for (const bankAccount of partner.bankAccounts) { const stored = storedById.get(bankAccount.id); if (!stored) continue; const changes = changedColumns(stored, toBankAccountRow(bankAccount, partner.id)); if (Object.keys(changes).length > 0) { await db .updateTable("PartnerBankAccount") .set(changes as Updateable<"PartnerBankAccount">) .where("id", "=", bankAccount.id) .execute(); } } const addedRows = partner.bankAccounts .filter((account) => !storedById.has(account.id)) .map((account) => toBankAccountRow(account, partner.id) as Insertable<"PartnerBankAccount">); if (addedRows.length > 0) { await db.insertInto("PartnerBankAccount").values(addedRows).execute(); } } export function createBusinessPartnerRepository(db: Transaction): BusinessPartnerRepository { async function loadById( id: string, opts?: { forUpdate?: boolean }, ): Promise { let query = db.selectFrom("BusinessPartner").selectAll().where("id", "=", id); if (opts?.forUpdate) query = query.forUpdate(); const row = await query.executeTakeFirst(); if (!row) return null; const addressRows = await db .selectFrom("PartnerAddress") .selectAll() .where("partnerId", "=", id) .execute(); const bankAccountRows = await db .selectFrom("PartnerBankAccount") .selectAll() .where("partnerId", "=", id) .execute(); return toBusinessPartner(row, addressRows, bankAccountRows); } return { findById: loadById, async findByAddressId(addressId, opts) { const address = await db .selectFrom("PartnerAddress") .select("partnerId") .where("id", "=", addressId) .executeTakeFirst(); return address ? loadById(address.partnerId, opts) : null; }, async findByBankAccountId(bankAccountId, opts) { const bankAccount = await db .selectFrom("PartnerBankAccount") .select("partnerId") .where("id", "=", bankAccountId) .executeTakeFirst(); return bankAccount ? loadById(bankAccount.partnerId, opts) : null; }, async save(partner) { const storedHeader = await db .selectFrom("BusinessPartner") .selectAll() .where("id", "=", partner.id) .forUpdate() .executeTakeFirst(); const headerRow = toHeaderRow(partner); if (!storedHeader) { await db .insertInto("BusinessPartner") .values(headerRow as Insertable<"BusinessPartner">) .execute(); if (partner.addresses.length > 0) { await db .insertInto("PartnerAddress") .values( partner.addresses.map( (address) => toAddressRow(address, partner.id) as Insertable<"PartnerAddress">, ), ) .execute(); } if (partner.bankAccounts.length > 0) { await db .insertInto("PartnerBankAccount") .values( partner.bankAccounts.map( (account) => toBankAccountRow(account, partner.id) as Insertable<"PartnerBankAccount">, ), ) .execute(); } return; } const headerChanges = changedColumns(storedHeader, headerRow); if (Object.keys(headerChanges).length > 0) { await db .updateTable("BusinessPartner") .set(headerChanges as Updateable<"BusinessPartner">) .where("id", "=", partner.id) .execute(); } const storedAddresses = await db .selectFrom("PartnerAddress") .selectAll() .where("partnerId", "=", partner.id) .execute(); const storedBankAccounts = await db .selectFrom("PartnerBankAccount") .selectAll() .where("partnerId", "=", partner.id) .execute(); await syncAddresses(db, partner, storedAddresses); await syncBankAccounts(db, partner, storedBankAccounts); }, }; }