#!/usr/bin/env -S npx tsx

import { createHash } from "node:crypto";
import { chmod, mkdir, readFile, rename, writeFile } from "node:fs/promises";
import path from "node:path";
import { fileURLToPath } from "node:url";

type Row = Record<string, unknown>;
type Query = (sql: string, parameters?: unknown[]) => Promise<Row[]>;

interface Options {
  outputRoot: string;
  pageSize: number;
  help: boolean;
}

interface TableDefinition {
  schema: string;
  table: string;
  relationKind: string;
  columns: string[];
  primaryKey: string[];
}

interface TableState {
  schema: string;
  table: string;
  count: number;
  fingerprintA: string;
  fingerprintB: string;
}

interface ArchiveFile {
  path: string;
  kind: "table" | "metadata";
  dataset: string;
  page: number;
  rowCount: number;
  bytes: number;
  sha256: string;
}

interface TableManifest {
  dataset: string;
  schema: string;
  table: string;
  relationKind: string;
  columns: string[];
  primaryKey: string[];
  pagination: "primary-key" | "physical-cursor";
  rowCountBefore: number;
  rowCountAfter: number;
  fingerprintBefore: [string, string];
  fingerprintAfter: [string, string];
  exportedRows: number;
  pages: number;
}

const scriptDirectory = path.dirname(fileURLToPath(import.meta.url));
const backendDirectory = path.resolve(scriptDirectory, "..");
const defaultOutputRoot = path.join(backendDirectory, "data", "migration", "source-archive");
const managementApiBaseUrl = "https://api.supabase.com/v1";
const envPath = path.join(backendDirectory, ".env");
const formatVersion = 1;
const allowedSchemas = ["public", "auth", "storage"] as const;

const catalogQueries: Record<string, string> = {
  tables: `
    SELECT n.nspname AS schema_name,
           c.relname AS table_name,
           c.relkind::text AS relation_kind,
           c.relpersistence::text AS persistence,
           c.relispartition AS is_partition,
           c.relrowsecurity AS row_security,
           c.relforcerowsecurity AS force_row_security,
           c.relreplident::text AS replica_identity,
           pg_get_userbyid(c.relowner) AS owner,
           obj_description(c.oid, 'pg_class') AS comment,
           COALESCE((
             SELECT json_agg(a.attname ORDER BY keys.ordinality)
             FROM pg_constraint pc
             CROSS JOIN LATERAL unnest(pc.conkey) WITH ORDINALITY AS keys(attnum, ordinality)
             JOIN pg_attribute a ON a.attrelid = pc.conrelid AND a.attnum = keys.attnum
             WHERE pc.conrelid = c.oid AND pc.contype = 'p'
           ), '[]'::json) AS primary_key
    FROM pg_class c
    JOIN pg_namespace n ON n.oid = c.relnamespace
    WHERE n.nspname IN ('public', 'auth', 'storage')
      AND c.relkind IN ('r', 'p')
      AND NOT c.relispartition
    ORDER BY n.nspname, c.relname`,
  columns: `
    SELECT n.nspname AS schema_name,
           c.relname AS table_name,
           a.attnum AS ordinal_position,
           a.attname AS column_name,
           tn.nspname AS type_schema,
           t.typname AS type_name,
           format_type(a.atttypid, a.atttypmod) AS formatted_type,
           a.attnotnull AS not_null,
           a.attidentity::text AS identity_kind,
           a.attgenerated::text AS generated_kind,
           coll.collname AS collation,
           pg_get_expr(ad.adbin, ad.adrelid) AS default_expression,
           col_description(a.attrelid, a.attnum) AS comment
    FROM pg_attribute a
    JOIN pg_class c ON c.oid = a.attrelid
    JOIN pg_namespace n ON n.oid = c.relnamespace
    JOIN pg_type t ON t.oid = a.atttypid
    JOIN pg_namespace tn ON tn.oid = t.typnamespace
    LEFT JOIN pg_attrdef ad ON ad.adrelid = a.attrelid AND ad.adnum = a.attnum
    LEFT JOIN pg_collation coll ON coll.oid = a.attcollation AND a.attcollation <> 0
    WHERE n.nspname IN ('public', 'auth', 'storage')
      AND c.relkind IN ('r', 'p', 'v', 'm')
      AND a.attnum > 0
      AND NOT a.attisdropped
    ORDER BY n.nspname, c.relname, a.attnum`,
  constraints: `
    SELECT n.nspname AS schema_name,
           c.relname AS table_name,
           pc.conname AS constraint_name,
           pc.contype::text AS constraint_type,
           pc.condeferrable AS deferrable,
           pc.condeferred AS initially_deferred,
           pc.convalidated AS validated,
           CASE WHEN pc.confrelid = 0 THEN NULL ELSE pc.confrelid::regclass::text END AS referenced_table,
           pg_get_constraintdef(pc.oid, true) AS definition
    FROM pg_constraint pc
    JOIN pg_class c ON c.oid = pc.conrelid
    JOIN pg_namespace n ON n.oid = c.relnamespace
    WHERE n.nspname IN ('public', 'auth', 'storage')
    ORDER BY n.nspname, c.relname, pc.conname`,
  indexes: `
    SELECT n.nspname AS schema_name,
           c.relname AS table_name,
           ic.relname AS index_name,
           i.indisprimary AS is_primary,
           i.indisunique AS is_unique,
           i.indisvalid AS is_valid,
           i.indisready AS is_ready,
           i.indisclustered AS is_clustered,
           am.amname AS access_method,
           pg_get_expr(i.indpred, i.indrelid) AS predicate,
           pg_get_indexdef(i.indexrelid) AS definition
    FROM pg_index i
    JOIN pg_class c ON c.oid = i.indrelid
    JOIN pg_namespace n ON n.oid = c.relnamespace
    JOIN pg_class ic ON ic.oid = i.indexrelid
    JOIN pg_am am ON am.oid = ic.relam
    WHERE n.nspname IN ('public', 'auth', 'storage')
    ORDER BY n.nspname, c.relname, ic.relname`,
  triggers: `
    SELECT n.nspname AS schema_name,
           c.relname AS table_name,
           tg.tgname AS trigger_name,
           tg.tgenabled::text AS enabled,
           tg.tgfoid::regprocedure::text AS function_name,
           pg_get_triggerdef(tg.oid, true) AS definition
    FROM pg_trigger tg
    JOIN pg_class c ON c.oid = tg.tgrelid
    JOIN pg_namespace n ON n.oid = c.relnamespace
    WHERE n.nspname IN ('public', 'auth', 'storage')
      AND NOT tg.tgisinternal
    ORDER BY n.nspname, c.relname, tg.tgname`,
  policies: `
    SELECT schemaname AS schema_name,
           tablename AS table_name,
           policyname AS policy_name,
           permissive,
           roles,
           cmd,
           qual,
           with_check
    FROM pg_policies
    WHERE schemaname IN ('public', 'auth', 'storage')
    ORDER BY schemaname, tablename, policyname`,
  functions: `
    SELECT n.nspname AS schema_name,
           p.proname AS function_name,
           p.prokind::text AS function_kind,
           pg_get_function_identity_arguments(p.oid) AS identity_arguments,
           pg_get_function_result(p.oid) AS result_type,
           l.lanname AS language,
           p.provolatile::text AS volatility,
           p.prosecdef AS security_definer,
           p.proleakproof AS leakproof,
           p.proconfig AS configuration,
           pg_get_userbyid(p.proowner) AS owner,
           obj_description(p.oid, 'pg_proc') AS comment,
           pg_get_functiondef(p.oid) AS definition
    FROM pg_proc p
    JOIN pg_namespace n ON n.oid = p.pronamespace
    JOIN pg_language l ON l.oid = p.prolang
    WHERE n.nspname IN ('public', 'auth', 'storage')
      AND p.prokind IN ('f', 'p')
    ORDER BY n.nspname, p.proname, pg_get_function_identity_arguments(p.oid)`,
  extensions: `
    SELECT e.extname AS extension_name,
           e.extversion AS version,
           n.nspname AS schema_name,
           e.extrelocatable AS relocatable,
           pg_get_userbyid(e.extowner) AS owner
    FROM pg_extension e
    JOIN pg_namespace n ON n.oid = e.extnamespace
    ORDER BY e.extname`,
  views: `
    SELECT n.nspname AS schema_name,
           c.relname AS view_name,
           c.relkind::text AS view_kind,
           pg_get_userbyid(c.relowner) AS owner,
           obj_description(c.oid, 'pg_class') AS comment,
           pg_get_viewdef(c.oid, true) AS definition
    FROM pg_class c
    JOIN pg_namespace n ON n.oid = c.relnamespace
    WHERE n.nspname IN ('public', 'auth', 'storage')
      AND c.relkind IN ('v', 'm')
    ORDER BY n.nspname, c.relname`,
  sequences: `
    SELECT schemaname AS schema_name,
           sequencename AS sequence_name,
           sequenceowner AS owner,
           data_type,
           start_value,
           min_value,
           max_value,
           increment_by,
           cycle,
           cache_size,
           last_value
    FROM pg_sequences
    WHERE schemaname IN ('public', 'auth', 'storage')
    ORDER BY schemaname, sequencename`,
  tablePrivileges: `
    SELECT grantor, grantee, table_schema AS schema_name, table_name, privilege_type, is_grantable
    FROM information_schema.role_table_grants
    WHERE table_schema IN ('public', 'auth', 'storage')
    ORDER BY table_schema, table_name, grantee, privilege_type, grantor`,
  routinePrivileges: `
    SELECT grantor, grantee, routine_schema AS schema_name, routine_name, specific_name, privilege_type, is_grantable
    FROM information_schema.role_routine_grants
    WHERE routine_schema IN ('public', 'auth', 'storage')
    ORDER BY routine_schema, routine_name, specific_name, grantee, privilege_type, grantor`,
  enumTypes: `
    SELECT n.nspname AS schema_name,
           t.typname AS type_name,
           json_agg(json_build_object('label', e.enumlabel, 'sortOrder', e.enumsortorder) ORDER BY e.enumsortorder) AS labels
    FROM pg_type t
    JOIN pg_namespace n ON n.oid = t.typnamespace
    JOIN pg_enum e ON e.enumtypid = t.oid
    WHERE n.nspname IN ('public', 'auth', 'storage')
    GROUP BY n.nspname, t.typname
    ORDER BY n.nspname, t.typname`,
  domains: `
    SELECT domain_schema AS schema_name,
           domain_name,
           data_type,
           udt_schema,
           udt_name,
           domain_default,
           character_maximum_length,
           numeric_precision,
           numeric_scale
    FROM information_schema.domains
    WHERE domain_schema IN ('public', 'auth', 'storage')
    ORDER BY domain_schema, domain_name`,
};

function printHelp(): void {
  console.log(`Usage: npx tsx backend/scripts/export-source-archive.ts [options]

Creates a read-only data and catalog archive of public, auth, and storage.

Options:
  --output-root <path>  Parent directory for timestamped archives
  --page-size <number>  Rows per table part (100-5000, default: 500)
  --help                Show this help`);
}

function parseOptions(arguments_: string[]): Options {
  const options: Options = { outputRoot: defaultOutputRoot, pageSize: 500, help: false };
  for (let index = 0; index < arguments_.length; index += 1) {
    const argument = arguments_[index];
    if (argument === "--help" || argument === "-h") {
      options.help = true;
    } else if (argument === "--output-root") {
      const value = arguments_[index + 1];
      if (!value) throw new Error("--output-root requires a path");
      options.outputRoot = path.resolve(value);
      index += 1;
    } else if (argument === "--page-size") {
      const value = Number(arguments_[index + 1]);
      if (!Number.isInteger(value) || value < 100 || value > 5000) {
        throw new Error("--page-size must be an integer between 100 and 5000");
      }
      options.pageSize = value;
      index += 1;
    } else {
      throw new Error(`Unknown argument: ${argument}`);
    }
  }
  return options;
}

function parseDotEnv(contents: string): Record<string, string> {
  const values: Record<string, string> = {};
  for (const rawLine of contents.replace(/^\uFEFF/, "").split(/\r?\n/)) {
    const line = rawLine.trim();
    if (!line || line.startsWith("#")) continue;
    const match = /^(?:export\s+)?([A-Za-z_][A-Za-z0-9_]*)\s*=\s*(.*)$/.exec(line);
    const key = match?.[1];
    const rawValue = match?.[2];
    if (key === undefined || rawValue === undefined) continue;
    let value = rawValue.trim();
    if (value.length >= 2 && ((value.startsWith('"') && value.endsWith('"')) ||
      (value.startsWith("'") && value.endsWith("'")))) {
      value = value.slice(1, -1);
    }
    values[key] = value;
  }
  return values;
}

async function loadEnvironment(): Promise<void> {
  try {
    const values = parseDotEnv(await readFile(envPath, "utf8"));
    for (const [key, value] of Object.entries(values)) {
      if (process.env[key] === undefined) process.env[key] = value;
    }
  } catch (error) {
    if ((error as NodeJS.ErrnoException).code !== "ENOENT") throw error;
  }
}

function requiredEnvironment(name: string, aliases: string[] = []): string {
  for (const key of [name, ...aliases]) {
    const value = process.env[key]?.trim();
    if (value) return value;
  }
  throw new Error(`Missing required environment variable: ${name}`);
}

function assertReadOnly(sql: string): void {
  const normalized = sql.trim();
  if (!/^(SELECT|WITH)\b/i.test(normalized) || normalized.includes(";")) {
    throw new Error("Refusing a source query that is not one SELECT statement");
  }
  if (/\b(ALTER|CALL|COPY|CREATE|DELETE|DO|DROP|GRANT|INSERT|REINDEX|REVOKE|TRUNCATE|UPDATE|VACUUM)\b/i.test(normalized)) {
    throw new Error("Refusing a source query containing a write operation");
  }
}

function rowsFromResponse(payload: unknown): Row[] {
  if (Array.isArray(payload)) return payload as Row[];
  if (payload && typeof payload === "object") {
    const record = payload as Record<string, unknown>;
    for (const key of ["data", "result", "rows"]) {
      if (Array.isArray(record[key])) return record[key] as Row[];
    }
  }
  throw new Error("Management API response did not contain query rows");
}

function sleep(milliseconds: number): Promise<void> {
  return new Promise((resolve) => setTimeout(resolve, milliseconds));
}

function safeApiMessage(payload: unknown): string {
  if (!payload || typeof payload !== "object") return "No structured details";
  const record = payload as Record<string, unknown>;
  const message = record.message ?? record.error ?? record.msg;
  if (typeof message !== "string") return "No structured details";
  return message
    .replace(/Bearer\s+\S+/gi, "Bearer [REDACTED]")
    .replace(/(token|password|secret)\s*[:=]\s*\S+/gi, "$1=[REDACTED]")
    .slice(0, 500);
}

function createQuery(accessToken: string, projectRef: string): Query {
  const endpoint = `${managementApiBaseUrl}/projects/${encodeURIComponent(projectRef)}/database/query`;
  return async (sql: string, parameters: unknown[] = []): Promise<Row[]> => {
    assertReadOnly(sql);
    for (let attempt = 1; attempt <= 5; attempt += 1) {
      try {
        const response = await fetch(endpoint, {
          method: "POST",
          headers: { Authorization: `Bearer ${accessToken}`, "Content-Type": "application/json" },
          body: JSON.stringify({ query: sql, parameters, read_only: true }),
          signal: AbortSignal.timeout(180_000),
        });
        const responseText = await response.text();
        let payload: unknown = null;
        try {
          payload = responseText ? JSON.parse(responseText) : null;
        } catch {
          payload = null;
        }
        if (response.ok) return rowsFromResponse(payload);

        const retryable = [408, 425, 429].includes(response.status) || response.status >= 500;
        if (!retryable || attempt === 5) {
          const requestId = response.headers.get("x-request-id");
          throw new Error(
            `Supabase Management API HTTP ${response.status}${requestId ? ` (${requestId})` : ""}: ${safeApiMessage(payload)}`,
          );
        }
      } catch (error) {
        if (attempt === 5) throw error;
        if (error instanceof Error && error.message.startsWith("Supabase Management API HTTP 4")) throw error;
      }
      await sleep(500 * (2 ** (attempt - 1)) + Math.floor(Math.random() * 250));
    }
    throw new Error("Management API retry loop ended unexpectedly");
  };
}

function quoteIdentifier(value: string): string {
  if (!value || value.includes("\0")) throw new Error("Invalid PostgreSQL identifier");
  return `"${value.replaceAll('"', '""')}"`;
}

function quoteLiteral(value: string): string {
  if (value.includes("\0")) throw new Error("Invalid PostgreSQL literal");
  return `'${value.replaceAll("'", "''")}'`;
}

function tableKey(schema: string, table: string): string {
  return JSON.stringify([schema, table]);
}

function datasetName(schema: string, table: string): string {
  return `table:${schema}.${table}`;
}

function tableOutputDirectory(schema: string, table: string): string {
  return path.posix.join("tables", schema, Buffer.from(table, "utf8").toString("base64url"));
}

function qualifiedName(schema: string, table: string): string {
  return `${quoteIdentifier(schema)}.${quoteIdentifier(table)}`;
}

function hash(contents: Uint8Array): string {
  return createHash("sha256").update(contents).digest("hex");
}

async function writePrivate(filePath: string, contents: Uint8Array): Promise<void> {
  await mkdir(path.dirname(filePath), { recursive: true, mode: 0o700 });
  const temporaryPath = `${filePath}.tmp-${process.pid}`;
  await writeFile(temporaryPath, contents, { flag: "wx", mode: 0o600 });
  await rename(temporaryPath, filePath);
  await chmod(filePath, 0o600).catch(() => undefined);
}

async function writeJsonArtifact(
  archiveDirectory: string,
  relativePath: string,
  payload: unknown,
  file: Omit<ArchiveFile, "path" | "bytes" | "sha256">,
): Promise<ArchiveFile> {
  const contents = Buffer.from(`${JSON.stringify(payload)}\n`, "utf8");
  await writePrivate(path.join(archiveDirectory, ...relativePath.split("/")), contents);
  return { ...file, path: relativePath, bytes: contents.byteLength, sha256: hash(contents) };
}

async function captureCatalog(query: Query): Promise<Record<string, Row[]>> {
  const catalog: Record<string, Row[]> = {};
  for (const [name, sql] of Object.entries(catalogQueries)) catalog[name] = await query(sql);
  return catalog;
}

function catalogHash(catalog: Record<string, Row[]>): string {
  return hash(Buffer.from(JSON.stringify(catalog), "utf8"));
}

function parseString(value: unknown, label: string): string {
  if (typeof value !== "string" || !value) throw new Error(`Invalid ${label} from catalog`);
  return value;
}

function parsePrimaryKey(value: unknown): string[] {
  if (!Array.isArray(value) || !value.every((item) => typeof item === "string" && item.length > 0)) {
    throw new Error("Invalid primary key metadata");
  }
  return value as string[];
}

function tableDefinitions(catalog: Record<string, Row[]>): TableDefinition[] {
  const tableRows = catalog.tables;
  const columnRows = catalog.columns;
  if (!tableRows || !columnRows) throw new Error("Catalog lacks table or column metadata");

  const columnsByTable = new Map<string, string[]>();
  for (const row of columnRows) {
    const schema = parseString(row.schema_name, "column schema");
    const table = parseString(row.table_name, "column table");
    const column = parseString(row.column_name, "column name");
    const key = tableKey(schema, table);
    const columns = columnsByTable.get(key) ?? [];
    columns.push(column);
    columnsByTable.set(key, columns);
  }

  return tableRows.map((row) => {
    const schema = parseString(row.schema_name, "table schema");
    const table = parseString(row.table_name, "table name");
    if (!allowedSchemas.includes(schema as typeof allowedSchemas[number])) {
      throw new Error(`Unexpected source schema: ${schema}`);
    }
    const columns = columnsByTable.get(tableKey(schema, table));
    if (!columns || columns.length === 0) throw new Error(`No columns discovered for ${schema}.${table}`);
    return {
      schema,
      table,
      relationKind: parseString(row.relation_kind, "relation kind"),
      columns,
      primaryKey: parsePrimaryKey(row.primary_key),
    };
  });
}

function buildStateQuery(tables: TableDefinition[]): string {
  if (tables.length === 0) throw new Error("No source tables discovered");
  const selections = tables.map((table) => `
    SELECT ${quoteLiteral(table.schema)}::text AS schema_name,
           ${quoteLiteral(table.table)}::text AS table_name,
           count(*)::text AS row_count,
           COALESCE(sum(hashtextextended(row_to_json(source_row)::text, 0)::numeric), 0)::text AS fingerprint_a,
           COALESCE(sum(hashtextextended(row_to_json(source_row)::text, 1)::numeric), 0)::text AS fingerprint_b
    FROM ${qualifiedName(table.schema, table.table)} AS source_row`);
  return `SELECT * FROM (${selections.join(" UNION ALL ")}) AS table_states ORDER BY schema_name, table_name`;
}

function safeCount(value: unknown, label: string): number {
  const count = typeof value === "number" ? value : Number(value);
  if (!Number.isSafeInteger(count) || count < 0) throw new Error(`Invalid count for ${label}`);
  return count;
}

async function captureStates(query: Query, tables: TableDefinition[]): Promise<Map<string, TableState>> {
  const rows = await query(buildStateQuery(tables));
  const states = new Map<string, TableState>();
  for (const row of rows) {
    const schema = parseString(row.schema_name, "state schema");
    const table = parseString(row.table_name, "state table");
    const state: TableState = {
      schema,
      table,
      count: safeCount(row.row_count, `${schema}.${table}`),
      fingerprintA: String(row.fingerprint_a),
      fingerprintB: String(row.fingerprint_b),
    };
    states.set(tableKey(schema, table), state);
  }
  if (states.size !== tables.length) throw new Error("Source state did not cover every discovered table");
  return states;
}

function stateFor(states: Map<string, TableState>, table: TableDefinition): TableState {
  const state = states.get(tableKey(table.schema, table.table));
  if (!state) throw new Error(`Missing state for ${table.schema}.${table.table}`);
  return state;
}

function buildSnapshotPageQuery(definitions: TableDefinition[], pageSize: number): string {
  const selections = definitions.map((definition, tableOrder) => {
    if (definition.primaryKey.length === 0) {
      throw new Error(`No primary key available for snapshot pagination of ${definition.schema}.${definition.table}`);
    }
    const keys = definition.primaryKey.map(quoteIdentifier);
    return `
      SELECT ${tableOrder}::integer AS archive_table_order,
             ${quoteLiteral(datasetName(definition.schema, definition.table))}::text AS archive_dataset,
             (((archive_position - 1) / ${pageSize}) + 1)::integer AS archive_page,
             json_agg(archive_row ORDER BY archive_position) AS archive_rows
      FROM (
        SELECT row_number() OVER (ORDER BY ${keys.join(", ")}) AS archive_position,
               row_to_json(source_row) AS archive_row
        FROM ${qualifiedName(definition.schema, definition.table)} AS source_row
      ) AS numbered_rows
      GROUP BY (((archive_position - 1) / ${pageSize}) + 1)::integer`;
  });
  return `SELECT * FROM (${selections.join(" UNION ALL ")}) AS archive_pages
          ORDER BY archive_table_order, archive_page`;
}

function archiveRows(value: unknown, dataset: string): Row[] {
  let parsed = value;
  if (typeof parsed === "string") {
    try {
      parsed = JSON.parse(parsed);
    } catch {
      throw new Error(`Invalid row JSON returned for ${dataset}`);
    }
  }
  if (!Array.isArray(parsed) || !parsed.every((row) => row && typeof row === "object" && !Array.isArray(row))) {
    throw new Error(`Invalid row array returned for ${dataset}`);
  }
  return parsed as Row[];
}

async function exportTablesBatched(
  query: Query,
  archiveDirectory: string,
  definitions: TableDefinition[],
  states: Map<string, TableState>,
  pageSize: number,
  files: ArchiveFile[],
): Promise<Map<string, { rows: number; pages: number }>> {
  for (const definition of definitions) {
    if (definition.primaryKey.length === 0) {
      throw new Error(
        `${definition.schema}.${definition.table} has no primary key; refusing an archive that cannot paginate deterministically`,
      );
    }
  }

  const definitionByDataset = new Map(
    definitions.map((definition) => [datasetName(definition.schema, definition.table), definition]),
  );
  const results = new Map<string, { rows: number; pages: number }>();
  for (const definition of definitions) {
    results.set(tableKey(definition.schema, definition.table), { rows: 0, pages: 0 });
  }

  const responseRows = await query(buildSnapshotPageQuery(definitions, pageSize));
  for (const responseRow of responseRows) {
    const dataset = parseString(responseRow.archive_dataset, "archive dataset");
    const definition = definitionByDataset.get(dataset);
    if (!definition) throw new Error(`Unexpected archive dataset returned: ${dataset}`);
    const page = safeCount(responseRow.archive_page, `${dataset} page`);
    if (page < 1) throw new Error(`Invalid page returned for ${dataset}`);
    const dataRows = archiveRows(responseRow.archive_rows, dataset);
    if (dataRows.length === 0 || dataRows.length > pageSize) throw new Error(`Invalid page size for ${dataset}`);
    const result = results.get(tableKey(definition.schema, definition.table));
    if (!result) throw new Error(`Missing result accumulator for ${dataset}`);
    if (page !== result.pages + 1) throw new Error(`Non-contiguous source pages for ${dataset}`);
    result.pages = page;
    result.rows += dataRows.length;
    const relativePath = path.posix.join(
      tableOutputDirectory(definition.schema, definition.table),
      `part-${String(page).padStart(6, "0")}.json`,
    );
    files.push(await writeJsonArtifact(
      archiveDirectory,
      relativePath,
      {
        format: "tonline-source-archive-table-page",
        formatVersion,
        dataset,
        schema: definition.schema,
        table: definition.table,
        page,
        rowCount: dataRows.length,
        rows: dataRows,
      },
      { kind: "table", dataset, page, rowCount: dataRows.length },
    ));
  }

  for (const definition of definitions) {
    const result = results.get(tableKey(definition.schema, definition.table));
    if (!result) throw new Error(`Missing result for ${definition.schema}.${definition.table}`);
    const expectedRows = stateFor(states, definition).count;
    if (result.rows !== expectedRows) {
      throw new Error(
        `${definition.schema}.${definition.table} exported ${result.rows} rows but the initial state had ${expectedRows}`,
      );
    }
    if (result.rows > 0) {
      console.log(`${definition.schema}.${definition.table}: ${result.rows} rows in ${result.pages} part(s)`);
    }
  }
  return results;
}

function sameState(left: TableState, right: TableState): boolean {
  return left.count === right.count && left.fingerprintA === right.fingerprintA &&
    left.fingerprintB === right.fingerprintB;
}

function timestamp(date: Date): string {
  return date.toISOString().replace(/[:.]/g, "-");
}

async function main(): Promise<void> {
  const options = parseOptions(process.argv.slice(2));
  if (options.help) {
    printHelp();
    return;
  }
  await loadEnvironment();
  const accessToken = requiredEnvironment("SUPABASE_ACCESS_TOKEN");
  const projectRef = requiredEnvironment("SUPABASE_SOURCE_PROJECT_REF", ["SUPABASE_PROJECT_REF"]);
  if (!/^[a-z0-9]+$/i.test(projectRef)) throw new Error("Invalid Supabase project reference");

  const startedAt = new Date();
  const archiveDirectory = path.join(options.outputRoot, timestamp(startedAt));
  await mkdir(options.outputRoot, { recursive: true, mode: 0o700 });
  await mkdir(archiveDirectory, { recursive: false, mode: 0o700 });
  await chmod(archiveDirectory, 0o700).catch(() => undefined);

  const query = createQuery(accessToken, projectRef);
  console.log("Starting complete read-only source archive; vault and Storage binaries are excluded.");
  const catalogBefore = await captureCatalog(query);
  const definitions = tableDefinitions(catalogBefore);
  const statesBefore = await captureStates(query, definitions);
  const files: ArchiveFile[] = [];
  const metadataManifest: Array<{ name: string; rowCount: number; path: string }> = [];

  for (const [name, rows] of Object.entries(catalogBefore)) {
    const relativePath = path.posix.join("metadata", `${name}.json`);
    files.push(await writeJsonArtifact(
      archiveDirectory,
      relativePath,
      { format: "tonline-source-archive-metadata", formatVersion, name, rowCount: rows.length, rows },
      { kind: "metadata", dataset: `metadata:${name}`, page: 1, rowCount: rows.length },
    ));
    metadataManifest.push({ name, rowCount: rows.length, path: relativePath });
  }

  const tableExports = await exportTablesBatched(
    query,
    archiveDirectory,
    definitions,
    statesBefore,
    options.pageSize,
    files,
  );

  const statesAfter = await captureStates(query, definitions);
  const catalogAfter = await captureCatalog(query);
  const catalogBeforeHash = catalogHash(catalogBefore);
  const catalogAfterHash = catalogHash(catalogAfter);
  if (catalogBeforeHash !== catalogAfterHash) {
    throw new Error("Source catalog changed during archive; archive refused and has no manifest");
  }

  const tables: TableManifest[] = [];
  for (const definition of definitions) {
    const before = stateFor(statesBefore, definition);
    const after = stateFor(statesAfter, definition);
    if (!sameState(before, after)) {
      throw new Error(
        `${definition.schema}.${definition.table} mutated during archive; archive refused and has no manifest`,
      );
    }
    const exported = tableExports.get(tableKey(definition.schema, definition.table));
    if (!exported) throw new Error(`Missing export result for ${definition.schema}.${definition.table}`);
    tables.push({
      dataset: datasetName(definition.schema, definition.table),
      schema: definition.schema,
      table: definition.table,
      relationKind: definition.relationKind,
      columns: definition.columns,
      primaryKey: definition.primaryKey,
      pagination: definition.primaryKey.length > 0 ? "primary-key" : "physical-cursor",
      rowCountBefore: before.count,
      rowCountAfter: after.count,
      fingerprintBefore: [before.fingerprintA, before.fingerprintB],
      fingerprintAfter: [after.fingerprintA, after.fingerprintB],
      exportedRows: exported.rows,
      pages: exported.pages,
    });
  }

  const completedAt = new Date();
  const manifest = {
    format: "tonline-supabase-source-archive",
    formatVersion,
    createdAt: startedAt.toISOString(),
    completedAt: completedAt.toISOString(),
    source: {
      provider: "supabase",
      projectRef,
      schemas: [...allowedSchemas],
      queryMode: "management-api-read-only",
    },
    archive: {
      pageSize: options.pageSize,
      tableCount: tables.length,
      nonEmptyTableCount: tables.filter((table) => table.exportedRows > 0).length,
      tableRowCount: tables.reduce((sum, table) => sum + table.exportedRows, 0),
      metadataRowCount: metadataManifest.reduce((sum, item) => sum + item.rowCount, 0),
      fileCount: files.length,
      containsSensitiveAuthData: true,
      storageBinariesDownloaded: false,
      transactionalSnapshot: false,
      mutationDetection: "global before/after counts plus two independent row fingerprints and catalog SHA-256",
      excludedSchemas: ["vault"],
      excludedRelations: ["vault.decrypted_secrets"],
    },
    catalogSha256: catalogBeforeHash,
    metadata: metadataManifest,
    tables,
    files,
  };
  const manifestContents = Buffer.from(`${JSON.stringify(manifest, null, 2)}\n`, "utf8");
  const manifestHash = hash(manifestContents);
  await writePrivate(path.join(archiveDirectory, "manifest.json"), manifestContents);
  await writePrivate(
    path.join(archiveDirectory, "manifest.sha256"),
    Buffer.from(`${manifestHash}  manifest.json\n`, "ascii"),
  );

  console.log(`Archive complete: ${archiveDirectory}`);
  console.log(`${tables.length} tables; ${manifest.archive.tableRowCount} data rows; ${files.length} files`);
  console.log(`Manifest SHA-256: ${manifestHash}`);
  console.log("No source writes, vault rows, decrypted secrets, or Storage binaries were requested.");
}

main().catch((error: unknown) => {
  const message = error instanceof Error ? error.message : "Unknown archive error";
  console.error(`Archive failed: ${message}`);
  if (process.env.ARCHIVE_DEBUG === "1" && error instanceof Error && error.stack) {
    console.error(error.stack);
  }
  process.exitCode = 1;
});
