fix(cli): show live SQLite migration progress
Report source scans, per-table copy milestones, checksum phases, verification outcomes, and unambiguous failure or finalization status during first-boot and manual migrations.
This commit is contained in:
7
.changeset/bright-migrations-report.md
Normal file
7
.changeset/bright-migrations-report.md
Normal file
@@ -0,0 +1,7 @@
|
|||||||
|
---
|
||||||
|
"@runfusion/fusion": patch
|
||||||
|
---
|
||||||
|
|
||||||
|
summary: Show live phase, table, row-copy, verification, and failure progress during SQLite migration.
|
||||||
|
category: fix
|
||||||
|
dev: Adds structured migration progress events to first-boot startup and `fn db migrate` terminal output.
|
||||||
@@ -4,6 +4,7 @@ import {
|
|||||||
vacuumAnalyze,
|
vacuumAnalyze,
|
||||||
resolveBackend,
|
resolveBackend,
|
||||||
migrateSqliteToPostgres,
|
migrateSqliteToPostgres,
|
||||||
|
formatMigrationProgress,
|
||||||
completeSqliteMigration,
|
completeSqliteMigration,
|
||||||
defaultMigrationSources,
|
defaultMigrationSources,
|
||||||
stampMigratedProjectRows,
|
stampMigratedProjectRows,
|
||||||
@@ -278,6 +279,9 @@ export async function runDbMigrate(
|
|||||||
projectId: registeredProjectId,
|
projectId: registeredProjectId,
|
||||||
migrationKey: `project:${registeredProjectId}`,
|
migrationKey: `project:${registeredProjectId}`,
|
||||||
deferCompletion: true,
|
deferCompletion: true,
|
||||||
|
onProgress: (event) => {
|
||||||
|
console.log(`fn db migrate: ${formatMigrationProgress(event)}`);
|
||||||
|
},
|
||||||
});
|
});
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
console.error(`fn db migrate: migration failed: ${(error as Error).message}`);
|
console.error(`fn db migrate: migration failed: ${(error as Error).message}`);
|
||||||
|
|||||||
@@ -34,8 +34,10 @@ import { tmpdir } from "node:os";
|
|||||||
import { join } from "node:path";
|
import { join } from "node:path";
|
||||||
import { DatabaseSync } from "../../sqlite-adapter.js";
|
import { DatabaseSync } from "../../sqlite-adapter.js";
|
||||||
import {
|
import {
|
||||||
|
formatMigrationProgress,
|
||||||
migrateSqliteToPostgres,
|
migrateSqliteToPostgres,
|
||||||
toSnakeCase,
|
toSnakeCase,
|
||||||
|
type MigrationProgressEvent,
|
||||||
} from "../../postgres/sqlite-migrator.js";
|
} from "../../postgres/sqlite-migrator.js";
|
||||||
import { applySchemaBaseline } from "../../postgres/schema-applier.js";
|
import { applySchemaBaseline } from "../../postgres/schema-applier.js";
|
||||||
|
|
||||||
@@ -46,6 +48,35 @@ const PG_AVAILABLE =
|
|||||||
|
|
||||||
const pgDescribe = PG_AVAILABLE ? describe : describe.skip;
|
const pgDescribe = PG_AVAILABLE ? describe : describe.skip;
|
||||||
|
|
||||||
|
describe("SQLite migration CLI progress", () => {
|
||||||
|
it("formats table copy progress with position, row counts, and percentage", () => {
|
||||||
|
expect(formatMigrationProgress({
|
||||||
|
phase: "table-progress",
|
||||||
|
sourceSchema: "project",
|
||||||
|
table: "run_audit_events",
|
||||||
|
tableIndex: 42,
|
||||||
|
tableCount: 124,
|
||||||
|
processedRows: 63_400,
|
||||||
|
sourceRows: 252_947,
|
||||||
|
})).toBe("[42/124] project.run_audit_events: processed 63,400/252,947 rows (25%)");
|
||||||
|
});
|
||||||
|
|
||||||
|
it("makes transaction rollback explicit on failure", () => {
|
||||||
|
expect(formatMigrationProgress({
|
||||||
|
phase: "failed",
|
||||||
|
error: "project.tasks failed verification",
|
||||||
|
})).toBe("FAILED — migration transaction rolled back: project.tasks failed verification");
|
||||||
|
});
|
||||||
|
|
||||||
|
it("distinguishes committed verification failures from transaction rollbacks", () => {
|
||||||
|
expect(formatMigrationProgress({
|
||||||
|
phase: "failed",
|
||||||
|
tableCount: 124,
|
||||||
|
failedTables: 2,
|
||||||
|
})).toBe("FAILED — 2/124 tables failed verification; migration will not be marked complete.");
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* FNXC:PostgresMigration 2026-06-24-09:05:
|
* FNXC:PostgresMigration 2026-06-24-09:05:
|
||||||
* Create a uniquely-named fresh PostgreSQL database. Mirrors the
|
* Create a uniquely-named fresh PostgreSQL database. Mirrors the
|
||||||
@@ -448,6 +479,90 @@ pgDescribe("SQLite-to-PostgreSQL migrator", () => {
|
|||||||
expect(archived.targetRows).toBe(1);
|
expect(archived.targetRows).toBe(1);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:CliMigrationProgress 2026-07-14-13:47:
|
||||||
|
CLI-triggered first-boot migrations must expose schema preparation, source discovery, per-table copy/verification, and terminal success so operators can distinguish a long-running copy from a stalled or failed startup.
|
||||||
|
*/
|
||||||
|
it("reports structured progress through planning, table verification, and completion", async () => {
|
||||||
|
const progress: MigrationProgressEvent[] = [];
|
||||||
|
let rejectFirstCallback = true;
|
||||||
|
const report = await migrateTest(
|
||||||
|
ctx!.db,
|
||||||
|
[{ sqlitePath: join(ctx!.fusionDir, "fusion.db"), pgSchema: "project" as const }],
|
||||||
|
{
|
||||||
|
onProgress: (event) => {
|
||||||
|
progress.push(event);
|
||||||
|
if (rejectFirstCallback) {
|
||||||
|
rejectFirstCallback = false;
|
||||||
|
return Promise.reject(new Error("test progress sink failure"));
|
||||||
|
}
|
||||||
|
},
|
||||||
|
},
|
||||||
|
);
|
||||||
|
|
||||||
|
expect(progress[0]?.phase).toBe("preparing-schema");
|
||||||
|
expect(progress.some((event) => event.phase === "scanning-source")).toBe(true);
|
||||||
|
expect(progress.some((event) => event.phase === "copy-started" && event.tableCount === report.tables.length)).toBe(true);
|
||||||
|
expect(progress.some((event) => event.phase === "table-started" && event.table === "tasks")).toBe(true);
|
||||||
|
expect(progress.some((event) => event.phase === "table-verifying" && event.table === "tasks")).toBe(true);
|
||||||
|
expect(progress.some((event) => event.phase === "table-complete" && event.table === "tasks")).toBe(true);
|
||||||
|
expect(progress.at(-1)?.phase).toBe("copy-complete");
|
||||||
|
expect(progress.filter((event) => event.phase === "table-complete")).toHaveLength(report.tables.length);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("reports bounded quarter progress for a multi-batch table without a redundant 100% event", async () => {
|
||||||
|
const sqlitePath = join(ctx!.fusionDir, "fusion.db");
|
||||||
|
const legacy = new DatabaseSync(sqlitePath);
|
||||||
|
try {
|
||||||
|
const insert = legacy.prepare(
|
||||||
|
`INSERT INTO tasks (id, title, description, "column", createdAt, updatedAt)
|
||||||
|
VALUES (?, ?, ?, ?, ?, ?)`,
|
||||||
|
);
|
||||||
|
for (let index = 0; index < 648; index += 1) {
|
||||||
|
insert.run(
|
||||||
|
`FN-PROGRESS-${index}`,
|
||||||
|
`Progress ${index}`,
|
||||||
|
"progress fixture",
|
||||||
|
"todo",
|
||||||
|
"2026-07-14T00:00:00Z",
|
||||||
|
"2026-07-14T00:00:00Z",
|
||||||
|
);
|
||||||
|
}
|
||||||
|
} finally {
|
||||||
|
legacy.close();
|
||||||
|
}
|
||||||
|
|
||||||
|
const progress: MigrationProgressEvent[] = [];
|
||||||
|
await migrateTest(
|
||||||
|
ctx!.db,
|
||||||
|
[{ sqlitePath, pgSchema: "project" as const }],
|
||||||
|
{ onProgress: (event) => { progress.push(event); } },
|
||||||
|
);
|
||||||
|
const taskProgress = progress.filter(
|
||||||
|
(event): event is Extract<MigrationProgressEvent, { phase: "table-progress" }> =>
|
||||||
|
event.phase === "table-progress" && event.table === "tasks",
|
||||||
|
);
|
||||||
|
|
||||||
|
expect(taskProgress.map((event) => Math.floor((event.processedRows / event.sourceRows) * 4))).toEqual([1, 2, 3]);
|
||||||
|
expect(taskProgress.every((event) => event.processedRows < event.sourceRows)).toBe(true);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("reports a terminal dry-run plan without implying that data was written", async () => {
|
||||||
|
const progress: MigrationProgressEvent[] = [];
|
||||||
|
const report = await migrateTest(
|
||||||
|
ctx!.db,
|
||||||
|
[{ sqlitePath: join(ctx!.fusionDir, "fusion.db"), pgSchema: "project" as const }],
|
||||||
|
{ dryRun: true, onProgress: (event) => { progress.push(event); } },
|
||||||
|
);
|
||||||
|
|
||||||
|
expect(report.dryRun).toBe(true);
|
||||||
|
expect(progress.at(-1)).toMatchObject({
|
||||||
|
phase: "dry-run-complete",
|
||||||
|
tableCount: report.tables.length,
|
||||||
|
});
|
||||||
|
expect(formatMigrationProgress(progress.at(-1)!)).toContain("no data written");
|
||||||
|
});
|
||||||
|
|
||||||
/*
|
/*
|
||||||
FNXC:PostgresMigration 2026-07-13-22:37:
|
FNXC:PostgresMigration 2026-07-13-22:37:
|
||||||
Every user table in a legacy SQLite database must be represented in the migration report. An unknown table is retained as an explicit failed verification so startup cannot claim a complete cutover while silently abandoning operator data.
|
Every user table in a legacy SQLite database must be represented in the migration report. An unknown table is retained as an explicit failed verification so startup cannot claim a complete cutover while silently abandoning operator data.
|
||||||
@@ -462,9 +577,12 @@ pgDescribe("SQLite-to-PostgreSQL migrator", () => {
|
|||||||
legacy.close();
|
legacy.close();
|
||||||
}
|
}
|
||||||
|
|
||||||
const report = await migrateTest(ctx!.db, [
|
const progress: MigrationProgressEvent[] = [];
|
||||||
{ sqlitePath, pgSchema: "project" as const },
|
const report = await migrateTest(
|
||||||
]);
|
ctx!.db,
|
||||||
|
[{ sqlitePath, pgSchema: "project" as const }],
|
||||||
|
{ onProgress: (event) => { progress.push(event); } },
|
||||||
|
);
|
||||||
|
|
||||||
expect(report.tables).toContainEqual(
|
expect(report.tables).toContainEqual(
|
||||||
expect.objectContaining({
|
expect.objectContaining({
|
||||||
@@ -476,6 +594,10 @@ pgDescribe("SQLite-to-PostgreSQL migrator", () => {
|
|||||||
skipped: false,
|
skipped: false,
|
||||||
}),
|
}),
|
||||||
);
|
);
|
||||||
|
expect(progress.at(-1)).toMatchObject({
|
||||||
|
phase: "failed",
|
||||||
|
failedTables: 1,
|
||||||
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
/*
|
/*
|
||||||
|
|||||||
@@ -2239,6 +2239,7 @@ export {
|
|||||||
isSqliteMigrationComplete,
|
isSqliteMigrationComplete,
|
||||||
completeSqliteMigration,
|
completeSqliteMigration,
|
||||||
defaultMigrationSources,
|
defaultMigrationSources,
|
||||||
|
formatMigrationProgress,
|
||||||
// FNXC:CentralProjectIdentity 2026-07-13-23:10:
|
// FNXC:CentralProjectIdentity 2026-07-13-23:10:
|
||||||
// Post-migration project-partition stamping, shared by the startup-factory
|
// Post-migration project-partition stamping, shared by the startup-factory
|
||||||
// first-boot auto-migration and `fn db migrate` so migrated rows are re-keyed
|
// first-boot auto-migration and `fn db migrate` so migrated rows are re-keyed
|
||||||
@@ -2281,6 +2282,8 @@ export type {
|
|||||||
SqliteMigrationSource,
|
SqliteMigrationSource,
|
||||||
SchemaName,
|
SchemaName,
|
||||||
MigrationReport,
|
MigrationReport,
|
||||||
|
MigrationProgressEvent,
|
||||||
|
MigrationProgressPhase,
|
||||||
TableMigrationResult,
|
TableMigrationResult,
|
||||||
StampMigratedProjectRowsInput,
|
StampMigratedProjectRowsInput,
|
||||||
StampMigratedProjectRowsResult,
|
StampMigratedProjectRowsResult,
|
||||||
|
|||||||
@@ -159,11 +159,14 @@ export {
|
|||||||
isSqliteMigrationComplete,
|
isSqliteMigrationComplete,
|
||||||
completeSqliteMigration,
|
completeSqliteMigration,
|
||||||
defaultMigrationSources,
|
defaultMigrationSources,
|
||||||
|
formatMigrationProgress,
|
||||||
toSnakeCase,
|
toSnakeCase,
|
||||||
type SqliteMigrationSource,
|
type SqliteMigrationSource,
|
||||||
type SchemaName,
|
type SchemaName,
|
||||||
type MigrationOptions,
|
type MigrationOptions,
|
||||||
type MigrationReport,
|
type MigrationReport,
|
||||||
|
type MigrationProgressEvent,
|
||||||
|
type MigrationProgressPhase,
|
||||||
type TableMigrationResult,
|
type TableMigrationResult,
|
||||||
} from "./sqlite-migrator.js";
|
} from "./sqlite-migrator.js";
|
||||||
|
|
||||||
|
|||||||
@@ -56,6 +56,7 @@ import { DatabaseSync } from "../sqlite-adapter.js";
|
|||||||
import type { PostgresJsDatabase } from "drizzle-orm/postgres-js";
|
import type { PostgresJsDatabase } from "drizzle-orm/postgres-js";
|
||||||
import { sql } from "drizzle-orm";
|
import { sql } from "drizzle-orm";
|
||||||
import { createHash } from "node:crypto";
|
import { createHash } from "node:crypto";
|
||||||
|
import { basename } from "node:path";
|
||||||
import { applySchemaBaseline } from "./schema-applier.js";
|
import { applySchemaBaseline } from "./schema-applier.js";
|
||||||
import {
|
import {
|
||||||
PROJECT_SCHEMA,
|
PROJECT_SCHEMA,
|
||||||
@@ -63,6 +64,7 @@ import {
|
|||||||
ARCHIVE_SCHEMA,
|
ARCHIVE_SCHEMA,
|
||||||
} from "./schema/_shared.js";
|
} from "./schema/_shared.js";
|
||||||
import { createLogger } from "../logger.js";
|
import { createLogger } from "../logger.js";
|
||||||
|
import { getErrorMessage } from "../error-message.js";
|
||||||
|
|
||||||
const log = createLogger("sqlite-migrator");
|
const log = createLogger("sqlite-migrator");
|
||||||
|
|
||||||
@@ -176,6 +178,110 @@ export interface MigrationReport {
|
|||||||
readonly appliedBaseline: boolean;
|
readonly appliedBaseline: boolean;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
interface TableProgressCoordinates {
|
||||||
|
readonly sourceSchema: SchemaName;
|
||||||
|
readonly table: string;
|
||||||
|
readonly tableIndex: number;
|
||||||
|
readonly tableCount: number;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Structured status emitted to CLI callers during a potentially long migration. */
|
||||||
|
export type MigrationProgressEvent =
|
||||||
|
| { readonly phase: "preparing-schema" }
|
||||||
|
| {
|
||||||
|
readonly phase: "scanning-source";
|
||||||
|
readonly sourcePath: string;
|
||||||
|
readonly sourceSchema: SchemaName;
|
||||||
|
readonly sourceIndex: number;
|
||||||
|
readonly sourceCount: number;
|
||||||
|
}
|
||||||
|
| { readonly phase: "copy-started"; readonly tableCount: number }
|
||||||
|
| ({ readonly phase: "table-started" } & TableProgressCoordinates)
|
||||||
|
| ({
|
||||||
|
readonly phase: "table-progress";
|
||||||
|
readonly processedRows: number;
|
||||||
|
readonly sourceRows: number;
|
||||||
|
} & TableProgressCoordinates)
|
||||||
|
| ({
|
||||||
|
readonly phase: "table-verifying";
|
||||||
|
readonly sourceRows: number;
|
||||||
|
readonly verificationStage: "target-count" | "source-content" | "target-content";
|
||||||
|
} & TableProgressCoordinates)
|
||||||
|
| ({
|
||||||
|
readonly phase: "table-complete";
|
||||||
|
readonly sourceRows: number;
|
||||||
|
readonly insertedRows: number;
|
||||||
|
readonly targetRows: number;
|
||||||
|
readonly verified: boolean;
|
||||||
|
readonly skipped: boolean;
|
||||||
|
readonly skipReason?: string;
|
||||||
|
} & TableProgressCoordinates)
|
||||||
|
| {
|
||||||
|
readonly phase: "copy-complete";
|
||||||
|
readonly tableCount: number;
|
||||||
|
readonly verifiedTables: number;
|
||||||
|
readonly failedTables: number;
|
||||||
|
readonly sequenceBumps: number;
|
||||||
|
}
|
||||||
|
| {
|
||||||
|
readonly phase: "dry-run-complete";
|
||||||
|
readonly tableCount: number;
|
||||||
|
readonly sourceRows: number;
|
||||||
|
}
|
||||||
|
| { readonly phase: "failed"; readonly error: string }
|
||||||
|
| {
|
||||||
|
readonly phase: "failed";
|
||||||
|
readonly tableCount: number;
|
||||||
|
readonly verifiedTables: number;
|
||||||
|
readonly failedTables: number;
|
||||||
|
};
|
||||||
|
|
||||||
|
export type MigrationProgressPhase = MigrationProgressEvent["phase"];
|
||||||
|
|
||||||
|
function formatTableProgressPrefix(event: TableProgressCoordinates): string {
|
||||||
|
return `[${event.tableIndex}/${event.tableCount}] ${event.sourceSchema}.${event.table}`;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Format a structured migration event for terminal output. */
|
||||||
|
export function formatMigrationProgress(event: MigrationProgressEvent): string {
|
||||||
|
switch (event.phase) {
|
||||||
|
case "preparing-schema":
|
||||||
|
return "Preparing PostgreSQL schema…";
|
||||||
|
case "scanning-source":
|
||||||
|
return `Scanning source ${event.sourceIndex}/${event.sourceCount}: ${basename(event.sourcePath)} → ${event.sourceSchema}…`;
|
||||||
|
case "copy-started":
|
||||||
|
return `Found ${event.tableCount} tables. Copying and verifying data…`;
|
||||||
|
case "table-started":
|
||||||
|
return `${formatTableProgressPrefix(event)}: starting…`;
|
||||||
|
case "table-progress": {
|
||||||
|
const percent = Math.floor((event.processedRows / event.sourceRows) * 100);
|
||||||
|
return `${formatTableProgressPrefix(event)}: processed ${event.processedRows.toLocaleString()}/${event.sourceRows.toLocaleString()} rows (${percent}%)`;
|
||||||
|
}
|
||||||
|
case "table-verifying": {
|
||||||
|
const stage = {
|
||||||
|
"target-count": "counting migrated rows",
|
||||||
|
"source-content": "checksumming SQLite source",
|
||||||
|
"target-content": "checksumming PostgreSQL target",
|
||||||
|
}[event.verificationStage];
|
||||||
|
return `${formatTableProgressPrefix(event)}: processed ${event.sourceRows.toLocaleString()} rows; ${stage}…`;
|
||||||
|
}
|
||||||
|
case "table-complete":
|
||||||
|
if (event.skipped) return `${formatTableProgressPrefix(event)}: skipped — ${event.skipReason ?? "not required"}`;
|
||||||
|
if (!event.verified) {
|
||||||
|
return `${formatTableProgressPrefix(event)}: VERIFICATION FAILED — source=${event.sourceRows}, target=${event.targetRows}${event.skipReason ? ` (${event.skipReason})` : ""}`;
|
||||||
|
}
|
||||||
|
return `${formatTableProgressPrefix(event)}: verified — ${event.sourceRows.toLocaleString()} rows (${event.insertedRows.toLocaleString()} inserted)`;
|
||||||
|
case "copy-complete":
|
||||||
|
return `Copy and verification complete — ${event.verifiedTables}/${event.tableCount} tables verified; ${event.sequenceBumps} sequences updated. Finalizing migration…`;
|
||||||
|
case "dry-run-complete":
|
||||||
|
return `Dry run complete — ${event.tableCount} tables and ${event.sourceRows.toLocaleString()} source rows planned; no data written.`;
|
||||||
|
case "failed":
|
||||||
|
return "error" in event
|
||||||
|
? `FAILED — migration transaction rolled back: ${event.error}`
|
||||||
|
: `FAILED — ${event.failedTables}/${event.tableCount} tables failed verification; migration will not be marked complete.`;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/** Options for the migration. */
|
/** Options for the migration. */
|
||||||
export interface MigrationOptions {
|
export interface MigrationOptions {
|
||||||
/** If true, report the planned copy without modifying PostgreSQL. */
|
/** If true, report the planned copy without modifying PostgreSQL. */
|
||||||
@@ -192,6 +298,21 @@ export interface MigrationOptions {
|
|||||||
readonly migrationKey?: string;
|
readonly migrationKey?: string;
|
||||||
/** Leave a verified migration running until caller-side project stamping succeeds. */
|
/** Leave a verified migration running until caller-side project stamping succeeds. */
|
||||||
readonly deferCompletion?: boolean;
|
readonly deferCompletion?: boolean;
|
||||||
|
/** Receives CLI-safe progress events; callback failures never abort migration. */
|
||||||
|
readonly onProgress?: (event: MigrationProgressEvent) => void | Promise<void>;
|
||||||
|
}
|
||||||
|
|
||||||
|
function emitMigrationProgress(options: MigrationOptions, event: MigrationProgressEvent): void {
|
||||||
|
try {
|
||||||
|
const pending = options.onProgress?.(event);
|
||||||
|
if (pending) {
|
||||||
|
void pending.catch((error) => {
|
||||||
|
log.warn(`Migration progress callback failed: ${getErrorMessage(error)}`);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
} catch (error) {
|
||||||
|
log.warn(`Migration progress callback failed: ${getErrorMessage(error)}`);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
const SQLITE_MIGRATION_STATE_TABLE = "fusion_sqlite_migrations";
|
const SQLITE_MIGRATION_STATE_TABLE = "fusion_sqlite_migrations";
|
||||||
@@ -281,6 +402,11 @@ export async function migrateSqliteToPostgres(
|
|||||||
),
|
),
|
||||||
);
|
);
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
|
const errorMessage = getErrorMessage(error);
|
||||||
|
emitMigrationProgress(options, {
|
||||||
|
phase: "failed",
|
||||||
|
error: errorMessage,
|
||||||
|
});
|
||||||
if (options.dryRun !== true) {
|
if (options.dryRun !== true) {
|
||||||
await migrationDb.transaction(async (tx) => {
|
await migrationDb.transaction(async (tx) => {
|
||||||
await tx.execute(sql`SELECT pg_advisory_xact_lock(hashtext('fusion:sqlite-migration-state'))`);
|
await tx.execute(sql`SELECT pg_advisory_xact_lock(hashtext('fusion:sqlite-migration-state'))`);
|
||||||
@@ -288,7 +414,7 @@ export async function migrateSqliteToPostgres(
|
|||||||
await tx.execute(sql`
|
await tx.execute(sql`
|
||||||
INSERT INTO public.${sql.identifier(SQLITE_MIGRATION_STATE_TABLE)}
|
INSERT INTO public.${sql.identifier(SQLITE_MIGRATION_STATE_TABLE)}
|
||||||
(migration_key, project_id, status, last_error, updated_at)
|
(migration_key, project_id, status, last_error, updated_at)
|
||||||
VALUES (${migrationKey}, ${options.projectId ?? null}, 'failed', ${error instanceof Error ? error.message : String(error)}, now())
|
VALUES (${migrationKey}, ${options.projectId ?? null}, 'failed', ${errorMessage}, now())
|
||||||
ON CONFLICT (migration_key) DO UPDATE
|
ON CONFLICT (migration_key) DO UPDATE
|
||||||
SET project_id = EXCLUDED.project_id, status = 'failed', last_error = EXCLUDED.last_error, updated_at = now()
|
SET project_id = EXCLUDED.project_id, status = 'failed', last_error = EXCLUDED.last_error, updated_at = now()
|
||||||
`);
|
`);
|
||||||
@@ -305,6 +431,7 @@ async function migrateSqliteToPostgresOnSession(
|
|||||||
): Promise<MigrationReport> {
|
): Promise<MigrationReport> {
|
||||||
const dryRun = options.dryRun === true;
|
const dryRun = options.dryRun === true;
|
||||||
const migrationKey = options.migrationKey ?? `project:${options.projectId ?? "unbound"}`;
|
const migrationKey = options.migrationKey ?? `project:${options.projectId ?? "unbound"}`;
|
||||||
|
emitMigrationProgress(options, { phase: "preparing-schema" });
|
||||||
|
|
||||||
/*
|
/*
|
||||||
* FNXC:PostgresMigration 2026-07-14-00:05:
|
* FNXC:PostgresMigration 2026-07-14-00:05:
|
||||||
@@ -345,7 +472,7 @@ async function migrateSqliteToPostgresOnSession(
|
|||||||
if (!dryRun) {
|
if (!dryRun) {
|
||||||
await migrationDb.execute(sql`
|
await migrationDb.execute(sql`
|
||||||
UPDATE public.${sql.identifier(SQLITE_MIGRATION_STATE_TABLE)}
|
UPDATE public.${sql.identifier(SQLITE_MIGRATION_STATE_TABLE)}
|
||||||
SET status = 'failed', last_error = ${error instanceof Error ? error.message : String(error)}, updated_at = now()
|
SET status = 'failed', last_error = ${getErrorMessage(error)}, updated_at = now()
|
||||||
WHERE migration_key = ${migrationKey}
|
WHERE migration_key = ${migrationKey}
|
||||||
`);
|
`);
|
||||||
}
|
}
|
||||||
@@ -375,7 +502,7 @@ async function migrateSqliteToPostgresOnSession(
|
|||||||
} catch (error) {
|
} catch (error) {
|
||||||
log.warn(
|
log.warn(
|
||||||
`Could not set session_replication_role = replica (FK deferral requires SUPERUSER/REPLICATION): ` +
|
`Could not set session_replication_role = replica (FK deferral requires SUPERUSER/REPLICATION): ` +
|
||||||
`${error instanceof Error ? error.message : String(error)}. ` +
|
`${getErrorMessage(error)}. ` +
|
||||||
`Tables will be copied in name order; FK violations may surface if order is wrong.`,
|
`Tables will be copied in name order; FK violations may surface if order is wrong.`,
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
@@ -383,11 +510,62 @@ async function migrateSqliteToPostgresOnSession(
|
|||||||
|
|
||||||
let copyError: unknown;
|
let copyError: unknown;
|
||||||
try {
|
try {
|
||||||
for (const source of sources) {
|
const plannedSources: Array<{ source: SqliteMigrationSource; plan: readonly TablePlan[] }> = [];
|
||||||
|
for (const [sourceOffset, source] of sources.entries()) {
|
||||||
|
emitMigrationProgress(options, {
|
||||||
|
phase: "scanning-source",
|
||||||
|
sourcePath: source.sqlitePath,
|
||||||
|
sourceSchema: source.pgSchema,
|
||||||
|
sourceIndex: sourceOffset + 1,
|
||||||
|
sourceCount: sources.length,
|
||||||
|
});
|
||||||
const plan = await buildMigrationPlan(migrationDb, source, options.projectId);
|
const plan = await buildMigrationPlan(migrationDb, source, options.projectId);
|
||||||
|
plannedSources.push({ source, plan });
|
||||||
|
}
|
||||||
|
const tableCount = plannedSources.reduce((sum, planned) => sum + planned.plan.length, 0);
|
||||||
|
emitMigrationProgress(options, { phase: "copy-started", tableCount });
|
||||||
|
let tableIndex = 0;
|
||||||
|
for (const { source, plan } of plannedSources) {
|
||||||
for (const tablePlan of plan) {
|
for (const tablePlan of plan) {
|
||||||
const result = await migrateTable(migrationDb, source, tablePlan, dryRun);
|
tableIndex += 1;
|
||||||
|
const progressBase = {
|
||||||
|
sourceSchema: source.pgSchema,
|
||||||
|
table: tablePlan.pgTable,
|
||||||
|
tableIndex,
|
||||||
|
tableCount,
|
||||||
|
} as const;
|
||||||
|
emitMigrationProgress(options, { phase: "table-started", ...progressBase });
|
||||||
|
const result = await migrateTable(
|
||||||
|
migrationDb,
|
||||||
|
source,
|
||||||
|
tablePlan,
|
||||||
|
dryRun,
|
||||||
|
{
|
||||||
|
onCopyProgress: (processedRows, sourceRows) => emitMigrationProgress(options, {
|
||||||
|
phase: "table-progress",
|
||||||
|
...progressBase,
|
||||||
|
processedRows,
|
||||||
|
sourceRows,
|
||||||
|
}),
|
||||||
|
onVerifying: (verificationStage, sourceRows) => emitMigrationProgress(options, {
|
||||||
|
phase: "table-verifying",
|
||||||
|
...progressBase,
|
||||||
|
sourceRows,
|
||||||
|
verificationStage,
|
||||||
|
}),
|
||||||
|
},
|
||||||
|
);
|
||||||
tableResults.push(result);
|
tableResults.push(result);
|
||||||
|
emitMigrationProgress(options, {
|
||||||
|
phase: "table-complete",
|
||||||
|
...progressBase,
|
||||||
|
sourceRows: result.sourceRows,
|
||||||
|
insertedRows: result.insertedRows,
|
||||||
|
targetRows: result.targetRows,
|
||||||
|
verified: result.verified,
|
||||||
|
skipped: result.skipped,
|
||||||
|
skipReason: result.skipReason,
|
||||||
|
});
|
||||||
|
|
||||||
// Bump identity sequences after a real (non-dry-run) copy.
|
// Bump identity sequences after a real (non-dry-run) copy.
|
||||||
if (!dryRun && !result.skipped && result.sourceRows > 0) {
|
if (!dryRun && !result.skipped && result.sourceRows > 0) {
|
||||||
@@ -440,7 +618,13 @@ async function migrateSqliteToPostgresOnSession(
|
|||||||
};
|
};
|
||||||
|
|
||||||
if (dryRun) {
|
if (dryRun) {
|
||||||
log.log(`[dry-run] Migration plan: ${tableResults.length} tables, ${tableResults.reduce((n, t) => n + t.sourceRows, 0)} source rows planned. No writes performed.`);
|
const sourceRows = tableResults.reduce((n, t) => n + t.sourceRows, 0);
|
||||||
|
log.log(`[dry-run] Migration plan: ${tableResults.length} tables, ${sourceRows} source rows planned. No writes performed.`);
|
||||||
|
emitMigrationProgress(options, {
|
||||||
|
phase: "dry-run-complete",
|
||||||
|
tableCount: tableResults.length,
|
||||||
|
sourceRows,
|
||||||
|
});
|
||||||
} else {
|
} else {
|
||||||
const ok = tableResults.filter((t) => t.verified).length;
|
const ok = tableResults.filter((t) => t.verified).length;
|
||||||
const bad = tableResults.length - ok;
|
const bad = tableResults.length - ok;
|
||||||
@@ -452,6 +636,18 @@ async function migrateSqliteToPostgresOnSession(
|
|||||||
WHERE migration_key = ${migrationKey}
|
WHERE migration_key = ${migrationKey}
|
||||||
`);
|
`);
|
||||||
log.log(`Migration complete: ${ok}/${tableResults.length} tables verified (${bad} failed verification). ${sequenceBumps.length} sequences bumped.`);
|
log.log(`Migration complete: ${ok}/${tableResults.length} tables verified (${bad} failed verification). ${sequenceBumps.length} sequences bumped.`);
|
||||||
|
emitMigrationProgress(options, bad === 0 ? {
|
||||||
|
phase: "copy-complete",
|
||||||
|
tableCount: tableResults.length,
|
||||||
|
verifiedTables: ok,
|
||||||
|
failedTables: bad,
|
||||||
|
sequenceBumps: sequenceBumps.length,
|
||||||
|
} : {
|
||||||
|
phase: "failed",
|
||||||
|
tableCount: tableResults.length,
|
||||||
|
verifiedTables: ok,
|
||||||
|
failedTables: bad,
|
||||||
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
return report;
|
return report;
|
||||||
@@ -903,6 +1099,30 @@ interface CanonicalLegacyRow {
|
|||||||
readonly json: string;
|
readonly json: string;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
interface TableMigrationProgressCallbacks {
|
||||||
|
readonly onCopyProgress?: (processedRows: number, sourceRows: number) => void;
|
||||||
|
readonly onVerifying?: (
|
||||||
|
stage: "target-count" | "source-content" | "target-content",
|
||||||
|
sourceRows: number,
|
||||||
|
) => void;
|
||||||
|
}
|
||||||
|
|
||||||
|
function createQuarterProgressReporter(
|
||||||
|
sourceRows: number,
|
||||||
|
onCopyProgress?: (processedRows: number, sourceRows: number) => void,
|
||||||
|
): (processedRows: number) => void {
|
||||||
|
let lastProgressQuarter = 0;
|
||||||
|
return (processedRows) => {
|
||||||
|
// The verifying and complete events represent the terminal 100% state.
|
||||||
|
if (processedRows >= sourceRows) return;
|
||||||
|
const progressQuarter = Math.floor((processedRows / sourceRows) * 4);
|
||||||
|
if (progressQuarter > lastProgressQuarter) {
|
||||||
|
lastProgressQuarter = progressQuarter;
|
||||||
|
onCopyProgress?.(processedRows, sourceRows);
|
||||||
|
}
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
function tagLegacyCell(value: unknown): TaggedLegacyCell {
|
function tagLegacyCell(value: unknown): TaggedLegacyCell {
|
||||||
if (value === null || value === undefined) return { type: "null" };
|
if (value === null || value === undefined) return { type: "null" };
|
||||||
if (Buffer.isBuffer(value)) {
|
if (Buffer.isBuffer(value)) {
|
||||||
@@ -940,6 +1160,7 @@ async function migrateLegacyPreservationTable(
|
|||||||
source: SqliteMigrationSource,
|
source: SqliteMigrationSource,
|
||||||
plan: TablePlan,
|
plan: TablePlan,
|
||||||
dryRun: boolean,
|
dryRun: boolean,
|
||||||
|
progress: TableMigrationProgressCallbacks,
|
||||||
): Promise<TableMigrationResult> {
|
): Promise<TableMigrationResult> {
|
||||||
const projectId = plan.partitionProjectId;
|
const projectId = plan.partitionProjectId;
|
||||||
const legacyPreservation = plan.legacyPreservation;
|
const legacyPreservation = plan.legacyPreservation;
|
||||||
@@ -969,6 +1190,7 @@ async function migrateLegacyPreservationTable(
|
|||||||
const sourceColumnOrder = (sqlite.prepare(
|
const sourceColumnOrder = (sqlite.prepare(
|
||||||
`PRAGMA table_info(${quoteIdent(plan.table)})`,
|
`PRAGMA table_info(${quoteIdent(plan.table)})`,
|
||||||
).all() as Array<{ name: string }>).map(({ name }) => quoteIdent(name)).join(", ");
|
).all() as Array<{ name: string }>).map(({ name }) => quoteIdent(name)).join(", ");
|
||||||
|
const reportCopyProgress = createQuarterProgressReporter(sourceRows, progress.onCopyProgress);
|
||||||
for (let offset = 0; offset < sourceRows; offset += INSERT_BATCH_SIZE) {
|
for (let offset = 0; offset < sourceRows; offset += INSERT_BATCH_SIZE) {
|
||||||
const rawBatch = sqlite.prepare(
|
const rawBatch = sqlite.prepare(
|
||||||
`SELECT * FROM ${quoteIdent(plan.table)} ORDER BY ${sourceColumnOrder} LIMIT ? OFFSET ?`,
|
`SELECT * FROM ${quoteIdent(plan.table)} ORDER BY ${sourceColumnOrder} LIMIT ? OFFSET ?`,
|
||||||
@@ -1005,8 +1227,11 @@ async function migrateLegacyPreservationTable(
|
|||||||
expectedBatch.get(row.legacy_row_hash) === canonicalizeCell(row.legacy_row) &&
|
expectedBatch.get(row.legacy_row_hash) === canonicalizeCell(row.legacy_row) &&
|
||||||
row.source_schema_sql === legacyPreservation.sourceSchemaSql,
|
row.source_schema_sql === legacyPreservation.sourceSchemaSql,
|
||||||
);
|
);
|
||||||
|
const processedRows = Math.min(offset + rawBatch.length, sourceRows);
|
||||||
|
reportCopyProgress(processedRows);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
progress.onVerifying?.("target-count", sourceRows);
|
||||||
const targetCountRows = await db.execute(sql`
|
const targetCountRows = await db.execute(sql`
|
||||||
SELECT COUNT(*)::int AS n
|
SELECT COUNT(*)::int AS n
|
||||||
FROM ${sql.raw(quoteIdent(plan.pgSchema))}.${sql.raw(quoteIdent(plan.pgTable))}
|
FROM ${sql.raw(quoteIdent(plan.pgSchema))}.${sql.raw(quoteIdent(plan.pgTable))}
|
||||||
@@ -1047,9 +1272,10 @@ async function migrateTable(
|
|||||||
source: SqliteMigrationSource,
|
source: SqliteMigrationSource,
|
||||||
plan: TablePlan,
|
plan: TablePlan,
|
||||||
dryRun: boolean,
|
dryRun: boolean,
|
||||||
|
progress: TableMigrationProgressCallbacks = {},
|
||||||
): Promise<TableMigrationResult> {
|
): Promise<TableMigrationResult> {
|
||||||
if (plan.legacyPreservation) {
|
if (plan.legacyPreservation) {
|
||||||
return migrateLegacyPreservationTable(db, source, plan, dryRun);
|
return migrateLegacyPreservationTable(db, source, plan, dryRun, progress);
|
||||||
}
|
}
|
||||||
if (plan.unmappedSourceColumns.length > 0) {
|
if (plan.unmappedSourceColumns.length > 0) {
|
||||||
const sqlite = openSqlite(source.sqlitePath);
|
const sqlite = openSqlite(source.sqlitePath);
|
||||||
@@ -1145,11 +1371,16 @@ async function migrateTable(
|
|||||||
// Stream rows in batches.
|
// Stream rows in batches.
|
||||||
const stmt = sqlite.prepare(`SELECT ${selectableCols} FROM ${quoteIdent(plan.table)}`);
|
const stmt = sqlite.prepare(`SELECT ${selectableCols} FROM ${quoteIdent(plan.table)}`);
|
||||||
const batch: Record<string, unknown>[] = [];
|
const batch: Record<string, unknown>[] = [];
|
||||||
|
let processedRows = 0;
|
||||||
|
const reportCopyProgress = createQuarterProgressReporter(sourceRows, progress.onCopyProgress);
|
||||||
const flush = async (): Promise<void> => {
|
const flush = async (): Promise<void> => {
|
||||||
if (batch.length === 0) return;
|
if (batch.length === 0) return;
|
||||||
|
const batchSize = batch.length;
|
||||||
const inserted = await insertBatch(db, plan, insertableCols, batch, hasIdentityCol);
|
const inserted = await insertBatch(db, plan, insertableCols, batch, hasIdentityCol);
|
||||||
insertedRows += inserted;
|
insertedRows += inserted;
|
||||||
|
processedRows += batchSize;
|
||||||
batch.length = 0;
|
batch.length = 0;
|
||||||
|
reportCopyProgress(processedRows);
|
||||||
};
|
};
|
||||||
|
|
||||||
for (const row of stmt.all() as Array<Record<string, unknown>>) {
|
for (const row of stmt.all() as Array<Record<string, unknown>>) {
|
||||||
@@ -1168,6 +1399,7 @@ async function migrateTable(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
await flush();
|
await flush();
|
||||||
|
progress.onVerifying?.("target-count", sourceRows);
|
||||||
|
|
||||||
// Verify the migration.
|
// Verify the migration.
|
||||||
// FNXC:PostgresMigration 2026-06-26-15:40 (fix migration-review P1 #15):
|
// FNXC:PostgresMigration 2026-06-26-15:40 (fix migration-review P1 #15):
|
||||||
@@ -1197,9 +1429,11 @@ async function migrateTable(
|
|||||||
: targetRows === sourceRows;
|
: targetRows === sourceRows;
|
||||||
let contentOk = true;
|
let contentOk = true;
|
||||||
if (rowCountOk && sourceRows > 0) {
|
if (rowCountOk && sourceRows > 0) {
|
||||||
|
progress.onVerifying?.("source-content", sourceRows);
|
||||||
const sourceCanonicalRows = computeSourceCanonicalRows(
|
const sourceCanonicalRows = computeSourceCanonicalRows(
|
||||||
sqlite, plan.table, insertableCols, plan.partitionProjectId,
|
sqlite, plan.table, insertableCols, plan.partitionProjectId,
|
||||||
);
|
);
|
||||||
|
progress.onVerifying?.("target-content", sourceRows);
|
||||||
const targetCanonicalRows = await computeTargetCanonicalRows(
|
const targetCanonicalRows = await computeTargetCanonicalRows(
|
||||||
db, plan.pgSchema, plan.pgTable, insertableCols, plan.partitionProjectId,
|
db, plan.pgSchema, plan.pgTable, insertableCols, plan.partitionProjectId,
|
||||||
);
|
);
|
||||||
|
|||||||
@@ -463,7 +463,7 @@ export async function createTaskStoreForBackend(
|
|||||||
considered.
|
considered.
|
||||||
*/
|
*/
|
||||||
const migrationKey = `project:${migrationProjectId ?? rootDir}`;
|
const migrationKey = `project:${migrationProjectId ?? rootDir}`;
|
||||||
const { migrateSqliteToPostgres, defaultMigrationSources, isSqliteMigrationComplete, completeSqliteMigration, recordSqliteMigrationComplete, CENTRAL_SQLITE_MIGRATION_KEY } = await import("./sqlite-migrator.js");
|
const { migrateSqliteToPostgres, defaultMigrationSources, formatMigrationProgress, isSqliteMigrationComplete, completeSqliteMigration, recordSqliteMigrationComplete, CENTRAL_SQLITE_MIGRATION_KEY } = await import("./sqlite-migrator.js");
|
||||||
const migrationComplete = await isSqliteMigrationComplete(connections.migration, migrationKey);
|
const migrationComplete = await isSqliteMigrationComplete(connections.migration, migrationKey);
|
||||||
if (!migrationComplete && isValidSqliteDatabaseFile(legacySqlitePath)) {
|
if (!migrationComplete && isValidSqliteDatabaseFile(legacySqlitePath)) {
|
||||||
// The central (global-dir) source is optional: when no global dir is
|
// The central (global-dir) source is optional: when no global dir is
|
||||||
@@ -486,6 +486,13 @@ export async function createTaskStoreForBackend(
|
|||||||
projectId: migrationProjectId,
|
projectId: migrationProjectId,
|
||||||
migrationKey,
|
migrationKey,
|
||||||
deferCompletion: true,
|
deferCompletion: true,
|
||||||
|
/*
|
||||||
|
FNXC:CliMigrationProgress 2026-07-14-13:47:
|
||||||
|
First-boot migration can copy hundreds of thousands of rows. Forward structured phase, table, quarter-copy, and terminal events to the CLI logger so an operator sees forward progress and an explicit rollback instead of a silent startup wait.
|
||||||
|
*/
|
||||||
|
onProgress: (event) => {
|
||||||
|
log.log(`startup-factory: SQLite migration — ${formatMigrationProgress(event)}`);
|
||||||
|
},
|
||||||
});
|
});
|
||||||
/*
|
/*
|
||||||
FNXC:PostgresMigrationVerification 2026-07-13-22:37:
|
FNXC:PostgresMigrationVerification 2026-07-13-22:37:
|
||||||
|
|||||||
Reference in New Issue
Block a user