#!/usr/bin/env -S deno run --allow-ffi --allow-net --allow-read --allow-write --allow-env /** * PANELS · Fase 5 · ETL único SQLite -> Postgres. * * Migra los datos de las dos SQLite históricas (data/app.db, data/platform.db) * a las 2 bases Postgres del monolito modular (panels_platform, * panels_product con esquemas iam/core). Corre en 3 fases: * * 1. Pre-flight: detecta de antemano lo que puede reventar el load o * corromper datos en silencio (duplicados CURP/RFC case-insensitive, * FKs huérfanas, tenant_id sin tenant, fechas con formato inválido). * 2. Carga: una transacción POR BASE/ESQUENA (platform, iam, core), * preservando los ids originales (OVERRIDING SYSTEM VALUE) para no * romper referencias cruzadas entre tablas. * 3. Verificación: compara conteos de filas origen/destino y sumas de * columnas de dinero con una tolerancia explícita (REAL -> NUMERIC * puede mover centavos). * * Uso (correr DESDE api/, con las credenciales _owner del ambiente destino * en el entorno -- DATABASE_URL_PLATFORM_OWNER, DATABASE_URL_IAM_OWNER, * DATABASE_URL_CORE_OWNER; el rol _app normal no alcanza a propósito): * cd api * set -a && source ../.env.dev-local && set +a # o el .env real del ambiente * deno run --allow-ffi --allow-net --allow-read --allow-write --allow-env \ * scripts/migrate-sqlite-to-postgres.ts \ * --app-db=../data/app.db --platform-db=../data/platform.db * * Sale con código != 0 si el pre-flight encuentra problemas o si la * verificación posterior no cuadra -- diseñado para un pipeline de * corte con criterios go/no-go, no para "correr y ya". */ import { Database as SqliteDatabase } from "jsr:@db/sqlite@0.12"; import { createPool } from "../pg.ts"; function arg(name: string, fallback: string): string { const prefix = `--${name}=`; const found = Deno.args.find((a) => a.startsWith(prefix)); return found ? found.slice(prefix.length) : fallback; } const APP_DB_PATH = arg("app-db", "../data/app.db"); const PLATFORM_DB_PATH = arg("platform-db", "../data/platform.db"); const DRY_RUN = Deno.args.includes("--dry-run"); const MONEY_TOLERANCE = 0.05; // pesos; ver nota de REAL -> NUMERIC en el plan let exitCode = 0; function fail(msg: string) { console.error(`[FAIL] ${msg}`); exitCode = 1; } function warn(msg: string) { console.warn(`[WARN] ${msg}`); } function ok(msg: string) { console.log(`[OK] ${msg}`); } // --------------------------------------------------------------------------- // Conexiones // --------------------------------------------------------------------------- const appDb = new SqliteDatabase(APP_DB_PATH, { readonly: true }); const platformDb = new SqliteDatabase(PLATFORM_DB_PATH, { readonly: true }); // El ETL necesita privilegios de owner (INSERT con OVERRIDING SYSTEM VALUE + // setval() de secuencias) -- el rol _app normal no alcanza a propósito // (least privilege). Cada rol _owner solo manda en SU esquema // (panels_iam_owner en iam, panels_core_owner en core, ver // db/provision/04-product-database.sql), así que hacen falta conexiones // separadas incluso dentro de panels_product -- no hay un "super owner" // que pueda escribir en ambos esquemas de una vez. function ownerUrl(name: string): string { const url = Deno.env.get(name) || ""; if (!url) throw new Error(`Falta ${name} -- el ETL requiere credenciales _owner, no _app`); return url; } const platformPg = createPool(ownerUrl("DATABASE_URL_PLATFORM_OWNER"), { max: 3 }); const iamPg = createPool(ownerUrl("DATABASE_URL_IAM_OWNER"), { max: 3 }); const corePg = createPool(ownerUrl("DATABASE_URL_CORE_OWNER"), { max: 3 }); // --------------------------------------------------------------------------- // Fase 1: Pre-flight // --------------------------------------------------------------------------- async function preflight(): Promise { console.log("\n=== Pre-flight ==="); // 1. Duplicados case-insensitive de CURP/RFC (SQLite los toleraba con // COLLATE NOCASE + índice único case-insensitive; CITEXT en Postgres // los rechazará igual, pero mejor detectarlo ANTES de la carga). const dupCurp = appDb.prepare( `SELECT LOWER(curp) AS c, COUNT(*) AS n FROM workers GROUP BY LOWER(curp) HAVING COUNT(*) > 1`, ).all() as { c: string; n: number }[]; if (dupCurp.length) fail(`${dupCurp.length} CURP duplicadas (case-insensitive): ${dupCurp.map((d) => d.c).join(", ")}`); else ok("Sin CURP duplicadas"); const dupRfc = appDb.prepare( `SELECT LOWER(rfc) AS c, COUNT(*) AS n FROM workers GROUP BY LOWER(rfc) HAVING COUNT(*) > 1`, ).all() as { c: string; n: number }[]; if (dupRfc.length) fail(`${dupRfc.length} RFC duplicados (case-insensitive): ${dupRfc.map((d) => d.c).join(", ")}`); else ok("Sin RFC duplicados"); // 2. FKs huérfanas (SQLite solo valida si PRAGMA foreign_keys estuvo ON // en cada escritura histórica -- puede haber huecos). const orphanChecks: { label: string; sql: string }[] = [ { label: "workers.company_id sin companies", sql: `SELECT COUNT(*) AS n FROM workers WHERE company_id IS NOT NULL AND company_id NOT IN (SELECT id FROM companies)` }, { label: "workers.risk_code sin risk_levels", sql: `SELECT COUNT(*) AS n FROM workers WHERE risk_code NOT IN (SELECT code FROM risk_levels)` }, { label: "documents.worker_id sin workers", sql: `SELECT COUNT(*) AS n FROM documents WHERE worker_id NOT IN (SELECT id FROM workers)` }, { label: "documents.uploaded_by sin users", sql: `SELECT COUNT(*) AS n FROM documents WHERE uploaded_by IS NOT NULL AND uploaded_by NOT IN (SELECT id FROM users)` }, { label: "assignments.worker_id sin workers", sql: `SELECT COUNT(*) AS n FROM assignments WHERE worker_id NOT IN (SELECT id FROM workers)` }, { label: "assignments.project_id sin projects", sql: `SELECT COUNT(*) AS n FROM assignments WHERE project_id NOT IN (SELECT id FROM projects)` }, { label: "projects.theme_id sin badge_themes", sql: `SELECT COUNT(*) AS n FROM projects WHERE theme_id NOT IN (SELECT id FROM badge_themes)` }, { label: "users.company_id sin companies", sql: `SELECT COUNT(*) AS n FROM users WHERE company_id IS NOT NULL AND company_id NOT IN (SELECT id FROM companies)` }, ]; for (const check of orphanChecks) { const row = appDb.prepare(check.sql).get() as { n: number }; if (row.n > 0) fail(`${row.n} filas: ${check.label}`); else ok(check.label.replace("sin", "OK ->")); } // 3. tenant_id nulo o sin tenant real. const tenantIds = new Set( (platformDb.prepare("SELECT id FROM tenants").all() as { id: number }[]).map((t) => t.id), ); for (const table of ["companies", "workers", "projects", "users"]) { const rows = appDb.prepare(`SELECT id, tenant_id FROM ${table}`).all() as { id: number; tenant_id: number | null }[]; const bad = rows.filter((r) => r.tenant_id == null || !tenantIds.has(r.tenant_id)); if (bad.length) fail(`${table}: ${bad.length} filas con tenant_id nulo o inexistente (ids: ${bad.slice(0, 10).map((r) => r.id).join(",")}${bad.length > 10 ? "..." : ""})`); else ok(`${table}.tenant_id: todas resuelven a un tenant real`); } // 4. Fechas con formato inválido en columnas que pasan a DATE. const dateCols: { table: string; col: string }[] = [ { table: "workers", col: "imss_alta_at" }, { table: "workers", col: "imss_baja_at" }, { table: "workers", col: "last_rehire_at" }, { table: "projects", col: "start_date" }, { table: "projects", col: "end_date" }, { table: "documents", col: "issued_at" }, { table: "documents", col: "expires_at" }, ]; for (const { table, col } of dateCols) { const rows = appDb.prepare( `SELECT id, ${col} AS v FROM ${table} WHERE ${col} IS NOT NULL AND ${col} != ''`, ).all() as { id: number; v: string }[]; const bad = rows.filter((r) => !/^\d{4}-\d{2}-\d{2}/.test(r.v)); if (bad.length) fail(`${table}.${col}: ${bad.length} fechas con formato inválido (ej. id=${bad[0].id} -> "${bad[0].v}")`); } if (!dateCols.some(() => false)) ok("Formato de fechas revisado"); } // --------------------------------------------------------------------------- // Fase 2: Carga // --------------------------------------------------------------------------- type Row = Record; function toBool(v: unknown): boolean { return v === 1 || v === true || v === "1"; } /** Copia una tabla completa preservando ids, con transform opcional por fila. * `identityColumn`: columna GENERATED ALWAYS AS IDENTITY a preservar via * OVERRIDING SYSTEM VALUE (null para catálogos con PK de texto/compuesta, * que no tienen identity y no la necesitan). `conflictColumns`: columna(s) * de conflicto para el ON CONFLICT ... DO NOTHING (idempotencia si se * corre el ETL dos veces). */ async function copyTable( sqlite: SqliteDatabase, pg: ReturnType, schema: string, table: string, opts: { sqliteTable?: string; transform?: (row: Row) => Row | null; columns?: string[]; identityColumn?: string | null; conflictColumns?: string[]; } = {}, ): Promise { const sqliteTable = opts.sqliteTable ?? table; const identityColumn = opts.identityColumn === undefined ? "id" : opts.identityColumn; const conflictColumns = opts.conflictColumns ?? (identityColumn ? [identityColumn] : []); const rows = sqlite.prepare(`SELECT * FROM ${sqliteTable}`).all() as Row[]; let inserted = 0; for (const raw of rows) { const row = opts.transform ? opts.transform(raw) : raw; if (!row) continue; // transform puede filtrar filas (ej. no migrables) const cols = opts.columns ?? Object.keys(row); const values = cols.map((c) => row[c]); const placeholders = cols.map((_, i) => `$${i + 1}`).join(", "); const overriding = identityColumn ? "OVERRIDING SYSTEM VALUE" : ""; const onConflict = conflictColumns.length ? `ON CONFLICT (${conflictColumns.join(", ")}) DO NOTHING` : ""; const text = `INSERT INTO ${schema}.${table} (${cols.join(", ")}) ${overriding} VALUES (${placeholders}) ${onConflict}`; if (!DRY_RUN) await pg.unsafe(text, values as never[]); inserted++; } console.log(` ${schema}.${table}: ${inserted}/${rows.length} filas`); return inserted; } async function setSequence(pg: ReturnType, schema: string, table: string): Promise { if (DRY_RUN) return; await pg.unsafe( `SELECT setval(pg_get_serial_sequence('${schema}.${table}', 'id'), COALESCE((SELECT MAX(id) FROM ${schema}.${table}), 1), (SELECT MAX(id) IS NOT NULL FROM ${schema}.${table}))`, ); } async function migratePlatform(): Promise { console.log("\n=== Carga: panels_platform ==="); await platformPg.begin(async (tx) => { await copyTable(platformDb, tx as never, "public", "tenants", { transform: (r) => ({ ...r, status: r.status === "inactivo" ? "suspendido" : r.status }), }); await copyTable(platformDb, tx as never, "public", "platform_users"); await copyTable(platformDb, tx as never, "public", "smtp_settings", { identityColumn: null, // id fijo = 1 (singleton), no es GENERATED conflictColumns: ["id"], // el baseline de Liquibase ya insertó la fila id=1 por defecto transform: (r) => ({ ...r, enabled: toBool(r.enabled) }), }); }); await setSequence(platformPg, "public", "tenants"); await setSequence(platformPg, "public", "platform_users"); } async function migrateIam(): Promise { console.log("\n=== Carga: panels_product.iam ==="); await iamPg.begin(async (tx) => { await copyTable(appDb, tx as never, "iam", "users", { columns: [ "id", "username", "password_hash", "display_name", "company_id", "tenant_id", "role_code", "must_change_password", "email", "created_at", ], transform: (r) => ({ ...r, role_code: r.role === "tenant_admin" ? "tenant_admin" : "user", must_change_password: toBool(r.must_change_password), }), }); }); await setSequence(iamPg, "iam", "users"); } /** users.id -> display_name, para el snapshot de uploaded_by/created_by * (Fase 4: documents/badge_jobs ya no tienen FK viva hacia iam). */ function userNameLookup(): Map { const rows = appDb.prepare("SELECT id, display_name FROM users").all() as { id: number; display_name: string }[]; return new Map(rows.map((r) => [r.id, r.display_name])); } async function migrateCore(): Promise { console.log("\n=== Carga: panels_product.core ==="); const names = userNameLookup(); await corePg.begin(async (tx) => { const t = tx as never as ReturnType; // Orden por dependencias de FK. risk_levels/badge_themes son catálogos // con PK de texto (code/id), no tienen columna identity que preservar. await copyTable(appDb, t, "core", "risk_levels", { identityColumn: null, conflictColumns: ["code"] }); await copyTable(appDb, t, "core", "badge_themes", { identityColumn: null, conflictColumns: ["id"] }); await copyTable(appDb, t, "core", "companies"); await copyTable(appDb, t, "core", "projects"); await copyTable(appDb, t, "core", "workers", { transform: (r) => ({ ...r, needs_badge: toBool(r.needs_badge) }), }); await copyTable(appDb, t, "core", "assignments", { transform: (r) => ({ ...r, active: toBool(r.active) }), }); await copyTable(appDb, t, "core", "document_types", { identityColumn: null, conflictColumns: ["code"], transform: (r) => ({ ...r, required: toBool(r.required), requires_issued_at: toBool(r.requires_issued_at), requires_expires_at: toBool(r.requires_expires_at), }), }); await copyTable(appDb, t, "core", "documents", { columns: [ "id", "worker_id", "type_code", "original_name", "mime", "size_bytes", "sha256", "iv", "storage_name", "is_current", "parse_status", "issued_at", "expires_at", "imss_company_id", "imss_alta_at", "uploaded_by_id", "uploaded_by_name", "uploaded_at", ], transform: (r) => ({ ...r, is_current: toBool(r.is_current), uploaded_by_id: r.uploaded_by ?? null, uploaded_by_name: r.uploaded_by ? names.get(Number(r.uploaded_by)) ?? "" : "", }), }); await copyTable(appDb, t, "core", "project_document_types", { identityColumn: null, conflictColumns: ["code"], transform: (r) => ({ ...r, required: toBool(r.required) }), }); await copyTable(appDb, t, "core", "project_documents", { columns: [ "id", "project_id", "type_code", "original_name", "mime", "size_bytes", "sha256", "iv", "storage_name", "is_current", "parse_status", "uploaded_by_id", "uploaded_by_name", "uploaded_at", ], transform: (r) => ({ ...r, is_current: toBool(r.is_current), uploaded_by_id: r.uploaded_by ?? null, uploaded_by_name: r.uploaded_by ? names.get(Number(r.uploaded_by)) ?? "" : "", }), }); await copyTable(appDb, t, "core", "company_document_types", { identityColumn: null, conflictColumns: ["code"], transform: (r) => ({ ...r, required: toBool(r.required) }), }); await copyTable(appDb, t, "core", "company_documents", { columns: [ "id", "company_id", "type_code", "original_name", "mime", "size_bytes", "sha256", "iv", "storage_name", "is_current", "parse_status", "uploaded_by_id", "uploaded_by_name", "uploaded_at", ], transform: (r) => ({ ...r, is_current: toBool(r.is_current), uploaded_by_id: r.uploaded_by ?? null, uploaded_by_name: r.uploaded_by ? names.get(Number(r.uploaded_by)) ?? "" : "", }), }); await copyTable(appDb, t, "core", "badge_jobs", { columns: ["id", "project_id", "status", "pdf_path", "created_by_id", "created_by_name", "created_at"], transform: (r) => ({ ...r, created_by_id: r.created_by ?? null, created_by_name: r.created_by ? names.get(Number(r.created_by)) ?? "" : "", }), }); await copyTable(appDb, t, "core", "badge_job_people", { identityColumn: null, conflictColumns: ["job_id", "worker_id"], transform: (r) => ({ ...r, delivered: toBool(r.delivered) }), }); await copyTable(appDb, t, "core", "loans"); await copyTable(appDb, t, "core", "attendance", { transform: (r) => ({ ...r, present: toBool(r.present) }), }); await copyTable(appDb, t, "core", "payroll_periods"); await copyTable(appDb, t, "core", "payroll_lines"); await copyTable(appDb, t, "core", "payroll_settings", { identityColumn: null, conflictColumns: ["tenant_id"], transform: (r) => ({ ...r, loan_commission_enabled: toBool(r.loan_commission_enabled) }), }); await copyTable(appDb, t, "core", "payroll_weeks"); await copyTable(appDb, t, "core", "payroll_sheets"); await copyTable(appDb, t, "core", "payroll_week_lines"); await copyTable(appDb, t, "core", "destajo_units"); await copyTable(appDb, t, "core", "destajo_periods"); await copyTable(appDb, t, "core", "destajo_jobs"); await copyTable(appDb, t, "core", "destajo_cut_lines"); await copyTable(appDb, t, "core", "loan_payments"); await copyTable(appDb, t, "core", "budget_chapters"); await copyTable(appDb, t, "core", "budget_items"); }); for ( const table of [ "companies", "projects", "workers", "assignments", "documents", "project_documents", "company_documents", "badge_jobs", "loans", "attendance", "payroll_periods", "payroll_lines", "payroll_weeks", "payroll_sheets", "payroll_week_lines", "destajo_units", "destajo_periods", "destajo_jobs", "destajo_cut_lines", "loan_payments", "budget_chapters", "budget_items", ] ) { await setSequence(corePg, "core", table); } } // --------------------------------------------------------------------------- // Fase 3: Verificación // --------------------------------------------------------------------------- async function verifyCounts(): Promise { console.log("\n=== Verificación: conteos de filas ==="); const checks: { sqlite: SqliteDatabase; table: string; pg: ReturnType; schema: string }[] = [ { sqlite: platformDb, table: "tenants", pg: platformPg, schema: "public" }, { sqlite: platformDb, table: "platform_users", pg: platformPg, schema: "public" }, { sqlite: appDb, table: "users", pg: iamPg, schema: "iam" }, { sqlite: appDb, table: "companies", pg: corePg, schema: "core" }, { sqlite: appDb, table: "workers", pg: corePg, schema: "core" }, { sqlite: appDb, table: "projects", pg: corePg, schema: "core" }, { sqlite: appDb, table: "documents", pg: corePg, schema: "core" }, { sqlite: appDb, table: "loans", pg: corePg, schema: "core" }, { sqlite: appDb, table: "budget_items", pg: corePg, schema: "core" }, ]; for (const c of checks) { const src = (c.sqlite.prepare(`SELECT COUNT(*) AS n FROM ${c.table}`).get() as { n: number }).n; const dstRows = await c.pg.unsafe(`SELECT COUNT(*)::int AS n FROM ${c.schema}.${c.table}`); const dst = (dstRows[0] as unknown as { n: number }).n; if (src !== dst) fail(`${c.schema}.${c.table}: origen=${src} destino=${dst}`); else ok(`${c.schema}.${c.table}: ${dst} filas en ambos lados`); } } async function verifyMoney(): Promise { console.log("\n=== Verificación: sumas de dinero (tolerancia $" + MONEY_TOLERANCE + ") ==="); const checks: { label: string; sqliteSql: string; pgSql: string }[] = [ { label: "workers.daily_wage", sqliteSql: "SELECT COALESCE(SUM(daily_wage),0) AS n FROM workers", pgSql: "SELECT COALESCE(SUM(daily_wage),0)::float AS n FROM core.workers" }, { label: "loans.balance", sqliteSql: "SELECT COALESCE(SUM(balance),0) AS n FROM loans", pgSql: "SELECT COALESCE(SUM(balance),0)::float AS n FROM core.loans" }, { label: "budget_items.amount", sqliteSql: "SELECT COALESCE(SUM(amount),0) AS n FROM budget_items", pgSql: "SELECT COALESCE(SUM(amount),0)::float AS n FROM core.budget_items" }, ]; for (const c of checks) { const src = (appDb.prepare(c.sqliteSql).get() as { n: number }).n; const dstRows = await corePg.unsafe(c.pgSql); const dst = (dstRows[0] as unknown as { n: number }).n; const diff = Math.abs(src - dst); if (diff > MONEY_TOLERANCE) fail(`${c.label}: origen=${src} destino=${dst} (diff=${diff.toFixed(4)} > tolerancia)`); else ok(`${c.label}: origen=${src} destino=${dst} (diff=${diff.toFixed(4)})`); } } // --------------------------------------------------------------------------- // main // --------------------------------------------------------------------------- try { await preflight(); if (exitCode !== 0) { console.error("\nPre-flight encontró problemas -- corrígelos antes de cargar. Abortando (no se tocó Postgres)."); Deno.exit(1); } if (DRY_RUN) { console.log("\n--dry-run: se detiene aquí (pre-flight OK, no se escribió nada)."); } else { await migratePlatform(); await migrateIam(); await migrateCore(); await verifyCounts(); await verifyMoney(); if (exitCode !== 0) { console.error("\nLa verificación posterior a la carga NO cuadra -- revisar antes de dar por buena la migración."); } else { console.log("\nMigración completa y verificada."); } } } finally { appDb.close(); platformDb.close(); await platformPg.end({ timeout: 5 }); await iamPg.end({ timeout: 5 }); await corePg.end({ timeout: 5 }); } Deno.exit(exitCode);