From 19a4bc6585f6bb0cd9b5eace9e4ec16e4927e70e Mon Sep 17 00:00:00 2001 From: Ben Senescu <44480372+bensenescu@users.noreply.github.com> Date: Mon, 29 Jun 2026 20:20:23 -0400 Subject: [PATCH] =?UTF-8?q?D1=20=E2=86=92=20Postgres=20data=20migration=20?= =?UTF-8?q?(ETL=20+=20runbook)=20(#274)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .gitignore | 5 +- drizzle-pg.config.ts | 5 + runbooks/d1-to-postgres-detailed.md | 149 +++++++++ runbooks/d1-to-postgres-simple.md | 81 +++++ scripts/migrate-d1-to-postgres.ts | 493 ++++++++++++++++++++++++++++ 5 files changed, 731 insertions(+), 2 deletions(-) create mode 100644 runbooks/d1-to-postgres-detailed.md create mode 100644 runbooks/d1-to-postgres-simple.md create mode 100644 scripts/migrate-d1-to-postgres.ts diff --git a/.gitignore b/.gitignore index c74a2d3..16de911 100644 --- a/.gitignore +++ b/.gitignore @@ -13,8 +13,9 @@ dist-sourcemaps/ .DS_Store .cache .env -.env.local -.env.production +.env. +.env.* +!.env.example .vercel .output .nitro diff --git a/drizzle-pg.config.ts b/drizzle-pg.config.ts index 2db2cd6..6ed0922 100644 --- a/drizzle-pg.config.ts +++ b/drizzle-pg.config.ts @@ -1,4 +1,9 @@ import { defineConfig } from "drizzle-kit"; +import { loadLocalEnv } from "./scripts/cli-utils"; + +// Pull POSTGRES_DATABASE_URL from .env.local (no-op if already in the shell env), +// so the migration runbooks' single .env.local works for `db:migrate:pg` too. +loadLocalEnv(); export default defineConfig({ dialect: "postgresql", diff --git a/runbooks/d1-to-postgres-detailed.md b/runbooks/d1-to-postgres-detailed.md new file mode 100644 index 0000000..66d5ac7 --- /dev/null +++ b/runbooks/d1-to-postgres-detailed.md @@ -0,0 +1,149 @@ +# D1 → Postgres migration — detailed runbook + +_Last updated: 2026-06-29._ + +When a hosted instance outgrows D1's storage ceiling, you switch it to the +Postgres backend (`DATABASE_PROVIDER=postgres`). The provider switch changes +where the app reads and writes — it does **not** move existing data. This runbook +copies the data into a freshly-migrated Postgres database using +`scripts/migrate-d1-to-postgres.ts`. For the condensed happy path see +[d1-to-postgres-simple.md](./d1-to-postgres-simple.md). + +The script reads **each table directly from D1 over the Cloudflare REST API** and +writes to Postgres — there is no SQL dump to download or reimport. (An earlier +dump-based approach was dropped after a `wrangler d1 export` download silently +truncated yet reported success; reading tables over the API avoids that whole +class of failure and reads live, complete data.) + +D1 is left untouched throughout, so **rollback is just flipping the provider flag +back to `d1`** until you've confirmed the Postgres cutover is healthy. + +> **Schema drift:** the script was authored 2026-06-29 and reflects the schema as +> of that date (it enumerates every table via Drizzle). If the schema changed +> after this was merged — a new table, a renamed timestamp column, a new table +> without a primary key — re-read the script (especially `deltaPredicate` and +> `conflictArbiterNames`) before relying on it. + +## What the script converts + +Some columns are dialect-native, so a raw copy doesn't work. Each raw D1 value is +run through the SQLite column's Drizzle codec and written through the Postgres +column, so the conversions happen with the same logic the app itself uses: + +- **better-auth tables** (`user`, `session`, `account`, …): SQLite stores + timestamps as integer epoch-ms and booleans as integer `0/1`; Postgres uses + `timestamptz` and real `boolean`. +- **App tables** (`projects`, `saved_keywords`, `billing_customer_status`, …): + timestamps are text in both, but D1's `current_timestamp` default + (`YYYY-MM-DD HH:MM:SS`) is rewritten to the ISO-8601 form the Postgres code + expects (`YYYY-MM-DDTHH:MM:SS.000Z`). Values the app already wrote in ISO are + left untouched. + +Tables are copied in foreign-key-safe order, and inserts use +`onConflictDoNothing`, so the script is safe to re-run. + +After the copy, the script advances the Postgres serial sequences +(`keyword_metrics.id`, `rank_snapshots.id`) to the max migrated id. Those ids are +auto-incremented in SQLite and copied verbatim, so without this step new inserts +after cutover would collide with migrated rows. + +## Prerequisites + +- The provider-aware build is deployed (this branch / its parent). +- Credentials in `.env.local` (the migration script auto-loads it; so does + `db:migrate:pg`): + ```sh + CLOUDFLARE_ACCOUNT_ID=... + CLOUDFLARE_API_TOKEN=... # needs D1 read + POSTGRES_DATABASE_URL=postgres://user:pass@host:5432/db + # CLOUDFLARE_D1_DATABASE_ID=... # optional; otherwise read from wrangler.jsonc + ``` + Get your account id from `wrangler whoami`. +- A Postgres database provisioned and **freshly migrated, empty**: + ```sh + pnpm db:migrate:pg + ``` + +## Steps + +1. **Dry run** to review per-table row counts before writing anything. This only + reads D1 (no Postgres connection is opened): + + ```sh + pnpm exec tsx scripts/migrate-d1-to-postgres.ts --dry-run + ``` + +2. **Migrate** into the empty Postgres. The script aborts if the target already + has rows (pass `--allow-nonempty` to override; it stays idempotent): + + ```sh + pnpm exec tsx scripts/migrate-d1-to-postgres.ts + ``` + + Confirm it ends with **"All row counts match."** + +3. **Cut over.** Point the deployment at Postgres (`DATABASE_PROVIDER=postgres` + plus a Hyperdrive binding or `POSTGRES_DATABASE_URL`) and deploy: + + ```sh + pnpm deploy:postgres + ``` + +4. **Verify.** Smoke-test the live app — load a project, save a keyword, check + billing. + +> **Optional: freeze writes for a consistent snapshot.** The copy reads each +> table as it goes, so writes that land mid-run can be missed. For a perfectly +> consistent one-shot copy, pause the scheduled rank-check cron and put the +> instance into a brief read-only window during steps 1–3, then re-enable after. +> This is **not required** — the low-downtime path below avoids it. + +## Low-downtime cutover (delta catch-up) + +The big tables (`keyword_metrics`, `audit_pages`, …) dominate the copy time, so a +full freeze can mean several minutes of downtime. To avoid it, do the bulk copy +**live**, then a short delta catch-up just before cutover: + +1. **Bulk copy, live** — run the full migration (step 2) while the app is still + serving. It's allowed to miss writes that land mid-copy; the delta pass below + reconciles them. +2. **Delta sync** — copy only what changed since the bulk copy and **upsert** it: + + ```sh + pnpm exec tsx scripts/migrate-d1-to-postgres.ts --update --since-hours 12 + ``` + + Set `--since-hours` to comfortably cover the gap between the bulk copy and now. + Genuinely large, append-mostly tables (`keyword_metrics`, `rank_snapshots`, + `audit_pages`, `audit_lighthouse_results`) are filtered to their recent rows; + everything else — including mutable tables that take in-place updates (audits, + rank-check runs, tracked keywords) and the small config/entity tables — is + fully re-upserted, so new signups and updates to existing rows (refreshed + tokens, archived projects, advanced schedules, completed runs) are all picked + up. Confirm **"All row counts match."** + +3. **Cut over** (step 3 above) and verify. + + For the smallest possible window you can briefly freeze writes between the + delta sync and cutover, but it's usually unnecessary. + +`--update` is safe to re-run, skips the empty-target preflight, and re-advances +the serial sequences. **Deletes are not synced** — a row deleted in D1 during the +window stays in Postgres (it surfaces as a count mismatch in the verify step, +e.g. a few expired sessions/verification rows). If deletes during the window must +be reflected, do the full copy under a write-freeze instead. + +## Rollback + +If anything looks wrong after cutover, set `DATABASE_PROVIDER=d1` (remove the +Postgres binding) and redeploy. D1 still holds the original data, untouched. + +## Notes + +- The migration reads D1 in pages (`--page-size`, default 5000) and writes in + FK-safe order; it holds at most one page in memory at a time. +- Local-only practice run: see + [`../docs/LOCAL_POSTGRES.md`](../docs/LOCAL_POSTGRES.md) to rehearse against a + Docker Postgres first. +- Flags: `--dry-run`, `--allow-nonempty`, `--page-size N`, `--update`, + `--since-hours N`. Run with no flags for a full one-time copy. diff --git a/runbooks/d1-to-postgres-simple.md b/runbooks/d1-to-postgres-simple.md new file mode 100644 index 0000000..2920f1f --- /dev/null +++ b/runbooks/d1-to-postgres-simple.md @@ -0,0 +1,81 @@ +# D1 → Postgres migration — simple runbook + +_Last updated: 2026-06-29._ + +The happy path for moving a hosted instance from D1 to Postgres. For the full +detail — what the script converts, the low-downtime delta sync, cutover and +rollback — see +[d1-to-postgres-detailed.md](./d1-to-postgres-detailed.md). + +> Switching `DATABASE_PROVIDER` to `postgres` changes where the app reads and +> writes — it does **not** move existing data. This copies the data. D1 is never +> written to, so rollback is just flipping the provider back to `d1`. + +## 1. Credentials → `.env.local` + +The migration script auto-loads `.env.local` (no inline env vars needed): + +```sh +CLOUDFLARE_ACCOUNT_ID=... +CLOUDFLARE_API_TOKEN=... # needs D1 read +POSTGRES_DATABASE_URL=postgres://user:pass@host:5432/db +# CLOUDFLARE_D1_DATABASE_ID=... # optional; otherwise read from wrangler.jsonc +``` + +## 2. Create the Postgres schema + +Provision an empty Postgres, then apply the schema: + +```sh +pnpm db:migrate:pg +``` + +## 3. Dry run (read-only) + +Reports per-table row counts and writes nothing: + +```sh +pnpm exec tsx scripts/migrate-d1-to-postgres.ts --dry-run +``` + +## 4. Migrate + +```sh +pnpm exec tsx scripts/migrate-d1-to-postgres.ts +``` + +Confirm it ends with **"All row counts match."** (Re-runnable; it aborts if the +target already has data — pass `--allow-nonempty` to override.) + +## 5. Catch-up sync (recommended) + +Right before cutover, run the script again with `--update` to copy anything +written or changed since the bulk copy (new signups, fresh rank checks, etc.) — +cheap insurance that nothing was missed: + +```sh +pnpm exec tsx scripts/migrate-d1-to-postgres.ts --update +``` + +Confirm **"All row counts match."** (A small mismatch from rows _deleted_ in D1 +during the window is expected — see the detailed runbook.) + +## 6. Cut over + +Point the deployment at Postgres (`DATABASE_PROVIDER=postgres` plus a Hyperdrive +binding or `POSTGRES_DATABASE_URL`) and deploy, then smoke-test (load a project, +save a keyword, check billing): + +```sh +pnpm deploy:postgres +``` + +> **Optional — freeze writes** for a perfectly consistent snapshot: pause the +> rank-check cron / put the app in a read-only window from the dry run through +> cutover. Not required — the catch-up sync in step 5 covers writes made during a +> live copy. + +## Rollback + +Set `DATABASE_PROVIDER=d1` (remove the Postgres binding) and redeploy. D1 still +holds the original data, untouched. diff --git a/scripts/migrate-d1-to-postgres.ts b/scripts/migrate-d1-to-postgres.ts new file mode 100644 index 0000000..b6a5ba8 --- /dev/null +++ b/scripts/migrate-d1-to-postgres.ts @@ -0,0 +1,493 @@ +/** + * One-time D1 (SQLite) -> Postgres data migration. + * + * Authored 2026-06-29. This script is coupled to the schema as of that date — it + * reflects every table over Drizzle, so a schema change merged later (new table, + * renamed timestamp column, a new no-PK/unique-index table, etc.) may require + * updating `deltaPredicate`/`conflictArbiterNames` before re-using it. Re-read it + * against the current schema before relying on it. + * + * The provider-aware DB layer (DATABASE_PROVIDER=postgres) ships the schema and + * runtime, but NOT the data. This script copies the data table-by-table, reading + * each table directly from D1 over the Cloudflare REST API and writing to a + * freshly-migrated Postgres database. Writes may be frozen during the copy for a + * perfectly consistent snapshot, but the `--update` delta pass (below) makes that + * optional — see the runbooks. + * + * Why read D1 directly (vs `wrangler d1 export` + reimport): a SQL-dump round + * trip is fragile — the export download can silently truncate yet report + * success, and reimporting 400MB+ of SQL text is its own failure surface. + * Querying each table over the API removes that entire class of problem and + * reads live, complete data. + * + * Type conversions (the reason a raw copy doesn't work) are delegated to + * Drizzle's own column codecs: we run each raw D1 value through the SQLite + * column's `mapFromDriverValue` (integer epoch-ms -> Date, 0/1 -> boolean) and + * write through the Postgres column (Date -> timestamptz, boolean -> boolean). + * App-table timestamps are TEXT in both dialects, but D1's `current_timestamp` + * default is `YYYY-MM-DD HH:MM:SS`; we rewrite those to the ISO-8601 form the + * Postgres code writes. Values already in ISO are left untouched. + * + * Tables are copied in foreign-key-safe order and inserts use + * `onConflictDoNothing`, so a failed run is safe to re-run from the top. + * + * Setup: put the credentials in `.env.local` (auto-loaded), then run the script + * directly with tsx. See the runbooks for the full procedure: + * - runbooks/d1-to-postgres-simple.md (happy path) + * - runbooks/d1-to-postgres-detailed.md (full detail + cutover) + * + * # .env.local + * CLOUDFLARE_ACCOUNT_ID=... + * CLOUDFLARE_API_TOKEN=... # D1 read + * POSTGRES_DATABASE_URL=postgres://user:pass@host:5432/db + * # CLOUDFLARE_D1_DATABASE_ID=... # optional; else read from wrangler.jsonc + * + * Usage: + * pnpm exec tsx scripts/migrate-d1-to-postgres.ts [flags] + * + * --dry-run Report D1 row counts only; write nothing. + * --allow-nonempty Proceed even if the target already has rows (default: abort). + * --page-size N Rows per D1 API page (default 5000). + * --update Delta/catch-up sync: only rows changed within the window + * (default 12h), UPSERTing instead of insert-or-ignore. Run + * this right before cutover to top up an already-migrated + * Postgres. Big time-series tables are filtered by their + * recency column; small tables are fully re-upserted (which + * also catches updates to old rows). Deletes are not synced. + * --since-hours N Window for --update (default 12). + */ +import { readFileSync } from "node:fs"; +import process from "node:process"; +import { + getTableColumns, + getTableName, + is, + sql, + Table, + Column, + type SQL, +} from "drizzle-orm"; +import { getTableConfig as getSqliteTableConfig } from "drizzle-orm/sqlite-core"; +import { getTableConfig as getPgTableConfig } from "drizzle-orm/pg-core"; +import { drizzle as drizzlePg } from "drizzle-orm/postgres-js"; +import postgres from "postgres"; +import { loadLocalEnv, parseArgs } from "./cli-utils"; +// The Node-safe raw barrels (not ../src/db/schema, the provider-aware one, which +// imports cloudflare:workers). Importing the barrels keeps the table list in sync +// automatically if a schema file is later added. +import * as sqliteSchema from "../src/db/d1/schema"; +import * as pgSchema from "../src/db/pg/schema"; + +loadLocalEnv(); + +const args = parseArgs(process.argv.slice(2)); +const dryRun = args["dry-run"] === "true"; +const allowNonEmpty = args["allow-nonempty"] === "true"; +const pageSize = Number(args["page-size"]) || 5000; +const updateMode = args["update"] === "true"; +const sinceHours = Number(args["since-hours"]) || 12; + +const accountId = process.env.CLOUDFLARE_ACCOUNT_ID; +const apiToken = process.env.CLOUDFLARE_API_TOKEN; +const connectionString = process.env.POSTGRES_DATABASE_URL; + +// D1's `current_timestamp` text is `YYYY-MM-DD HH:MM:SS`. App code elsewhere +// writes `new Date().toISOString()`, which already has a `T` separator and never +// matches this pattern, so it is left untouched. +const SPACE_TIMESTAMP = /^\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2}$/u; +const toIsoText = (value: string) => `${value.replace(" ", "T")}.000Z`; + +function resolveDatabaseId(): string { + const override = process.env.CLOUDFLARE_D1_DATABASE_ID?.trim(); + if (override) return override; + const wrangler = readFileSync("wrangler.jsonc", "utf8"); + const match = wrangler.match(/"database_id"\s*:\s*"([^"]+)"/u); + if (!match) { + throw new Error( + "Could not find database_id in wrangler.jsonc; set CLOUDFLARE_D1_DATABASE_ID.", + ); + } + return match[1]; +} + +async function d1Query>( + databaseId: string, + statement: string, +): Promise { + const response = await fetch( + `https://api.cloudflare.com/client/v4/accounts/${accountId}/d1/database/${databaseId}/query`, + { + method: "POST", + headers: { + Authorization: `Bearer ${apiToken}`, + "Content-Type": "application/json", + }, + body: JSON.stringify({ sql: statement }), + }, + ); + const json = (await response.json()) as { + success: boolean; + errors?: unknown; + result?: { results: T[] }[]; + }; + if (!response.ok || !json.success) { + throw new Error(`D1 query failed: ${JSON.stringify(json.errors ?? json)}`); + } + return json.result?.[0]?.results ?? []; +} + +function tablesByName(modules: Record[]): Map { + const out = new Map(); + for (const mod of modules) { + for (const value of Object.values(mod)) { + if (is(value, Table)) out.set(getTableName(value), value); + } + } + return out; +} + +// Order tables so each is copied after the tables its foreign keys reference +// (Kahn's algorithm over the SQLite FK graph). Self-references are ignored. +function fkSafeOrder(tables: Map): string[] { + const deps = new Map>(); + for (const [name, table] of tables) { + const referenced = new Set(); + for (const fk of getSqliteTableConfig(table).foreignKeys) { + const target = getTableName(fk.reference().foreignTable); + if (target !== name && tables.has(target)) referenced.add(target); + } + deps.set(name, referenced); + } + + const ordered: string[] = []; + const placed = new Set(); + while (ordered.length < tables.size) { + const ready = [...deps] + .filter(([name]) => !placed.has(name)) + .filter(([, refs]) => [...refs].every((r) => placed.has(r))) + .map(([name]) => name) + .sort(); + if (ready.length === 0) { + // No FK cycle exists in this schema, so this is unreachable. Fail loud + // rather than copy the remaining tables in an arbitrary order that could + // violate FK constraints on insert. + const remaining = [...tables.keys()].filter((n) => !placed.has(n)); + throw new Error( + `Foreign-key cycle detected among: ${remaining.join(", ")}. Cannot order these tables for a constraint-safe copy.`, + ); + } + for (const name of ready) { + ordered.push(name); + placed.add(name); + } + } + return ordered; +} + +function primaryKeyColumnNames(table: Table): string[] { + const config = getSqliteTableConfig(table); + const cols = new Set(); + for (const col of Object.values(getTableColumns(table))) { + if (col.primary) cols.add(col.name); + } + for (const composite of config.primaryKeys) { + for (const col of composite.columns) cols.add(col.name); + } + return [...cols]; +} + +// Stable pagination order: the table's primary key. These tables append at the +// tail, so OFFSET paging is stable during a frozen full copy; any mid-table +// churn during a live bulk copy is reconciled by the --update delta pass. +function primaryKeyOrder(table: Table): string { + const cols = primaryKeyColumnNames(table); + return cols.length ? `ORDER BY ${cols.map((c) => `"${c}"`).join(",")}` : ""; +} + +// Conflict arbiter for --update upserts: the table's PK if it has one, else the +// columns of its first unique index. Join tables (saved_keyword_tag_assignments) +// have no PK — their identity is a unique index — so a PK-only arbiter would emit +// `ON CONFLICT ()` and fail to parse. +function conflictArbiterNames(table: Table): string[] { + const pk = primaryKeyColumnNames(table); + if (pk.length) return pk; + for (const idx of getSqliteTableConfig(table).indexes) { + if (idx.config.unique) { + const names: string[] = []; + for (const col of idx.config.columns) { + // Index entries can be raw SQL; only plain columns have a name. + if (is(col, Column)) names.push(col.name); + } + return names; + } + } + return []; +} + +// --update mode: WHERE fragment selecting rows changed within the window. Only +// genuinely large, append-mostly tables are filtered here. Everything else falls +// through to a full re-upsert each run — which is the ONLY way to catch in-place +// updates to pre-window rows, since those rows' timestamps don't move. That +// matters for mutable tables that aren't huge: audits / rank_check_runs (status, +// counts, completedAt updated after creation), rank_tracking_keywords (metrics +// refreshed in place), plus the small config/entity tables. Filtering them by a +// creation timestamp would silently drop those updates, and the row-count verify +// can't catch it. SQLite datetime() normalizes both ISO and D1's legacy +// space-format text, so the comparison is format-agnostic. The cutoff is +// script-generated, not user input. Deletes are never synced by insert/upsert. +function deltaPredicate(table: string, cutoffIso: string): string | null { + const since = `datetime('${cutoffIso}')`; + switch (table) { + // keyword_metrics is mutable but its fetched_at advances on every upsert, so + // a recency filter still catches in-place metric refreshes. + case "keyword_metrics": + return `datetime("fetched_at") >= ${since}`; + case "rank_snapshots": + return `datetime("checked_at") >= ${since}`; + // Audit child rows carry no timestamp of their own — scope by parent audit. + case "audit_pages": + case "audit_lighthouse_results": + return `"audit_id" IN (SELECT id FROM audits WHERE datetime("started_at") >= ${since})`; + default: + return null; + } +} + +// ON CONFLICT DO UPDATE target (the conflict arbiter columns) + set (every other +// column to the incoming "excluded" value), so --update mode overwrites changed +// rows. The target is sourced from the pg table config so it carries the pg +// column type onConflictDoUpdate requires; the set is keyed by JS field name to +// match the values() payload. +function buildUpsert(pgTable: Table, arbiterNames: string[]) { + const pk = new Set(arbiterNames); + const target = getPgTableConfig(pgTable).columns.filter((c) => + pk.has(c.name), + ); + const set: Record = {}; + for (const [jsKey, col] of Object.entries(getTableColumns(pgTable)) as [ + string, + Column, + ][]) { + if (!pk.has(col.name)) { + set[jsKey] = sql`excluded.${sql.identifier(col.name)}`; + } + } + return { target, set }; +} + +// Decode a raw D1 row into the shape Drizzle's Postgres insert expects, reusing +// each SQLite column's codec (epoch-ms -> Date, 0/1 -> boolean), then normalize +// legacy space-format timestamp text to ISO. +function convertRow( + sqliteTable: Table, + row: Record, +): Record { + const out: Record = {}; + for (const [jsKey, column] of Object.entries( + getTableColumns(sqliteTable), + ) as [string, Column][]) { + const raw = row[column.name]; + if (raw === null || raw === undefined) { + out[jsKey] = null; + continue; + } + let value = column.mapFromDriverValue(raw); + if (typeof value === "string" && SPACE_TIMESTAMP.test(value)) { + value = toIsoText(value); + } + out[jsKey] = value; + } + return out; +} + +async function countD1( + databaseId: string, + table: string, + predicate?: string | null, +): Promise { + const where = predicate ? ` WHERE ${predicate}` : ""; + const [row] = await d1Query<{ c: number }>( + databaseId, + `SELECT count(*) AS c FROM "${table}"${where}`, + ); + return Number(row?.c ?? 0); +} + +type PgDb = ReturnType; + +// Copy one table from D1 to Postgres, paginating in PK order. `predicate` scopes +// the read (null = whole table); `upsert` overwrites on PK conflict (--update) +// vs insert-or-ignore (full migration). Returns the row count read from D1. +async function copyTable( + databaseId: string, + dest: PgDb, + name: string, + sqliteTable: Table, + pgTable: Table, + predicate: string | null, + upsert: boolean, +): Promise { + const orderBy = primaryKeyOrder(sqliteTable); + const where = predicate ? `WHERE ${predicate} ` : ""; + // Keep each INSERT's bound-parameter count (rows × columns) well under + // Postgres's 65535 limit. + const colCount = Object.keys(getTableColumns(pgTable)).length; + const insertBatch = Math.max(1, Math.min(5000, Math.floor(50000 / colCount))); + const { target, set } = buildUpsert( + pgTable, + conflictArbiterNames(sqliteTable), + ); + // Upsert needs an arbiter and at least one column to set; without both (e.g. a + // table with no PK/unique index, or all-arbiter columns) fall back to ignore. + const canUpsert = upsert && target.length > 0 && Object.keys(set).length > 0; + + let offset = 0; + let count = 0; + for (;;) { + const page = await d1Query>( + databaseId, + `SELECT * FROM "${name}" ${where}${orderBy} LIMIT ${pageSize} OFFSET ${offset}`, + ); + if (page.length === 0) break; + + const rows = page.map((row) => convertRow(sqliteTable, row)); + for (let i = 0; i < rows.length; i += insertBatch) { + const chunk = rows.slice(i, i + insertBatch); + await (canUpsert + ? dest.insert(pgTable).values(chunk).onConflictDoUpdate({ target, set }) + : dest.insert(pgTable).values(chunk).onConflictDoNothing()); + } + count += page.length; + offset += pageSize; + if (page.length < pageSize) break; + } + return count; +} + +async function main() { + if (!accountId || !apiToken) { + throw new Error( + "CLOUDFLARE_ACCOUNT_ID and CLOUDFLARE_API_TOKEN are required.", + ); + } + if (!dryRun && !connectionString) { + throw new Error("POSTGRES_DATABASE_URL is required."); + } + const databaseId = resolveDatabaseId(); + + const sqliteTables = tablesByName([sqliteSchema]); + const pgTables = tablesByName([pgSchema]); + const order = fkSafeOrder(sqliteTables); + const cutoffIso = new Date(Date.now() - sinceHours * 3_600_000).toISOString(); + + const mode = updateMode ? `UPDATE — last ${sinceHours}h` : "FULL"; + console.log( + `Migrating ${order.length} tables [${mode}]: D1 ${databaseId} -> Postgres${dryRun ? " (DRY RUN)" : ""}\n`, + ); + + if (dryRun) { + let total = 0; + for (const name of order) { + const predicate = updateMode ? deltaPredicate(name, cutoffIso) : null; + const count = await countD1(databaseId, name, predicate); + total += count; + console.log(` ${name.padEnd(34)} ${String(count).padStart(7)}`); + } + console.log( + `\nDone. ${total} ${updateMode ? "rows to sync" : "source rows"}. (dry run — nothing written)`, + ); + return; + } + + const pgClient = postgres(connectionString!, { max: 1 }); + const dest = drizzlePg(pgClient); + + // Preflight: refuse to write into a populated target unless opted in. + // Skipped in --update mode, where the target is expected to already hold the + // full migration this run only tops up. + if (!allowNonEmpty && !updateMode) { + for (const name of order) { + const [{ count }] = await dest + .select({ count: sql`count(*)::int` }) + .from(pgTables.get(name)!); + if (Number(count) > 0) { + await pgClient.end(); + throw new Error( + `Target table "${name}" already has ${count} rows. Migrate into a freshly-created Postgres, or pass --allow-nonempty (idempotent via onConflictDoNothing).`, + ); + } + } + } + + const written = new Map(); + for (const name of order) { + const predicate = updateMode ? deltaPredicate(name, cutoffIso) : null; + const count = await copyTable( + databaseId, + dest, + name, + sqliteTables.get(name)!, + pgTables.get(name)!, + predicate, + updateMode, + ); + written.set(name, count); + const verb = updateMode ? "synced" : "wrote"; + console.log(` ${name.padEnd(34)} ${verb} ${String(count).padStart(7)}`); + } + + // Serial PKs (keyword_metrics.id, rank_snapshots.id) were copied with their + // explicit D1 ids, but the Postgres sequences still sit at 1. Without advancing + // them, post-cutover inserts that omit id and rely on the serial default would + // collide with migrated rows — silently dropped for rank_snapshots (untargeted + // onConflictDoNothing) and a hard duplicate-key error for keyword_metrics. Bump + // each serial column's sequence to max(id) (3-arg setval keeps next id at 1 when + // the table is empty). + console.log("\nResetting serial sequences..."); + for (const name of order) { + const pgTable = pgTables.get(name)!; + for (const column of Object.values(getTableColumns(pgTable))) { + if (column.getSQLType() !== "serial") continue; + const col = sql.identifier(column.name); + const tbl = sql.identifier(name); + await dest.execute(sql` + SELECT setval( + pg_get_serial_sequence(${name}, ${column.name}), + COALESCE((SELECT max(${col}) FROM ${tbl}), 1), + (SELECT max(${col}) FROM ${tbl}) IS NOT NULL + ) + `); + console.log(` ${name}.${column.name} sequence reset`); + } + } + + // Verify: D1 count vs Postgres count per table. + console.log("\nVerifying row counts..."); + let mismatches = 0; + for (const name of order) { + const sourceCount = await countD1(databaseId, name); + const [{ count }] = await dest + .select({ count: sql`count(*)::int` }) + .from(pgTables.get(name)!); + if (Number(count) !== sourceCount) { + mismatches += 1; + console.log(` MISMATCH ${name}: D1 ${sourceCount} vs postgres ${count}`); + } + } + + await pgClient.end(); + + const total = [...written.values()].reduce((n, v) => n + v, 0); + console.log( + `\nDone. ${total} rows across ${order.length} tables.` + + (mismatches === 0 + ? " All row counts match." + : ` ${mismatches} table(s) mismatched — investigate before cutover.`), + ); + if (mismatches > 0) process.exitCode = 1; +} + +main().catch((error) => { + console.error(error); + process.exit(1); +});