import crypto from "node:crypto"; import { eq, inArray } from "drizzle-orm"; import { db } from "../client.js"; import { circuitDeviceRows } from "../schema/circuit-device-rows.js"; import { circuits } from "../schema/circuits.js"; import { legacyConsumerCircuitMigrations } from "../schema/legacy-consumer-circuit-migrations.js"; import { legacyConsumerMigrationReports } from "../schema/legacy-consumer-migration-report.js"; import { toCircuitDeviceRowCreateValues, type CircuitDeviceRowCreateInput, } from "./circuit-device-row.repository.js"; import { toCircuitCreateValues, type CircuitCreatePersistenceInput, } from "./circuit.repository.js"; export interface LegacyMigrationCircuitPersistenceInput { circuit: CircuitCreatePersistenceInput; deviceRows: Array< Omit & { legacyConsumerId: string } >; } export interface LegacyMigrationReportPersistenceInput { legacyConsumerCount: number; createdCircuitCount: number; createdDeviceRowCount: number; duplicateGroupedCount: number; generatedIdentifierCount: number; unassignedRowCount: number; warningsJson: string; generatedIdentifiersJson: string; duplicateGroupsJson: string; } export class LegacyConsumerMigrationRepository { async listMigratedConsumerIds(circuitListId: string) { const rows = await db .select({ consumerId: legacyConsumerCircuitMigrations.consumerId }) .from(legacyConsumerCircuitMigrations) .where(eq(legacyConsumerCircuitMigrations.circuitListId, circuitListId)); return rows.map((row) => row.consumerId); } persistCircuitListMigration(input: { circuitListId: string; circuits: LegacyMigrationCircuitPersistenceInput[]; report: LegacyMigrationReportPersistenceInput; }) { if (input.circuits.some((entry) => entry.circuit.circuitListId !== input.circuitListId)) { throw new Error("All migrated circuits must belong to the target circuit list."); } const consumerIds = input.circuits.flatMap((entry) => entry.deviceRows.map((row) => row.legacyConsumerId) ); if (new Set(consumerIds).size !== consumerIds.length) { throw new Error("A legacy consumer may only be migrated once per operation."); } const preparedCircuits = input.circuits.map((entry) => { const circuitId = crypto.randomUUID(); return { circuitId, circuitValues: toCircuitCreateValues(circuitId, entry.circuit), rows: entry.deviceRows.map((row) => { const rowId = crypto.randomUUID(); return { consumerId: row.legacyConsumerId, rowId, rowValues: toCircuitDeviceRowCreateValues(rowId, { ...row, circuitId, }), }; }), }; }); const createdAtIso = new Date().toISOString(); db.transaction((tx) => { if (consumerIds.length > 0) { const existingMappings = tx .select({ consumerId: legacyConsumerCircuitMigrations.consumerId }) .from(legacyConsumerCircuitMigrations) .where(inArray(legacyConsumerCircuitMigrations.consumerId, consumerIds)) .all(); if (existingMappings.length > 0) { throw new Error("A legacy consumer was migrated before this operation completed."); } } for (const preparedCircuit of preparedCircuits) { tx.insert(circuits).values(preparedCircuit.circuitValues).run(); for (const row of preparedCircuit.rows) { tx.insert(circuitDeviceRows).values(row.rowValues).run(); tx .insert(legacyConsumerCircuitMigrations) .values({ consumerId: row.consumerId, circuitId: preparedCircuit.circuitId, circuitDeviceRowId: row.rowId, circuitListId: input.circuitListId, createdAtIso, }) .run(); } } const reportValues = { ...input.report, createdAtIso, }; const existingReport = tx .select({ id: legacyConsumerMigrationReports.id }) .from(legacyConsumerMigrationReports) .where(eq(legacyConsumerMigrationReports.circuitListId, input.circuitListId)) .limit(1) .all(); if (existingReport.length > 0) { tx .update(legacyConsumerMigrationReports) .set(reportValues) .where(eq(legacyConsumerMigrationReports.id, existingReport[0].id)) .run(); } else { tx .insert(legacyConsumerMigrationReports) .values({ id: crypto.randomUUID(), circuitListId: input.circuitListId, ...reportValues, }) .run(); } }); } }