From c0eafc2fcee8113eb68b4af2966fd59b9207b6d1 Mon Sep 17 00:00:00 2001 From: Drew Donaldson <49219012+Automata-intelligentsia@users.noreply.github.com> Date: Sun, 19 Jul 2026 13:15:30 -0400 Subject: [PATCH] fix(postgres): widen overflow-prone token/usage counters to bigint (#2331) SQLite INTEGER is effectively int64, but the PostgreSQL baseline mapped unbounded token/usage counters on `project.tasks` and `project.chat_token_usage` to `integer` (int4). Real data contains values > 2,147,483,647, causing the SQLite-to-PostgreSQL migration to fail with `value ... is out of range for type integer`. Changes: - Change baseline DDL to `bigint` for the affected columns. - Update Drizzle schema to `bigint({ mode: "number" })` to preserve JS `number` semantics. - Add forward migration `0024_bigint_counters.sql` for existing clusters. - Bump `SCHEMA_BASELINE_VERSION` to `0024`. Fixes the int4 overflow observed during migration of large token/usage counters. ## Summary by CodeRabbit * **Bug Fixes** * Expanded token-usage and activity/lease counters to 64-bit integers to prevent overflow on large workloads. * Improved distributed task ID state/reservations to be isolated per project and to merge/update conflicting entries more reliably. * **Chores** * Added an idempotent PostgreSQL migration for bigint counter support and advanced schema baseline tracking. * Updated dashboard build support by adding `html2canvas` type definitions and the production dependency. * **Tests** * Updated schema-applier migration checks to include the new bigint-counters baseline identity. --------- Co-authored-by: gsxdsm --- .../__tests__/postgres/schema-applier.test.ts | 26 + .../src/postgres/migrations/0000_initial.sql | 42 +- .../migrations/0026_bigint_counters.sql | 46 + packages/core/src/postgres/schema-applier.ts | 1473 +++++++++-------- packages/core/src/postgres/schema/project.ts | 24 +- packages/core/src/postgres/sqlite-migrator.ts | 148 +- 6 files changed, 999 insertions(+), 760 deletions(-) create mode 100644 packages/core/src/postgres/migrations/0026_bigint_counters.sql diff --git a/packages/core/src/__tests__/postgres/schema-applier.test.ts b/packages/core/src/__tests__/postgres/schema-applier.test.ts index c15c3c590c..318cd64907 100644 --- a/packages/core/src/__tests__/postgres/schema-applier.test.ts +++ b/packages/core/src/__tests__/postgres/schema-applier.test.ts @@ -66,6 +66,7 @@ import { SESSION_ADVISOR_ENABLED_SCHEMA_VERSION, SQLITE_SCHEMA_PARITY_VERSION, SYMBOL_LOCKS_SCHEMA_VERSION, + BIGINT_COUNTERS_VERSION, TASK_VERIFICATION_REQUEST_VERSION, } from "../../postgres/schema-applier.js"; import { rekeyFallbackProjectPartition } from "../../postgres/migration-stamping.js"; @@ -178,6 +179,26 @@ describe("schema-applier: immutable migration identities", () => { expect(Number(SCHEMA_BASELINE_VERSION)).toBeGreaterThanOrEqual(Number(SYMBOL_LOCKS_SCHEMA_VERSION)); }); + /* + FNXC:PostgresBigintCounters 2026-07-19-12:00: + 0026 widens overflow-prone counters to bigint. Keep identity fixed and at-or-before SCHEMA_BASELINE_VERSION. + + FNXC:PostgresBigintCounters 2026-07-19-08:40: + Also assert the authoritative applier registry wires 0026_bigint_counters.sql — + constant identity alone does not prove applySchemaBaseline will run the migration. + */ + it("registers bigint counters at migration version 0026", () => { + expect(BIGINT_COUNTERS_VERSION).toBe("0026"); + expect(Number(SCHEMA_BASELINE_VERSION)).toBeGreaterThanOrEqual(Number(BIGINT_COUNTERS_VERSION)); + const applierSource = readFileSync( + fileURLToPath(new URL("../../postgres/schema-applier.ts", import.meta.url)), + "utf8", + ); + expect(applierSource).toContain("0026_bigint_counters.sql"); + expect(applierSource).toContain("BIGINT_COUNTERS_VERSION"); + expect(applierSource).toMatch(/applied\.includes\(\s*BIGINT_COUNTERS_VERSION\s*\)/); + }); + }); /* @@ -1272,6 +1293,7 @@ pgDescribe("schema-applier: automation project-isolation upgrade", () => { RESEARCH_FEATURE_PROVENANCE_VERSION, TASK_VERIFICATION_REQUEST_VERSION, SYMBOL_LOCKS_SCHEMA_VERSION, + BIGINT_COUNTERS_VERSION, ]); expect((await applySchemaBaseline(ctx.db, { pluginHooks: [] })).applied).toBe(false); }); @@ -1323,6 +1345,7 @@ pgDescribe("schema-applier: automation project-isolation upgrade", () => { RESEARCH_FEATURE_PROVENANCE_VERSION, TASK_VERIFICATION_REQUEST_VERSION, SYMBOL_LOCKS_SCHEMA_VERSION, + BIGINT_COUNTERS_VERSION, ]); }); @@ -1507,6 +1530,7 @@ pgDescribe("schema-applier: automation project-isolation upgrade", () => { RESEARCH_FEATURE_PROVENANCE_VERSION, TASK_VERIFICATION_REQUEST_VERSION, SYMBOL_LOCKS_SCHEMA_VERSION, + BIGINT_COUNTERS_VERSION, ]); }); @@ -1572,6 +1596,7 @@ pgDescribe("schema-applier: automation project-isolation upgrade", () => { RESEARCH_FEATURE_PROVENANCE_VERSION, TASK_VERIFICATION_REQUEST_VERSION, SYMBOL_LOCKS_SCHEMA_VERSION, + BIGINT_COUNTERS_VERSION, ]); }); @@ -1637,6 +1662,7 @@ pgDescribe("schema-applier: automation project-isolation upgrade", () => { RESEARCH_FEATURE_PROVENANCE_VERSION, TASK_VERIFICATION_REQUEST_VERSION, SYMBOL_LOCKS_SCHEMA_VERSION, + BIGINT_COUNTERS_VERSION, ]); }); }); diff --git a/packages/core/src/postgres/migrations/0000_initial.sql b/packages/core/src/postgres/migrations/0000_initial.sql index dae8fc52e2..ac640340a0 100644 --- a/packages/core/src/postgres/migrations/0000_initial.sql +++ b/packages/core/src/postgres/migrations/0000_initial.sql @@ -103,11 +103,11 @@ CREATE TABLE IF NOT EXISTS project.tasks ( validator_thinking_level text, planning_thinking_level text, execution_mode text DEFAULT 'standard', - token_usage_input_tokens integer, - token_usage_output_tokens integer, - token_usage_cached_tokens integer, - token_usage_cache_write_tokens integer, - token_usage_total_tokens integer, + token_usage_input_tokens bigint, + token_usage_output_tokens bigint, + token_usage_cached_tokens bigint, + token_usage_cache_write_tokens bigint, + token_usage_total_tokens bigint, token_usage_first_used_at text, token_usage_last_used_at text, token_usage_model_provider text, @@ -120,7 +120,7 @@ CREATE TABLE IF NOT EXISTS project.tasks ( updated_at text NOT NULL, column_moved_at text, first_execution_at text, - cumulative_active_ms integer, + cumulative_active_ms bigint, execution_started_at text, execution_completed_at text, dependencies jsonb DEFAULT '[]', @@ -174,7 +174,7 @@ CREATE TABLE IF NOT EXISTS project.tasks ( checkout_node_id text, checkout_run_id text, checkout_lease_renewed_at text, - checkout_lease_epoch integer DEFAULT 0, + checkout_lease_epoch bigint DEFAULT 0, deleted_at text, allow_resurrection integer DEFAULT 0, transition_pending text, @@ -209,15 +209,18 @@ CREATE TABLE IF NOT EXISTS project.config ( ); CREATE TABLE IF NOT EXISTS project.distributed_task_id_state ( - prefix text PRIMARY KEY, + project_id text NOT NULL DEFAULT current_setting('fusion.project_id', true), + prefix text NOT NULL, next_sequence integer NOT NULL, committed_cluster_task_count integer NOT NULL, last_committed_task_id text, - updated_at text NOT NULL + updated_at text NOT NULL, + PRIMARY KEY (project_id, prefix) ); CREATE TABLE IF NOT EXISTS project.distributed_task_id_reservations ( - reservation_id text PRIMARY KEY, + project_id text NOT NULL DEFAULT current_setting('fusion.project_id', true), + reservation_id text NOT NULL, prefix text NOT NULL, node_id text NOT NULL, sequence integer NOT NULL, @@ -229,17 +232,18 @@ CREATE TABLE IF NOT EXISTS project.distributed_task_id_reservations ( aborted_at text, created_at text NOT NULL, updated_at text NOT NULL, + PRIMARY KEY (project_id, reservation_id), CONSTRAINT distributed_task_id_reservations_prefix_fkey - FOREIGN KEY (prefix) REFERENCES project.distributed_task_id_state(prefix) ON DELETE CASCADE, + FOREIGN KEY (project_id, prefix) REFERENCES project.distributed_task_id_state(project_id, prefix) ON DELETE CASCADE, CONSTRAINT distributed_task_id_reservations_status_check CHECK (status IN ('reserved', 'committed', 'aborted', 'expired')), CONSTRAINT distributed_task_id_reservations_reason_check CHECK (reason IS NULL OR reason IN ('abort', 'expired', 'failed-create')), - CONSTRAINT distributed_task_id_reservations_prefix_sequence_unique UNIQUE (prefix, sequence), - CONSTRAINT distributed_task_id_reservations_prefix_task_id_unique UNIQUE (prefix, task_id) + CONSTRAINT distributed_task_id_reservations_prefix_sequence_unique UNIQUE (project_id, prefix, sequence), + CONSTRAINT distributed_task_id_reservations_prefix_task_id_unique UNIQUE (project_id, prefix, task_id) ); CREATE INDEX IF NOT EXISTS "idxDistributedTaskIdReservationsPrefixStatus" - ON project.distributed_task_id_reservations(prefix, status); + ON project.distributed_task_id_reservations(project_id, prefix, status); CREATE INDEX IF NOT EXISTS "idxDistributedTaskIdReservationsExpiry" ON project.distributed_task_id_reservations(status, expires_at); @@ -1391,11 +1395,11 @@ CREATE TABLE IF NOT EXISTS project.chat_token_usage ( agent_id text, model_provider text, model_id text, - input_tokens integer NOT NULL, - output_tokens integer NOT NULL, - cached_tokens integer NOT NULL, - cache_write_tokens integer NOT NULL, - total_tokens integer NOT NULL, + input_tokens bigint NOT NULL, + output_tokens bigint NOT NULL, + cached_tokens bigint NOT NULL, + cache_write_tokens bigint NOT NULL, + total_tokens bigint NOT NULL, created_at text NOT NULL ); diff --git a/packages/core/src/postgres/migrations/0026_bigint_counters.sql b/packages/core/src/postgres/migrations/0026_bigint_counters.sql new file mode 100644 index 0000000000..faa8d451f2 --- /dev/null +++ b/packages/core/src/postgres/migrations/0026_bigint_counters.sql @@ -0,0 +1,46 @@ +/* +FNXC:PostgresBigintCounters 2026-07-18-21:45: +SQLite INTEGER is a 1-8 byte signed integer (effectively int64), but the PostgreSQL +baseline mapped several open-ended counters to integer (int4). Real task data +contains values that exceed 2,147,483,647 (e.g. cached token counts and cumulative +active millisecond timers), which caused the SQLite-to-PostgreSQL migration to fail +with "value ... is out of range for type integer". Upgrade the affected columns to +bigint without changing nullability or defaults. + +Affected columns: + project.tasks.token_usage_input_tokens + project.tasks.token_usage_output_tokens + project.tasks.token_usage_cached_tokens + project.tasks.token_usage_cache_write_tokens + project.tasks.token_usage_total_tokens + project.tasks.cumulative_active_ms + project.tasks.checkout_lease_epoch + project.chat_token_usage.input_tokens + project.chat_token_usage.output_tokens + project.chat_token_usage.cached_tokens + project.chat_token_usage.cache_write_tokens + project.chat_token_usage.total_tokens +*/ +DO $$ +BEGIN + IF to_regclass('project.tasks') IS NOT NULL THEN + ALTER TABLE project.tasks + ALTER COLUMN token_usage_input_tokens TYPE bigint, + ALTER COLUMN token_usage_output_tokens TYPE bigint, + ALTER COLUMN token_usage_cached_tokens TYPE bigint, + ALTER COLUMN token_usage_cache_write_tokens TYPE bigint, + ALTER COLUMN token_usage_total_tokens TYPE bigint, + ALTER COLUMN cumulative_active_ms TYPE bigint, + ALTER COLUMN checkout_lease_epoch TYPE bigint; + END IF; + + IF to_regclass('project.chat_token_usage') IS NOT NULL THEN + ALTER TABLE project.chat_token_usage + ALTER COLUMN input_tokens TYPE bigint, + ALTER COLUMN output_tokens TYPE bigint, + ALTER COLUMN cached_tokens TYPE bigint, + ALTER COLUMN cache_write_tokens TYPE bigint, + ALTER COLUMN total_tokens TYPE bigint; + END IF; +END +$$; diff --git a/packages/core/src/postgres/schema-applier.ts b/packages/core/src/postgres/schema-applier.ts index 579c4c0535..ff53b5469e 100644 --- a/packages/core/src/postgres/schema-applier.ts +++ b/packages/core/src/postgres/schema-applier.ts @@ -1,725 +1,748 @@ -/** - * PostgreSQL schema applier. - * - * FNXC:PostgresSchema 2026-06-24-03:40: - * Applies the fresh Drizzle migration baseline to a PostgreSQL connection - * and records it in a migration bookkeeping table. The baseline migration - * (migrations/0000_initial.sql) is the snapshot of the final SQLite schema - * (SCHEMA_VERSION=128) translated to PostgreSQL — applying it to an empty - * database yields final-schema parity (VAL-SCHEMA-001). - * - * After the baseline lands, plugin-owned tables are materialized via the - * schema-init hook (VAL-SCHEMA-007). The applier calls each registered plugin - * hook so plugins evolve their own tables independently of the core migration. - * - * Migration tracking uses a single-row bookkeeping table in the public schema - * so the applier is idempotent: re-running against an already-migrated database - * is a no-op. The version-gate discipline (the institutional learning that - * fresh-DB tests cannot catch a skipped-on-upgrade migration) is carried - * forward via the applier's explicit baseline marker. - */ - -import { readFile } from "node:fs/promises"; -import { existsSync } from "node:fs"; -import { dirname, join } from "node:path"; -import { fileURLToPath } from "node:url"; -import type { PostgresJsDatabase } from "drizzle-orm/postgres-js"; -import { sql } from "drizzle-orm"; -import { runPluginSchemaInitHooks, DEFAULT_PLUGIN_SCHEMA_INIT_HOOKS, type PluginSchemaInitHook } from "./plugin-schema-hook.js"; -import { acquireSchemaMutationLocks } from "./advisory-locks.js"; - -/** The latest PostgreSQL schema version known to this applier. */ -/* -FNXC:GitHubImportTranslate 2026-07-17-23:48: -Advances to 0019 for the import-translation legacy-partition backfill. Per-migration identities above stay fixed; only this latest-version marker moves. -*/ -export const SCHEMA_BASELINE_VERSION = "0025"; -const INITIAL_SCHEMA_VERSION = "0000"; -const AUTOMATION_ISOLATION_SCHEMA_VERSION = "0001"; -const ANALYTICS_ISOLATION_SCHEMA_VERSION = "0002"; -/** - * FNXC:PostgresMigrationIdentity 2026-07-14-01:41: - * Each migration keeps an immutable bookkeeping identity even as SCHEMA_BASELINE_VERSION advances to newer migrations. Upgrade checks and inserts must use this dedicated 0003 identifier so a later latest-version marker cannot make an unrecorded monitor/approval migration look applied. - */ -export const MONITOR_APPROVAL_ISOLATION_SCHEMA_VERSION = "0003"; -export const LEGACY_CUTOVER_PRESERVATION_SCHEMA_VERSION = "0004"; -export const MULTI_PROJECT_CUTOVER_SCHEMA_VERSION = "0005"; -export const PROJECT_OWNERSHIP_SCHEMA_VERSION = "0006"; -export const SQLITE_SCHEMA_PARITY_VERSION = "0007"; -/** - * FNXC:PlannerOversight 2026-07-14-18:49: - * Version 0008 adds project.tasks.session_advisor_enabled for per-task session - * advisor overrides. Keep this identity fixed when SCHEMA_BASELINE_VERSION advances. - */ -export const SESSION_ADVISOR_ENABLED_SCHEMA_VERSION = "0008"; -export const MISSION_FIX_IDEMPOTENCY_VERSION = "0009"; -/* -FNXC:GitHubImportTranslate 2026-07-15-09:30: -Import-translation cache advances to 0010. Migrations are registered here explicitly (not auto-discovered from the migrations dir), so a new .sql file that is not wired through a version constant + bookkeeping check silently never runs. -*/ -export const IMPORT_TRANSLATION_CACHE_VERSION = "0010"; -/** - * FNXC:GitHubImportTranslate 2026-07-16-23:30: - * Existing databases already recorded 0010, so the cache scope correction is - * deliberately a new forward migration rather than a retroactive SQL edit. - */ -export const IMPORT_TRANSLATION_CACHE_SCOPE_FIX_VERSION = "0016"; -/* -FNXC:GitHubImportTranslate 2026-07-17-23:48: -0016 aligned future cache writes with the normalized legacy partition, but a -pre-0016 cache row can still carry a historic blank project_id. Migration 0019 -backfills that durable data before a restarted store scopes cache reads to -__legacy_unscoped__, preventing an avoidable re-translation. -*/ -export const IMPORT_TRANSLATION_CACHE_LEGACY_PARTITION_BACKFILL_VERSION = "0019"; -/* -FNXC:MultiProjectIsolation 2026-07-15-23:40: -Version 0011 splits the domain "project" field from the RLS partition on the tables -that conflated them: `project_id` stays the trigger/GUC-owned isolation partition, -`owner_project_id` becomes the caller-supplied domain field. Writing domain values -into the partition put parent rows and child rows in different partitions and broke -the composite FKs (SQLSTATE 23503). Keep this identity fixed when -SCHEMA_BASELINE_VERSION advances. -*/ -export const OWNER_PROJECT_ID_SPLIT_VERSION = "0011"; -/* -FNXC:ChatPinned 2026-07-16-12:30: -Version 0012 makes the persisted pin timestamp available on databases that -already applied the baseline before Direct conversations can be pinned. -*/ -export const CHAT_SESSION_PINS_VERSION = "0012"; -/** FNXC:ExecutorToolFailureRetry 2026-07-16-12:00: upgrades existing PostgreSQL task rows before retry-state reads. */ -export const EXECUTOR_TOOL_FAILURE_RETRY_VERSION = "0013"; -/** FNXC:ExecutorEscalation 2026-07-16-21:00: Existing clusters need the durable single-shot latch before executor reads it during post-FN-7996 escalation. */ -export const EXECUTOR_ESCALATION_ATTEMPT_VERSION = "0014"; -/** FNXC:PostgresSchema 2026-07-16-22:00: central global routines follow main's already-landed 0014 migration. */ -export const GLOBAL_ROUTINES_SCHEMA_VERSION = "0015"; -/** FNXC:Settings-MergerModel 2026-07-16-12:00: per-task merger lane is an additive upgrade. */ -export const TASK_MERGER_MODEL_LANE_VERSION = "0017"; -/** - * FNXC:Lifecycle 2026-07-16-22:35: - * Version 0018 lands project.tasks.bulk_completion_refusal_at (FN-8141) on - * existing clusters. PR #2260 added the column to the model + 0000 baseline but - * forgot the forward migration, so every pre-#2260 database crashed on its first - * TaskStore SELECT. Keep this identity fixed when SCHEMA_BASELINE_VERSION advances. - */ -export const BULK_COMPLETION_REFUSAL_AT_VERSION = "0018"; -/** FNXC:EphemeralAgentTaskCreation 2026-07-30-12:00: durable project-scoped proposal key/index protects task creation across crash and reclaim races. */ -export const TASK_PROPOSAL_CLAIM_VERSION = "0020"; -/** FNXC:ConfigVersioning 2026-07-18-00:00: existing clusters need immutable configuration history before write paths use it. */ -export const CONFIGURATION_REVISIONS_VERSION = "0021"; -/** FNXC:Ideation 2026-07-30-15:30: Persisted ideation needs its own forward migration because configuration revisions already own 0021. */ -export const IDEATION_SCHEMA_VERSION = "0022"; -/** FNXC:ResearchMissionBridge 2026-07-18-12:00: forward migration stores stable research finding provenance on canonical features. */ -export const RESEARCH_FEATURE_PROVENANCE_VERSION = "0023"; -/** FNXC:TaskVerificationRequest 2026-07-30-00:00: upgrades need the project-scoped chat-to-executor verification queue. */ -export const TASK_VERIFICATION_REQUEST_VERSION = "0024"; -/** FNXC:SymbolLock 2026-07-30-14:10: upgraded projects need the durable lock table and RLS contract before scheduler admission can use it. */ -export const SYMBOL_LOCKS_SCHEMA_VERSION = "0025"; - -/** Bookkeeping table for the fresh Drizzle migration history. */ -export const MIGRATION_BOOKKEEPING_TABLE = "fusion_schema_migrations"; - -const __dirname = dirname(fileURLToPath(import.meta.url)); - -/* -FNXC:StandaloneExeMigrations 2026-07-17-13:30: -The bun-compiled standalone `fn` binary runs its bundled code from the virtual -/$bunfs/root filesystem, so the historical `join(__dirname, "migrations")` -resolution points at a path that does not exist on disk (bun --compile does not -embed readFile assets) and every DATABASE_URL boot died with -ENOENT /$bunfs/root/migrations/0000_initial.sql. Resolution order: - 1. FUSION_MIGRATIONS_DIR env override — always wins when set (operator escape hatch). - 2. join(__dirname, "migrations") — the npm/tsup and desktop layout; kept first - among the probes so nothing changes for existing installs. - 3. join(dirname(process.execPath), "migrations") — the standalone-exe layout, - where build.ts / the release tarball stage migrations/ next to the binary. -The existsSync probe (not a runtime-detection heuristic) picks between (2) and -(3): inside the compiled binary the module-relative dir simply does not exist, -while for node-based installs it always does, so npm/desktop behavior is untouched. -*/ -function resolveMigrationsDir(): string { - const envDir = process.env.FUSION_MIGRATIONS_DIR; - if (envDir) return envDir; - const moduleDir = join(__dirname, "migrations"); - if (existsSync(join(moduleDir, "0000_initial.sql"))) return moduleDir; - const execDir = join(dirname(process.execPath), "migrations"); - if (existsSync(join(execDir, "0000_initial.sql"))) return execDir; - // Preserve the historical default (and its historical error message) when - // neither location exists — the readFile ENOENT remains the diagnostic. - return moduleDir; -} - -const MIGRATIONS_DIR = resolveMigrationsDir(); -const BASELINE_MIGRATION_PATH = join(MIGRATIONS_DIR, "0000_initial.sql"); -const AUTOMATION_ISOLATION_MIGRATION_PATH = join( - MIGRATIONS_DIR, - "0001_automation_project_isolation.sql", -); -const ANALYTICS_ISOLATION_MIGRATION_PATH = join( - MIGRATIONS_DIR, - "0002_analytics_project_isolation.sql", -); -const MONITOR_APPROVAL_ISOLATION_MIGRATION_PATH = join( - MIGRATIONS_DIR, - "0003_monitor_approval_project_isolation.sql", -); -const LEGACY_CUTOVER_PRESERVATION_MIGRATION_PATH = join( - MIGRATIONS_DIR, - "0004_legacy_cutover_preservation.sql", -); -const MULTI_PROJECT_CUTOVER_MIGRATION_PATH = join( - MIGRATIONS_DIR, - "0005_multi_project_cutover.sql", -); -const PROJECT_OWNERSHIP_MIGRATION_PATH = join( - MIGRATIONS_DIR, - "0006_project_ownership.sql", -); -const SQLITE_SCHEMA_PARITY_MIGRATION_PATH = join( - MIGRATIONS_DIR, - "0007_sqlite_schema_parity.sql", -); -const SESSION_ADVISOR_ENABLED_MIGRATION_PATH = join( - MIGRATIONS_DIR, - "0008_session_advisor_enabled.sql", -); -const MISSION_FIX_IDEMPOTENCY_MIGRATION_PATH = join( - MIGRATIONS_DIR, - "0009_mission_fix_idempotency.sql", -); -const IMPORT_TRANSLATION_CACHE_MIGRATION_PATH = join( - MIGRATIONS_DIR, - "0010_import_translation_cache.sql", -); -const IMPORT_TRANSLATION_CACHE_SCOPE_FIX_MIGRATION_PATH = join( - MIGRATIONS_DIR, - "0016_import_translation_cache_scope_fix.sql", -); -const IMPORT_TRANSLATION_CACHE_LEGACY_PARTITION_BACKFILL_MIGRATION_PATH = join( - MIGRATIONS_DIR, - "0019_import_translation_cache_legacy_partition_backfill.sql", -); -const OWNER_PROJECT_ID_SPLIT_MIGRATION_PATH = join( - MIGRATIONS_DIR, - "0011_owner_project_id.sql", -); -const CHAT_SESSION_PINS_MIGRATION_PATH = join( - MIGRATIONS_DIR, - "0012_chat_session_pins.sql", -); -const EXECUTOR_TOOL_FAILURE_RETRY_MIGRATION_PATH = join( - MIGRATIONS_DIR, - "0013_executor_tool_failure_retry.sql", -); -const EXECUTOR_ESCALATION_ATTEMPT_MIGRATION_PATH = join( - MIGRATIONS_DIR, - "0014_executor_escalation_attempt.sql", -); -const GLOBAL_ROUTINES_MIGRATION_PATH = join( - MIGRATIONS_DIR, - "0015_global_routines.sql", -); -const TASK_MERGER_MODEL_LANE_MIGRATION_PATH = join(MIGRATIONS_DIR, "0017_task_merger_model_lane.sql"); -const BULK_COMPLETION_REFUSAL_AT_MIGRATION_PATH = join(MIGRATIONS_DIR, "0018_bulk_completion_refusal_at.sql"); -const TASK_PROPOSAL_CLAIM_MIGRATION_PATH = join(MIGRATIONS_DIR, "0020_task_proposal_claim.sql"); -const CONFIGURATION_REVISIONS_MIGRATION_PATH = join(MIGRATIONS_DIR, "0021_configuration_revisions.sql"); -const IDEATION_MIGRATION_PATH = join(MIGRATIONS_DIR, "0022_ideation.sql"); -const RESEARCH_FEATURE_PROVENANCE_MIGRATION_PATH = join(MIGRATIONS_DIR, "0023_research_feature_provenance.sql"); -const TASK_VERIFICATION_REQUEST_MIGRATION_PATH = join(MIGRATIONS_DIR, "0024_task_verification_request.sql"); -const SYMBOL_LOCKS_MIGRATION_PATH = join(MIGRATIONS_DIR, "0025_symbol_locks.sql"); - -/** - * Ensure the migration bookkeeping table exists. Lives in the public schema so - * it survives across the three application schemas and is queryable without - * search_path qualification. - */ -async function ensureBookkeepingTable(db: PostgresJsDatabase>): Promise { - await db.execute(sql.raw(` - CREATE TABLE IF NOT EXISTS public.${MIGRATION_BOOKKEEPING_TABLE} ( - version text PRIMARY KEY, - applied_at timestamptz NOT NULL DEFAULT now() - ) - `)); -} - -/** Read the baseline migration SQL from disk. Exported for tests. */ -export async function readBaselineMigrationSql(): Promise { - return readFile(BASELINE_MIGRATION_PATH, "utf8"); -} - -/** Return the set of already-applied migration versions, or empty if none. */ -export async function getAppliedMigrations( - db: PostgresJsDatabase>, -): Promise { - await ensureBookkeepingTable(db); - const rows = (await db.execute( - sql`SELECT version FROM public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} ORDER BY version`, - )) as unknown as Array<{ version: string }>; - return rows.map((row) => row.version); -} - -/** - * Apply the fresh baseline migration to the given connection. - * - * Idempotent: if the baseline version is already recorded, this is a no-op. - * After the baseline lands, all registered plugin schema-init hooks run so - * plugin-owned tables (e.g. roadmap) materialize (VAL-SCHEMA-007). - * - * The baseline SQL is applied as a single batch via postgres.js's file/unsafe - * execution path. It uses CREATE TABLE IF NOT EXISTS and CREATE INDEX IF NOT - * EXISTS throughout, so a partial prior apply is safe to resume. - */ -export async function applySchemaBaseline( - db: PostgresJsDatabase>, - options: { pluginHooks?: readonly PluginSchemaInitHook[] } = {}, -): Promise<{ applied: boolean; pluginHooksRun: number }> { - /* - * FNXC:PostgresSchema 2026-07-14-00:05: - * Schema versions are a cluster-wide invariant. Serialize version discovery, - * DDL, and bookkeeping in one transaction so concurrent Fusion processes - * cannot both apply a version or race its primary-key marker. - */ - return db.transaction(async (tx) => { - await acquireSchemaMutationLocks(tx); - await ensureBookkeepingTable(tx); - /* - FNXC:PostgresSchema 2026-07-16-00:55: - FN-8051 requires project, central, and archive to exist before plugin schema-init hooks run. - Hooks run even when migration markers are already recorded and target project tables, so - ensure the namespaces unconditionally inside the advisory-locked transaction rather than - relying on the baseline batch that a marker-present database skips. - */ - await tx.execute(sql.raw(` - CREATE SCHEMA IF NOT EXISTS project; - CREATE SCHEMA IF NOT EXISTS central; - CREATE SCHEMA IF NOT EXISTS archive; - `)); - const applied = await getAppliedMigrations(tx); - const baselineAlreadyApplied = applied.includes(INITIAL_SCHEMA_VERSION); - const automationIsolationAlreadyApplied = applied.includes(AUTOMATION_ISOLATION_SCHEMA_VERSION); - const analyticsIsolationAlreadyApplied = applied.includes(ANALYTICS_ISOLATION_SCHEMA_VERSION); - const monitorApprovalIsolationAlreadyApplied = applied.includes(MONITOR_APPROVAL_ISOLATION_SCHEMA_VERSION); - const legacyCutoverPreservationAlreadyApplied = applied.includes(LEGACY_CUTOVER_PRESERVATION_SCHEMA_VERSION); - const multiProjectCutoverAlreadyApplied = applied.includes(MULTI_PROJECT_CUTOVER_SCHEMA_VERSION); - const projectOwnershipAlreadyApplied = applied.includes(PROJECT_OWNERSHIP_SCHEMA_VERSION); - const sqliteSchemaParityAlreadyApplied = applied.includes(SQLITE_SCHEMA_PARITY_VERSION); - const sessionAdvisorEnabledAlreadyApplied = applied.includes(SESSION_ADVISOR_ENABLED_SCHEMA_VERSION); - const missionFixIdempotencyAlreadyApplied = applied.includes(MISSION_FIX_IDEMPOTENCY_VERSION); - const importTranslationCacheAlreadyApplied = applied.includes(IMPORT_TRANSLATION_CACHE_VERSION); - const importTranslationCacheScopeFixAlreadyApplied = applied.includes(IMPORT_TRANSLATION_CACHE_SCOPE_FIX_VERSION); - const importTranslationCacheLegacyPartitionBackfillAlreadyApplied = applied.includes(IMPORT_TRANSLATION_CACHE_LEGACY_PARTITION_BACKFILL_VERSION); - const ownerProjectIdSplitAlreadyApplied = applied.includes(OWNER_PROJECT_ID_SPLIT_VERSION); - const chatSessionPinsAlreadyApplied = applied.includes(CHAT_SESSION_PINS_VERSION); - const executorToolFailureRetryAlreadyApplied = applied.includes(EXECUTOR_TOOL_FAILURE_RETRY_VERSION); - const executorEscalationAttemptAlreadyApplied = applied.includes(EXECUTOR_ESCALATION_ATTEMPT_VERSION); - const globalRoutinesAlreadyApplied = applied.includes(GLOBAL_ROUTINES_SCHEMA_VERSION); - const taskMergerModelLaneAlreadyApplied = applied.includes(TASK_MERGER_MODEL_LANE_VERSION); - const bulkCompletionRefusalAtAlreadyApplied = applied.includes(BULK_COMPLETION_REFUSAL_AT_VERSION); - const taskProposalClaimAlreadyApplied = applied.includes(TASK_PROPOSAL_CLAIM_VERSION); - const configurationRevisionsAlreadyApplied = applied.includes(CONFIGURATION_REVISIONS_VERSION); - const ideationAlreadyApplied = applied.includes(IDEATION_SCHEMA_VERSION); - const researchFeatureProvenanceAlreadyApplied = applied.includes(RESEARCH_FEATURE_PROVENANCE_VERSION); - const taskVerificationRequestAlreadyApplied = applied.includes(TASK_VERIFICATION_REQUEST_VERSION); - const symbolLocksAlreadyApplied = applied.includes(SYMBOL_LOCKS_SCHEMA_VERSION); - let schemaChanged = false; - - if (!baselineAlreadyApplied) { - const baselineSql = await readBaselineMigrationSql(); - // The baseline contains multiple statements including CREATE SCHEMA, CREATE - // TABLE, CREATE INDEX, and seed INSERTs. postgres.js executes a single - // query string as one batch (simple query protocol when unparameterized). - await tx.execute(sql.raw(baselineSql)); - await tx.execute( - sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${INITIAL_SCHEMA_VERSION}) ON CONFLICT (version) DO NOTHING`, - ); - schemaChanged = true; - } - - /* - * FNXC:AutomationIsolation 2026-07-13-22:37: - * A database that already recorded the initial PostgreSQL baseline must still receive project-scoped automation storage. Apply this version independently of 0000; ambiguous legacy ownership fails closed before any bound cron runner can silently omit those schedules. - */ - if (!automationIsolationAlreadyApplied) { - const migrationSql = await readFile(AUTOMATION_ISOLATION_MIGRATION_PATH, "utf8"); - await tx.execute(sql.raw(migrationSql)); - await tx.execute( - sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${AUTOMATION_ISOLATION_SCHEMA_VERSION}) ON CONFLICT (version) DO NOTHING`, - ); - schemaChanged = true; - } - - /* - FNXC:AnalyticsIsolation 2026-07-14-00:05: - Existing PostgreSQL databases that already recorded 0001 must independently receive analytics project partitions before project-scoped readers and writers start. Keep 0002 versioned so a fresh baseline cannot hide a skipped upgrade path. - */ - if (!analyticsIsolationAlreadyApplied) { - const migrationSql = await readFile(ANALYTICS_ISOLATION_MIGRATION_PATH, "utf8"); - await tx.execute(sql.raw(migrationSql)); - await tx.execute( - sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${ANALYTICS_ISOLATION_SCHEMA_VERSION}) ON CONFLICT (version) DO NOTHING`, - ); - schemaChanged = true; - } - - /* - FNXC:CommandCenterTenantIsolation 2026-07-14-01:04: - Version 0003 supplies durable ownership for monitor and approval analytics. It must run independently after 0002 so databases that already accepted the earlier analytics migration cannot silently skip the remaining tenant partitions. - */ - if (!monitorApprovalIsolationAlreadyApplied) { - const migrationSql = await readFile(MONITOR_APPROVAL_ISOLATION_MIGRATION_PATH, "utf8"); - await tx.execute(sql.raw(migrationSql)); - await tx.execute( - sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${MONITOR_APPROVAL_ISOLATION_SCHEMA_VERSION}) ON CONFLICT (version) DO NOTHING`, - ); - schemaChanged = true; - } - - /* - FNXC:PostgresMigrationCompleteness 2026-07-14-09:27: - Apply retired-table preservation independently of 0000 so an older partial migration target can retry without losing board, project-auth, or task-reviewer rows. The DDL is additive and idempotent. - */ - if (!legacyCutoverPreservationAlreadyApplied) { - const migrationSql = await readFile(LEGACY_CUTOVER_PRESERVATION_MIGRATION_PATH, "utf8"); - await tx.execute(sql.raw(migrationSql)); - await tx.execute( - sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${LEGACY_CUTOVER_PRESERVATION_SCHEMA_VERSION}) ON CONFLICT (version) DO NOTHING`, - ); - schemaChanged = true; - } - - /* - FNXC:PostgresMultiProjectCutover 2026-07-14-11:18: - Existing targets may already contain one completed project plus a partially copied second project. Apply metadata partitioning and collision-safe revision identity before any retry builds its migration plan. - */ - if (!multiProjectCutoverAlreadyApplied) { - const migrationSql = await readFile(MULTI_PROJECT_CUTOVER_MIGRATION_PATH, "utf8"); - await tx.execute(sql.raw(migrationSql)); - await tx.execute( - sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${MULTI_PROJECT_CUTOVER_SCHEMA_VERSION}) ON CONFLICT (version) DO NOTHING`, - ); - schemaChanged = true; - } - - // Run plugin schema-init hooks regardless of whether the baseline was just - // applied or already present — plugin tables must exist on every connection - // the applier touches. The hooks are themselves idempotent (CREATE TABLE IF - // NOT EXISTS), so re-running is safe. - const pluginHooks = options.pluginHooks ?? DEFAULT_PLUGIN_SCHEMA_INIT_HOOKS; - await runPluginSchemaInitHooks(tx, pluginHooks); - - /* - FNXC:ProjectDataIsolation 2026-07-14-12:10: - Run universal ownership once, after plugin hooks, so first application covers core and plugin tables without duplicate DDL. Later boots validate that every newly introduced plugin table declared the same ownership contract instead of rebuilding primary keys, foreign keys, and policies on every startup. - - FNXC:ProjectArchiveIsolation 2026-07-14-14:31: - The steady-state audit includes archive.archived_tasks because archived task IDs are project-local and must retain the same forced-RLS boundary as live task rows. - */ - if (!projectOwnershipAlreadyApplied) { - const projectOwnershipSql = await readFile(PROJECT_OWNERSHIP_MIGRATION_PATH, "utf8"); - await tx.execute(sql.raw(projectOwnershipSql)); - await tx.execute( - sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${PROJECT_OWNERSHIP_SCHEMA_VERSION}) ON CONFLICT (version) DO NOTHING`, - ); - schemaChanged = true; - } else { - const ownershipGaps = (await tx.execute(sql` - SELECT n.nspname || '.' || c.relname AS table_name - FROM pg_class c - JOIN pg_namespace n ON n.oid = c.relnamespace - LEFT JOIN information_schema.columns col - ON col.table_schema = n.nspname - AND col.table_name = c.relname - AND col.column_name = 'project_id' - WHERE (n.nspname = 'project' OR (n.nspname = 'archive' AND c.relname = 'archived_tasks')) - AND c.relkind = 'r' - AND ( - col.column_name IS NULL - OR NOT c.relrowsecurity - OR NOT c.relforcerowsecurity - OR NOT EXISTS ( - SELECT 1 FROM pg_policy p - WHERE p.polrelid = c.oid AND p.polname = 'fusion_project_isolation' - ) - ) - ORDER BY c.relname - `)) as unknown as Array<{ table_name: string }>; - if (ownershipGaps.length > 0) { - throw new Error( - `Project-owned tables are missing required isolation: ${ownershipGaps.map(({ table_name }) => table_name).join(", ")}`, - ); - } - const relationalGaps = (await tx.execute(sql` - SELECT c.conrelid::regclass::text AS object_name, c.conname AS detail - FROM pg_constraint c - JOIN pg_class t ON t.oid = c.conrelid - JOIN pg_namespace n ON n.oid = t.relnamespace - WHERE (n.nspname = 'project' OR (n.nspname = 'archive' AND t.relname = 'archived_tasks')) - AND c.contype IN ('p', 'u') - AND NOT EXISTS ( - SELECT 1 FROM unnest(c.conkey) key_attnum - JOIN pg_attribute a ON a.attrelid = c.conrelid AND a.attnum = key_attnum - WHERE a.attname = 'project_id' - ) - UNION ALL - SELECT t.oid::regclass::text, idx.relname - FROM pg_index i - JOIN pg_class t ON t.oid = i.indrelid - JOIN pg_class idx ON idx.oid = i.indexrelid - JOIN pg_namespace n ON n.oid = t.relnamespace - WHERE n.nspname = 'project' AND i.indisunique - AND NOT EXISTS ( - SELECT 1 FROM unnest(i.indkey::smallint[]) key_attnum - JOIN pg_attribute a ON a.attrelid = t.oid AND a.attnum = key_attnum - WHERE a.attname = 'project_id' - ) - UNION ALL - SELECT c.conrelid::regclass::text, c.conname - FROM pg_constraint c - JOIN pg_class child ON child.oid = c.conrelid - JOIN pg_namespace child_ns ON child_ns.oid = child.relnamespace - JOIN pg_class parent ON parent.oid = c.confrelid - JOIN pg_namespace parent_ns ON parent_ns.oid = parent.relnamespace - WHERE c.contype = 'f' AND child_ns.nspname = 'project' AND parent_ns.nspname = 'project' - AND ( - NOT EXISTS ( - SELECT 1 FROM unnest(c.conkey) key_attnum - JOIN pg_attribute a ON a.attrelid = c.conrelid AND a.attnum = key_attnum - WHERE a.attname = 'project_id' - ) OR NOT EXISTS ( - SELECT 1 FROM unnest(c.confkey) key_attnum - JOIN pg_attribute a ON a.attrelid = c.confrelid AND a.attnum = key_attnum - WHERE a.attname = 'project_id' - ) - ) - ORDER BY 1, 2 - `)) as unknown as Array<{ object_name: string; detail: string }>; - if (relationalGaps.length > 0) { - throw new Error( - `Project-owned keys or relationships are globally scoped: ${relationalGaps.map(({ object_name, detail }) => `${object_name}.${detail}`).join(", ")}`, - ); - } - } - - /* - FNXC:PostgresMigrationColumnCoverage 2026-07-14-13:17: - Apply SQLite schema parity independently of the original baseline and ownership migration. Existing partial cutover targets must gain all late source columns before the idempotent migration retry rebuilds its table plan. - */ - if (!sqliteSchemaParityAlreadyApplied) { - const sqliteSchemaParitySql = await readFile(SQLITE_SCHEMA_PARITY_MIGRATION_PATH, "utf8"); - await tx.execute(sql.raw(sqliteSchemaParitySql)); - await tx.execute( - sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${SQLITE_SCHEMA_PARITY_VERSION}) ON CONFLICT (version) DO NOTHING`, - ); - schemaChanged = true; - } - - /* - FNXC:PlannerOversight 2026-07-14-18:49: - Apply session_advisor_enabled independently of 0007 so databases that already - recorded SQLite schema parity still gain the per-task session-advisor column - before TaskStore/Drizzle SELECT paths run on boot. - */ - if (!sessionAdvisorEnabledAlreadyApplied) { - const sessionAdvisorEnabledSql = await readFile(SESSION_ADVISOR_ENABLED_MIGRATION_PATH, "utf8"); - await tx.execute(sql.raw(sessionAdvisorEnabledSql)); - await tx.execute( - sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${SESSION_ADVISOR_ENABLED_SCHEMA_VERSION}) ON CONFLICT (version) DO NOTHING`, - ); - schemaChanged = true; - } - - /* - FNXC:MissionFixIdempotency 2026-07-14-18:55: - Existing PostgreSQL databases receive the validator-run lineage uniqueness invariant independently of earlier schema versions. Duplicate historical rows fail the migration visibly instead of being silently discarded. - - FNXC:PostgresConflictResolution 2026-07-14-20:52: - Main assigned migration 0008 to session-advisor state before the cutover landed, so mission lineage uniqueness advances to 0009. Both migrations must run in order; sharing a bookkeeping version would silently skip one invariant. - */ - if (!missionFixIdempotencyAlreadyApplied) { - const migrationSql = await readFile(MISSION_FIX_IDEMPOTENCY_MIGRATION_PATH, "utf8"); - await tx.execute(sql.raw(migrationSql)); - await tx.execute( - sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${MISSION_FIX_IDEMPOTENCY_VERSION}) ON CONFLICT (version) DO NOTHING`, - ); - schemaChanged = true; - } - - /* - FNXC:GitHubImportTranslate 2026-07-15-09:30: - Create the import-translation cache table independently of earlier schema versions so existing databases gain it on boot before any Import Tasks translate/import read runs against it. - */ - if (!importTranslationCacheAlreadyApplied) { - const importTranslationCacheSql = await readFile(IMPORT_TRANSLATION_CACHE_MIGRATION_PATH, "utf8"); - await tx.execute(sql.raw(importTranslationCacheSql)); - await tx.execute( - sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${IMPORT_TRANSLATION_CACHE_VERSION}) ON CONFLICT (version) DO NOTHING`, - ); - schemaChanged = true; - } - - /* - FNXC:MultiProjectIsolation 2026-07-15-23:40: - Apply the owner_project_id domain/partition split independently of earlier - schema versions so existing databases gain the domain column (backfilled from - the previously conflated partition value) before any store read/write path - that now targets owner_project_id runs on boot. - */ - if (!ownerProjectIdSplitAlreadyApplied) { - const ownerProjectIdSplitSql = await readFile(OWNER_PROJECT_ID_SPLIT_MIGRATION_PATH, "utf8"); - await tx.execute(sql.raw(ownerProjectIdSplitSql)); - await tx.execute( - sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${OWNER_PROJECT_ID_SPLIT_VERSION}) ON CONFLICT (version) DO NOTHING`, - ); - schemaChanged = true; - } - - /* - FNXC:ChatPinned 2026-07-16-12:30: - Apply the pin timestamp separately from the baseline so all pre-existing - databases can safely read and write Direct chat pins after this rollout. - */ - if (!chatSessionPinsAlreadyApplied) { - const chatSessionPinsSql = await readFile(CHAT_SESSION_PINS_MIGRATION_PATH, "utf8"); - await tx.execute(sql.raw(chatSessionPinsSql)); - await tx.execute( - sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${CHAT_SESSION_PINS_VERSION}) ON CONFLICT (version) DO NOTHING`, - ); - schemaChanged = true; - } - - if (!executorToolFailureRetryAlreadyApplied) { - const executorToolFailureRetrySql = await readFile(EXECUTOR_TOOL_FAILURE_RETRY_MIGRATION_PATH, "utf8"); - await tx.execute(sql.raw(executorToolFailureRetrySql)); - await tx.execute( - sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${EXECUTOR_TOOL_FAILURE_RETRY_VERSION}) ON CONFLICT (version) DO NOTHING`, - ); - schemaChanged = true; - } - - if (!executorEscalationAttemptAlreadyApplied) { - const executorEscalationAttemptSql = await readFile(EXECUTOR_ESCALATION_ATTEMPT_MIGRATION_PATH, "utf8"); - await tx.execute(sql.raw(executorEscalationAttemptSql)); - await tx.execute( - sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${EXECUTOR_ESCALATION_ATTEMPT_VERSION}) ON CONFLICT (version) DO NOTHING`, - ); - schemaChanged = true; - } - - if (!globalRoutinesAlreadyApplied) { - const migrationSql = await readFile(GLOBAL_ROUTINES_MIGRATION_PATH, "utf8"); - await tx.execute(sql.raw(migrationSql)); - await tx.execute( - sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${GLOBAL_ROUTINES_SCHEMA_VERSION}) ON CONFLICT (version) DO NOTHING`, - ); - schemaChanged = true; - } - if (!taskMergerModelLaneAlreadyApplied) { - const migrationSql = await readFile(TASK_MERGER_MODEL_LANE_MIGRATION_PATH, "utf8"); - await tx.execute(sql.raw(migrationSql)); - await tx.execute( - sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${TASK_MERGER_MODEL_LANE_VERSION}) ON CONFLICT (version) DO NOTHING`, - ); - schemaChanged = true; - } - /* - FNXC:Lifecycle 2026-07-16-22:35: - FN-8141 bulk-completion-refusal taint marker. PR #2260 added the column to - the model + 0000 baseline but no forward migration, so existing clusters - (baseline marker already present) never gained it and crashed on the first - TaskStore SELECT. Apply it as a forward migration so those clusters recover. - */ - if (!bulkCompletionRefusalAtAlreadyApplied) { - const migrationSql = await readFile(BULK_COMPLETION_REFUSAL_AT_MIGRATION_PATH, "utf8"); - await tx.execute(sql.raw(migrationSql)); - await tx.execute( - sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${BULK_COMPLETION_REFUSAL_AT_VERSION}) ON CONFLICT (version) DO NOTHING`, - ); - schemaChanged = true; - } - - if (!taskProposalClaimAlreadyApplied) { - const migrationSql = await readFile(TASK_PROPOSAL_CLAIM_MIGRATION_PATH, "utf8"); - await tx.execute(sql.raw(migrationSql)); - await tx.execute(sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${TASK_PROPOSAL_CLAIM_VERSION}) ON CONFLICT (version) DO NOTHING`); - schemaChanged = true; - } - - /* - FNXC:GitHubImportTranslate 2026-07-16-23:30: - 0010's marker prevents its corrected fresh-install definition from running - on upgrades. Apply 0016 separately before runtime cache reads so existing - rows, RLS, and unbound compatibility stores share one partition contract. - */ - /* - FNXC:ConfigVersioning 2026-07-18-00:00: - Migrations are explicitly registered rather than discovered. Keep 0021's - bookkeeping check adjacent to its apply block so upgrades cannot silently - omit configuration history while fresh installs appear healthy. - */ - if (!configurationRevisionsAlreadyApplied) { - const migrationSql = await readFile(CONFIGURATION_REVISIONS_MIGRATION_PATH, "utf8"); - await tx.execute(sql.raw(migrationSql)); - await tx.execute(sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${CONFIGURATION_REVISIONS_VERSION}) ON CONFLICT (version) DO NOTHING`); - schemaChanged = true; - } - - /* - FNXC:Ideation 2026-07-30-15:30: - Register the ideation migration explicitly. New SQL files are never auto-discovered, - and upgrades must receive project-scoped session/candidate tables before the store opens. - */ - if (!ideationAlreadyApplied) { - const migrationSql = await readFile(IDEATION_MIGRATION_PATH, "utf8"); - await tx.execute(sql.raw(migrationSql)); - await tx.execute(sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${IDEATION_SCHEMA_VERSION}) ON CONFLICT (version) DO NOTHING`); - schemaChanged = true; - } - - if (!researchFeatureProvenanceAlreadyApplied) { - const migrationSql = await readFile(RESEARCH_FEATURE_PROVENANCE_MIGRATION_PATH, "utf8"); - await tx.execute(sql.raw(migrationSql)); - await tx.execute(sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${RESEARCH_FEATURE_PROVENANCE_VERSION}) ON CONFLICT (version) DO NOTHING`); - schemaChanged = true; - } - - if (!importTranslationCacheScopeFixAlreadyApplied) { - const migrationSql = await readFile(IMPORT_TRANSLATION_CACHE_SCOPE_FIX_MIGRATION_PATH, "utf8"); - await tx.execute(sql.raw(migrationSql)); - await tx.execute( - sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${IMPORT_TRANSLATION_CACHE_SCOPE_FIX_VERSION}) ON CONFLICT (version) DO NOTHING`, - ); - schemaChanged = true; - } - - if (!taskVerificationRequestAlreadyApplied) { - const migrationSql = await readFile(TASK_VERIFICATION_REQUEST_MIGRATION_PATH, "utf8"); - await tx.execute(sql.raw(migrationSql)); - await tx.execute(sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${TASK_VERIFICATION_REQUEST_VERSION}) ON CONFLICT (version) DO NOTHING`); - schemaChanged = true; - } - - /* - FNXC:SymbolLock 2026-07-30-14:10: - Migration files are manually registered. Apply 0025 independently so both - fresh and upgraded installations gain forced RLS after 0006 owns its setup. - */ - if (!symbolLocksAlreadyApplied) { - const migrationSql = await readFile(SYMBOL_LOCKS_MIGRATION_PATH, "utf8"); - await tx.execute(sql.raw(migrationSql)); - await tx.execute(sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${SYMBOL_LOCKS_SCHEMA_VERSION}) ON CONFLICT (version) DO NOTHING`); - schemaChanged = true; - } - - if (!importTranslationCacheLegacyPartitionBackfillAlreadyApplied) { - const migrationSql = await readFile(IMPORT_TRANSLATION_CACHE_LEGACY_PARTITION_BACKFILL_MIGRATION_PATH, "utf8"); - await tx.execute(sql.raw(migrationSql)); - await tx.execute( - sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${IMPORT_TRANSLATION_CACHE_LEGACY_PARTITION_BACKFILL_VERSION}) ON CONFLICT (version) DO NOTHING`, - ); - schemaChanged = true; - } - - return { applied: schemaChanged, pluginHooksRun: pluginHooks.length }; - }); -} +/** + * PostgreSQL schema applier. + * + * FNXC:PostgresSchema 2026-06-24-03:40: + * Applies the fresh Drizzle migration baseline to a PostgreSQL connection + * and records it in a migration bookkeeping table. The baseline migration + * (migrations/0000_initial.sql) is the snapshot of the final SQLite schema + * (SCHEMA_VERSION=128) translated to PostgreSQL — applying it to an empty + * database yields final-schema parity (VAL-SCHEMA-001). + * + * After the baseline lands, plugin-owned tables are materialized via the + * schema-init hook (VAL-SCHEMA-007). The applier calls each registered plugin + * hook so plugins evolve their own tables independently of the core migration. + * + * Migration tracking uses a single-row bookkeeping table in the public schema + * so the applier is idempotent: re-running against an already-migrated database + * is a no-op. The version-gate discipline (the institutional learning that + * fresh-DB tests cannot catch a skipped-on-upgrade migration) is carried + * forward via the applier's explicit baseline marker. + */ + +import { readFile } from "node:fs/promises"; +import { existsSync } from "node:fs"; +import { dirname, join } from "node:path"; +import { fileURLToPath } from "node:url"; +import type { PostgresJsDatabase } from "drizzle-orm/postgres-js"; +import { sql } from "drizzle-orm"; +import { runPluginSchemaInitHooks, DEFAULT_PLUGIN_SCHEMA_INIT_HOOKS, type PluginSchemaInitHook } from "./plugin-schema-hook.js"; +import { acquireSchemaMutationLocks } from "./advisory-locks.js"; + +/** The latest PostgreSQL schema version known to this applier. */ +/* +FNXC:GitHubImportTranslate 2026-07-17-23:48: +Advances to 0019 for the import-translation legacy-partition backfill. Per-migration identities above stay fixed; only this latest-version marker moves. + +FNXC:PostgresBigintCounters 2026-07-19-12:00: +SCHEMA_BASELINE_VERSION advances to 0026 for the bigint counters migration. +Per-migration identities above stay fixed; only this latest-version marker moves. +*/ +export const SCHEMA_BASELINE_VERSION = "0026"; +const INITIAL_SCHEMA_VERSION = "0000"; +const AUTOMATION_ISOLATION_SCHEMA_VERSION = "0001"; +const ANALYTICS_ISOLATION_SCHEMA_VERSION = "0002"; +/** + * FNXC:PostgresMigrationIdentity 2026-07-14-01:41: + * Each migration keeps an immutable bookkeeping identity even as SCHEMA_BASELINE_VERSION advances to newer migrations. Upgrade checks and inserts must use this dedicated 0003 identifier so a later latest-version marker cannot make an unrecorded monitor/approval migration look applied. + */ +export const MONITOR_APPROVAL_ISOLATION_SCHEMA_VERSION = "0003"; +export const LEGACY_CUTOVER_PRESERVATION_SCHEMA_VERSION = "0004"; +export const MULTI_PROJECT_CUTOVER_SCHEMA_VERSION = "0005"; +export const PROJECT_OWNERSHIP_SCHEMA_VERSION = "0006"; +export const SQLITE_SCHEMA_PARITY_VERSION = "0007"; +/** + * FNXC:PlannerOversight 2026-07-14-18:49: + * Version 0008 adds project.tasks.session_advisor_enabled for per-task session + * advisor overrides. Keep this identity fixed when SCHEMA_BASELINE_VERSION advances. + */ +export const SESSION_ADVISOR_ENABLED_SCHEMA_VERSION = "0008"; +export const MISSION_FIX_IDEMPOTENCY_VERSION = "0009"; +/* +FNXC:GitHubImportTranslate 2026-07-15-09:30: +Import-translation cache advances to 0010. Migrations are registered here explicitly (not auto-discovered from the migrations dir), so a new .sql file that is not wired through a version constant + bookkeeping check silently never runs. +*/ +export const IMPORT_TRANSLATION_CACHE_VERSION = "0010"; +/** + * FNXC:GitHubImportTranslate 2026-07-16-23:30: + * Existing databases already recorded 0010, so the cache scope correction is + * deliberately a new forward migration rather than a retroactive SQL edit. + */ +export const IMPORT_TRANSLATION_CACHE_SCOPE_FIX_VERSION = "0016"; +/* +FNXC:GitHubImportTranslate 2026-07-17-23:48: +0016 aligned future cache writes with the normalized legacy partition, but a +pre-0016 cache row can still carry a historic blank project_id. Migration 0019 +backfills that durable data before a restarted store scopes cache reads to +__legacy_unscoped__, preventing an avoidable re-translation. +*/ +export const IMPORT_TRANSLATION_CACHE_LEGACY_PARTITION_BACKFILL_VERSION = "0019"; +/* +FNXC:MultiProjectIsolation 2026-07-15-23:40: +Version 0011 splits the domain "project" field from the RLS partition on the tables +that conflated them: `project_id` stays the trigger/GUC-owned isolation partition, +`owner_project_id` becomes the caller-supplied domain field. Writing domain values +into the partition put parent rows and child rows in different partitions and broke +the composite FKs (SQLSTATE 23503). Keep this identity fixed when +SCHEMA_BASELINE_VERSION advances. +*/ +export const OWNER_PROJECT_ID_SPLIT_VERSION = "0011"; +/* +FNXC:ChatPinned 2026-07-16-12:30: +Version 0012 makes the persisted pin timestamp available on databases that +already applied the baseline before Direct conversations can be pinned. +*/ +export const CHAT_SESSION_PINS_VERSION = "0012"; +/** FNXC:ExecutorToolFailureRetry 2026-07-16-12:00: upgrades existing PostgreSQL task rows before retry-state reads. */ +export const EXECUTOR_TOOL_FAILURE_RETRY_VERSION = "0013"; +/** FNXC:ExecutorEscalation 2026-07-16-21:00: Existing clusters need the durable single-shot latch before executor reads it during post-FN-7996 escalation. */ +export const EXECUTOR_ESCALATION_ATTEMPT_VERSION = "0014"; +/** FNXC:PostgresSchema 2026-07-16-22:00: central global routines follow main's already-landed 0014 migration. */ +export const GLOBAL_ROUTINES_SCHEMA_VERSION = "0015"; +/** FNXC:Settings-MergerModel 2026-07-16-12:00: per-task merger lane is an additive upgrade. */ +export const TASK_MERGER_MODEL_LANE_VERSION = "0017"; +/** + * FNXC:Lifecycle 2026-07-16-22:35: + * Version 0018 lands project.tasks.bulk_completion_refusal_at (FN-8141) on + * existing clusters. PR #2260 added the column to the model + 0000 baseline but + * forgot the forward migration, so every pre-#2260 database crashed on its first + * TaskStore SELECT. Keep this identity fixed when SCHEMA_BASELINE_VERSION advances. + */ +export const BULK_COMPLETION_REFUSAL_AT_VERSION = "0018"; +/** FNXC:EphemeralAgentTaskCreation 2026-07-30-12:00: durable project-scoped proposal key/index protects task creation across crash and reclaim races. */ +export const TASK_PROPOSAL_CLAIM_VERSION = "0020"; +/** FNXC:ConfigVersioning 2026-07-18-00:00: existing clusters need immutable configuration history before write paths use it. */ +export const CONFIGURATION_REVISIONS_VERSION = "0021"; +/** FNXC:Ideation 2026-07-30-15:30: Persisted ideation needs its own forward migration because configuration revisions already own 0021. */ +export const IDEATION_SCHEMA_VERSION = "0022"; +/** FNXC:ResearchMissionBridge 2026-07-18-12:00: forward migration stores stable research finding provenance on canonical features. */ +export const RESEARCH_FEATURE_PROVENANCE_VERSION = "0023"; +/** FNXC:TaskVerificationRequest 2026-07-30-00:00: upgrades need the project-scoped chat-to-executor verification queue. */ +export const TASK_VERIFICATION_REQUEST_VERSION = "0024"; +/** FNXC:SymbolLock 2026-07-30-14:10: upgraded projects need the durable lock table and RLS contract before scheduler admission can use it. */ +export const SYMBOL_LOCKS_SCHEMA_VERSION = "0025"; +/** FNXC:PostgresBigintCounters 2026-07-18-21:45: widen overflow-prone counters to bigint before SQLite migration. */ +export const BIGINT_COUNTERS_VERSION = "0026"; + +/** Bookkeeping table for the fresh Drizzle migration history. */ +export const MIGRATION_BOOKKEEPING_TABLE = "fusion_schema_migrations"; + +const __dirname = dirname(fileURLToPath(import.meta.url)); + +/* +FNXC:StandaloneExeMigrations 2026-07-17-13:30: +The bun-compiled standalone `fn` binary runs its bundled code from the virtual +/$bunfs/root filesystem, so the historical `join(__dirname, "migrations")` +resolution points at a path that does not exist on disk (bun --compile does not +embed readFile assets) and every DATABASE_URL boot died with +ENOENT /$bunfs/root/migrations/0000_initial.sql. Resolution order: + 1. FUSION_MIGRATIONS_DIR env override — always wins when set (operator escape hatch). + 2. join(__dirname, "migrations") — the npm/tsup and desktop layout; kept first + among the probes so nothing changes for existing installs. + 3. join(dirname(process.execPath), "migrations") — the standalone-exe layout, + where build.ts / the release tarball stage migrations/ next to the binary. +The existsSync probe (not a runtime-detection heuristic) picks between (2) and +(3): inside the compiled binary the module-relative dir simply does not exist, +while for node-based installs it always does, so npm/desktop behavior is untouched. +*/ +function resolveMigrationsDir(): string { + const envDir = process.env.FUSION_MIGRATIONS_DIR; + if (envDir) return envDir; + const moduleDir = join(__dirname, "migrations"); + if (existsSync(join(moduleDir, "0000_initial.sql"))) return moduleDir; + const execDir = join(dirname(process.execPath), "migrations"); + if (existsSync(join(execDir, "0000_initial.sql"))) return execDir; + // Preserve the historical default (and its historical error message) when + // neither location exists — the readFile ENOENT remains the diagnostic. + return moduleDir; +} + +const MIGRATIONS_DIR = resolveMigrationsDir(); +const BASELINE_MIGRATION_PATH = join(MIGRATIONS_DIR, "0000_initial.sql"); +const AUTOMATION_ISOLATION_MIGRATION_PATH = join( + MIGRATIONS_DIR, + "0001_automation_project_isolation.sql", +); +const ANALYTICS_ISOLATION_MIGRATION_PATH = join( + MIGRATIONS_DIR, + "0002_analytics_project_isolation.sql", +); +const MONITOR_APPROVAL_ISOLATION_MIGRATION_PATH = join( + MIGRATIONS_DIR, + "0003_monitor_approval_project_isolation.sql", +); +const LEGACY_CUTOVER_PRESERVATION_MIGRATION_PATH = join( + MIGRATIONS_DIR, + "0004_legacy_cutover_preservation.sql", +); +const MULTI_PROJECT_CUTOVER_MIGRATION_PATH = join( + MIGRATIONS_DIR, + "0005_multi_project_cutover.sql", +); +const PROJECT_OWNERSHIP_MIGRATION_PATH = join( + MIGRATIONS_DIR, + "0006_project_ownership.sql", +); +const SQLITE_SCHEMA_PARITY_MIGRATION_PATH = join( + MIGRATIONS_DIR, + "0007_sqlite_schema_parity.sql", +); +const SESSION_ADVISOR_ENABLED_MIGRATION_PATH = join( + MIGRATIONS_DIR, + "0008_session_advisor_enabled.sql", +); +const MISSION_FIX_IDEMPOTENCY_MIGRATION_PATH = join( + MIGRATIONS_DIR, + "0009_mission_fix_idempotency.sql", +); +const IMPORT_TRANSLATION_CACHE_MIGRATION_PATH = join( + MIGRATIONS_DIR, + "0010_import_translation_cache.sql", +); +const IMPORT_TRANSLATION_CACHE_SCOPE_FIX_MIGRATION_PATH = join( + MIGRATIONS_DIR, + "0016_import_translation_cache_scope_fix.sql", +); +const IMPORT_TRANSLATION_CACHE_LEGACY_PARTITION_BACKFILL_MIGRATION_PATH = join( + MIGRATIONS_DIR, + "0019_import_translation_cache_legacy_partition_backfill.sql", +); +const OWNER_PROJECT_ID_SPLIT_MIGRATION_PATH = join( + MIGRATIONS_DIR, + "0011_owner_project_id.sql", +); +const CHAT_SESSION_PINS_MIGRATION_PATH = join( + MIGRATIONS_DIR, + "0012_chat_session_pins.sql", +); +const EXECUTOR_TOOL_FAILURE_RETRY_MIGRATION_PATH = join( + MIGRATIONS_DIR, + "0013_executor_tool_failure_retry.sql", +); +const EXECUTOR_ESCALATION_ATTEMPT_MIGRATION_PATH = join( + MIGRATIONS_DIR, + "0014_executor_escalation_attempt.sql", +); +const GLOBAL_ROUTINES_MIGRATION_PATH = join( + MIGRATIONS_DIR, + "0015_global_routines.sql", +); +const TASK_MERGER_MODEL_LANE_MIGRATION_PATH = join(MIGRATIONS_DIR, "0017_task_merger_model_lane.sql"); +const BULK_COMPLETION_REFUSAL_AT_MIGRATION_PATH = join(MIGRATIONS_DIR, "0018_bulk_completion_refusal_at.sql"); +const TASK_PROPOSAL_CLAIM_MIGRATION_PATH = join(MIGRATIONS_DIR, "0020_task_proposal_claim.sql"); +const CONFIGURATION_REVISIONS_MIGRATION_PATH = join(MIGRATIONS_DIR, "0021_configuration_revisions.sql"); +const IDEATION_MIGRATION_PATH = join(MIGRATIONS_DIR, "0022_ideation.sql"); +const RESEARCH_FEATURE_PROVENANCE_MIGRATION_PATH = join(MIGRATIONS_DIR, "0023_research_feature_provenance.sql"); +const TASK_VERIFICATION_REQUEST_MIGRATION_PATH = join(MIGRATIONS_DIR, "0024_task_verification_request.sql"); +const SYMBOL_LOCKS_MIGRATION_PATH = join(MIGRATIONS_DIR, "0025_symbol_locks.sql"); +const BIGINT_COUNTERS_MIGRATION_PATH = join(MIGRATIONS_DIR, "0026_bigint_counters.sql"); + +/** + * Ensure the migration bookkeeping table exists. Lives in the public schema so + * it survives across the three application schemas and is queryable without + * search_path qualification. + */ +async function ensureBookkeepingTable(db: PostgresJsDatabase>): Promise { + await db.execute(sql.raw(` + CREATE TABLE IF NOT EXISTS public.${MIGRATION_BOOKKEEPING_TABLE} ( + version text PRIMARY KEY, + applied_at timestamptz NOT NULL DEFAULT now() + ) + `)); +} + +/** Read the baseline migration SQL from disk. Exported for tests. */ +export async function readBaselineMigrationSql(): Promise { + return readFile(BASELINE_MIGRATION_PATH, "utf8"); +} + +/** Return the set of already-applied migration versions, or empty if none. */ +export async function getAppliedMigrations( + db: PostgresJsDatabase>, +): Promise { + await ensureBookkeepingTable(db); + const rows = (await db.execute( + sql`SELECT version FROM public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} ORDER BY version`, + )) as unknown as Array<{ version: string }>; + return rows.map((row) => row.version); +} + +/** + * Apply the fresh baseline migration to the given connection. + * + * Idempotent: if the baseline version is already recorded, this is a no-op. + * After the baseline lands, all registered plugin schema-init hooks run so + * plugin-owned tables (e.g. roadmap) materialize (VAL-SCHEMA-007). + * + * The baseline SQL is applied as a single batch via postgres.js's file/unsafe + * execution path. It uses CREATE TABLE IF NOT EXISTS and CREATE INDEX IF NOT + * EXISTS throughout, so a partial prior apply is safe to resume. + */ +export async function applySchemaBaseline( + db: PostgresJsDatabase>, + options: { pluginHooks?: readonly PluginSchemaInitHook[] } = {}, +): Promise<{ applied: boolean; pluginHooksRun: number }> { + /* + * FNXC:PostgresSchema 2026-07-14-00:05: + * Schema versions are a cluster-wide invariant. Serialize version discovery, + * DDL, and bookkeeping in one transaction so concurrent Fusion processes + * cannot both apply a version or race its primary-key marker. + */ + return db.transaction(async (tx) => { + await acquireSchemaMutationLocks(tx); + await ensureBookkeepingTable(tx); + /* + FNXC:PostgresSchema 2026-07-16-00:55: + FN-8051 requires project, central, and archive to exist before plugin schema-init hooks run. + Hooks run even when migration markers are already recorded and target project tables, so + ensure the namespaces unconditionally inside the advisory-locked transaction rather than + relying on the baseline batch that a marker-present database skips. + */ + await tx.execute(sql.raw(` + CREATE SCHEMA IF NOT EXISTS project; + CREATE SCHEMA IF NOT EXISTS central; + CREATE SCHEMA IF NOT EXISTS archive; + `)); + const applied = await getAppliedMigrations(tx); + const baselineAlreadyApplied = applied.includes(INITIAL_SCHEMA_VERSION); + const automationIsolationAlreadyApplied = applied.includes(AUTOMATION_ISOLATION_SCHEMA_VERSION); + const analyticsIsolationAlreadyApplied = applied.includes(ANALYTICS_ISOLATION_SCHEMA_VERSION); + const monitorApprovalIsolationAlreadyApplied = applied.includes(MONITOR_APPROVAL_ISOLATION_SCHEMA_VERSION); + const legacyCutoverPreservationAlreadyApplied = applied.includes(LEGACY_CUTOVER_PRESERVATION_SCHEMA_VERSION); + const multiProjectCutoverAlreadyApplied = applied.includes(MULTI_PROJECT_CUTOVER_SCHEMA_VERSION); + const projectOwnershipAlreadyApplied = applied.includes(PROJECT_OWNERSHIP_SCHEMA_VERSION); + const sqliteSchemaParityAlreadyApplied = applied.includes(SQLITE_SCHEMA_PARITY_VERSION); + const sessionAdvisorEnabledAlreadyApplied = applied.includes(SESSION_ADVISOR_ENABLED_SCHEMA_VERSION); + const missionFixIdempotencyAlreadyApplied = applied.includes(MISSION_FIX_IDEMPOTENCY_VERSION); + const importTranslationCacheAlreadyApplied = applied.includes(IMPORT_TRANSLATION_CACHE_VERSION); + const importTranslationCacheScopeFixAlreadyApplied = applied.includes(IMPORT_TRANSLATION_CACHE_SCOPE_FIX_VERSION); + const importTranslationCacheLegacyPartitionBackfillAlreadyApplied = applied.includes(IMPORT_TRANSLATION_CACHE_LEGACY_PARTITION_BACKFILL_VERSION); + const ownerProjectIdSplitAlreadyApplied = applied.includes(OWNER_PROJECT_ID_SPLIT_VERSION); + const chatSessionPinsAlreadyApplied = applied.includes(CHAT_SESSION_PINS_VERSION); + const executorToolFailureRetryAlreadyApplied = applied.includes(EXECUTOR_TOOL_FAILURE_RETRY_VERSION); + const executorEscalationAttemptAlreadyApplied = applied.includes(EXECUTOR_ESCALATION_ATTEMPT_VERSION); + const globalRoutinesAlreadyApplied = applied.includes(GLOBAL_ROUTINES_SCHEMA_VERSION); + const taskMergerModelLaneAlreadyApplied = applied.includes(TASK_MERGER_MODEL_LANE_VERSION); + const bulkCompletionRefusalAtAlreadyApplied = applied.includes(BULK_COMPLETION_REFUSAL_AT_VERSION); + const taskProposalClaimAlreadyApplied = applied.includes(TASK_PROPOSAL_CLAIM_VERSION); + const configurationRevisionsAlreadyApplied = applied.includes(CONFIGURATION_REVISIONS_VERSION); + const ideationAlreadyApplied = applied.includes(IDEATION_SCHEMA_VERSION); + const researchFeatureProvenanceAlreadyApplied = applied.includes(RESEARCH_FEATURE_PROVENANCE_VERSION); + const taskVerificationRequestAlreadyApplied = applied.includes(TASK_VERIFICATION_REQUEST_VERSION); + const symbolLocksAlreadyApplied = applied.includes(SYMBOL_LOCKS_SCHEMA_VERSION); + const bigintCountersAlreadyApplied = applied.includes(BIGINT_COUNTERS_VERSION); + let schemaChanged = false; + + if (!baselineAlreadyApplied) { + const baselineSql = await readBaselineMigrationSql(); + // The baseline contains multiple statements including CREATE SCHEMA, CREATE + // TABLE, CREATE INDEX, and seed INSERTs. postgres.js executes a single + // query string as one batch (simple query protocol when unparameterized). + await tx.execute(sql.raw(baselineSql)); + await tx.execute( + sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${INITIAL_SCHEMA_VERSION}) ON CONFLICT (version) DO NOTHING`, + ); + schemaChanged = true; + } + + /* + * FNXC:AutomationIsolation 2026-07-13-22:37: + * A database that already recorded the initial PostgreSQL baseline must still receive project-scoped automation storage. Apply this version independently of 0000; ambiguous legacy ownership fails closed before any bound cron runner can silently omit those schedules. + */ + if (!automationIsolationAlreadyApplied) { + const migrationSql = await readFile(AUTOMATION_ISOLATION_MIGRATION_PATH, "utf8"); + await tx.execute(sql.raw(migrationSql)); + await tx.execute( + sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${AUTOMATION_ISOLATION_SCHEMA_VERSION}) ON CONFLICT (version) DO NOTHING`, + ); + schemaChanged = true; + } + + /* + FNXC:AnalyticsIsolation 2026-07-14-00:05: + Existing PostgreSQL databases that already recorded 0001 must independently receive analytics project partitions before project-scoped readers and writers start. Keep 0002 versioned so a fresh baseline cannot hide a skipped upgrade path. + */ + if (!analyticsIsolationAlreadyApplied) { + const migrationSql = await readFile(ANALYTICS_ISOLATION_MIGRATION_PATH, "utf8"); + await tx.execute(sql.raw(migrationSql)); + await tx.execute( + sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${ANALYTICS_ISOLATION_SCHEMA_VERSION}) ON CONFLICT (version) DO NOTHING`, + ); + schemaChanged = true; + } + + /* + FNXC:CommandCenterTenantIsolation 2026-07-14-01:04: + Version 0003 supplies durable ownership for monitor and approval analytics. It must run independently after 0002 so databases that already accepted the earlier analytics migration cannot silently skip the remaining tenant partitions. + */ + if (!monitorApprovalIsolationAlreadyApplied) { + const migrationSql = await readFile(MONITOR_APPROVAL_ISOLATION_MIGRATION_PATH, "utf8"); + await tx.execute(sql.raw(migrationSql)); + await tx.execute( + sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${MONITOR_APPROVAL_ISOLATION_SCHEMA_VERSION}) ON CONFLICT (version) DO NOTHING`, + ); + schemaChanged = true; + } + + /* + FNXC:PostgresMigrationCompleteness 2026-07-14-09:27: + Apply retired-table preservation independently of 0000 so an older partial migration target can retry without losing board, project-auth, or task-reviewer rows. The DDL is additive and idempotent. + */ + if (!legacyCutoverPreservationAlreadyApplied) { + const migrationSql = await readFile(LEGACY_CUTOVER_PRESERVATION_MIGRATION_PATH, "utf8"); + await tx.execute(sql.raw(migrationSql)); + await tx.execute( + sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${LEGACY_CUTOVER_PRESERVATION_SCHEMA_VERSION}) ON CONFLICT (version) DO NOTHING`, + ); + schemaChanged = true; + } + + /* + FNXC:PostgresMultiProjectCutover 2026-07-14-11:18: + Existing targets may already contain one completed project plus a partially copied second project. Apply metadata partitioning and collision-safe revision identity before any retry builds its migration plan. + */ + if (!multiProjectCutoverAlreadyApplied) { + const migrationSql = await readFile(MULTI_PROJECT_CUTOVER_MIGRATION_PATH, "utf8"); + await tx.execute(sql.raw(migrationSql)); + await tx.execute( + sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${MULTI_PROJECT_CUTOVER_SCHEMA_VERSION}) ON CONFLICT (version) DO NOTHING`, + ); + schemaChanged = true; + } + + // Run plugin schema-init hooks regardless of whether the baseline was just + // applied or already present — plugin tables must exist on every connection + // the applier touches. The hooks are themselves idempotent (CREATE TABLE IF + // NOT EXISTS), so re-running is safe. + const pluginHooks = options.pluginHooks ?? DEFAULT_PLUGIN_SCHEMA_INIT_HOOKS; + await runPluginSchemaInitHooks(tx, pluginHooks); + + /* + FNXC:ProjectDataIsolation 2026-07-14-12:10: + Run universal ownership once, after plugin hooks, so first application covers core and plugin tables without duplicate DDL. Later boots validate that every newly introduced plugin table declared the same ownership contract instead of rebuilding primary keys, foreign keys, and policies on every startup. + + FNXC:ProjectArchiveIsolation 2026-07-14-14:31: + The steady-state audit includes archive.archived_tasks because archived task IDs are project-local and must retain the same forced-RLS boundary as live task rows. + */ + if (!projectOwnershipAlreadyApplied) { + const projectOwnershipSql = await readFile(PROJECT_OWNERSHIP_MIGRATION_PATH, "utf8"); + await tx.execute(sql.raw(projectOwnershipSql)); + await tx.execute( + sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${PROJECT_OWNERSHIP_SCHEMA_VERSION}) ON CONFLICT (version) DO NOTHING`, + ); + schemaChanged = true; + } else { + const ownershipGaps = (await tx.execute(sql` + SELECT n.nspname || '.' || c.relname AS table_name + FROM pg_class c + JOIN pg_namespace n ON n.oid = c.relnamespace + LEFT JOIN information_schema.columns col + ON col.table_schema = n.nspname + AND col.table_name = c.relname + AND col.column_name = 'project_id' + WHERE (n.nspname = 'project' OR (n.nspname = 'archive' AND c.relname = 'archived_tasks')) + AND c.relkind = 'r' + AND ( + col.column_name IS NULL + OR NOT c.relrowsecurity + OR NOT c.relforcerowsecurity + OR NOT EXISTS ( + SELECT 1 FROM pg_policy p + WHERE p.polrelid = c.oid AND p.polname = 'fusion_project_isolation' + ) + ) + ORDER BY c.relname + `)) as unknown as Array<{ table_name: string }>; + if (ownershipGaps.length > 0) { + throw new Error( + `Project-owned tables are missing required isolation: ${ownershipGaps.map(({ table_name }) => table_name).join(", ")}`, + ); + } + const relationalGaps = (await tx.execute(sql` + SELECT c.conrelid::regclass::text AS object_name, c.conname AS detail + FROM pg_constraint c + JOIN pg_class t ON t.oid = c.conrelid + JOIN pg_namespace n ON n.oid = t.relnamespace + WHERE (n.nspname = 'project' OR (n.nspname = 'archive' AND t.relname = 'archived_tasks')) + AND c.contype IN ('p', 'u') + AND NOT EXISTS ( + SELECT 1 FROM unnest(c.conkey) key_attnum + JOIN pg_attribute a ON a.attrelid = c.conrelid AND a.attnum = key_attnum + WHERE a.attname = 'project_id' + ) + UNION ALL + SELECT t.oid::regclass::text, idx.relname + FROM pg_index i + JOIN pg_class t ON t.oid = i.indrelid + JOIN pg_class idx ON idx.oid = i.indexrelid + JOIN pg_namespace n ON n.oid = t.relnamespace + WHERE n.nspname = 'project' AND i.indisunique + AND NOT EXISTS ( + SELECT 1 FROM unnest(i.indkey::smallint[]) key_attnum + JOIN pg_attribute a ON a.attrelid = t.oid AND a.attnum = key_attnum + WHERE a.attname = 'project_id' + ) + UNION ALL + SELECT c.conrelid::regclass::text, c.conname + FROM pg_constraint c + JOIN pg_class child ON child.oid = c.conrelid + JOIN pg_namespace child_ns ON child_ns.oid = child.relnamespace + JOIN pg_class parent ON parent.oid = c.confrelid + JOIN pg_namespace parent_ns ON parent_ns.oid = parent.relnamespace + WHERE c.contype = 'f' AND child_ns.nspname = 'project' AND parent_ns.nspname = 'project' + AND ( + NOT EXISTS ( + SELECT 1 FROM unnest(c.conkey) key_attnum + JOIN pg_attribute a ON a.attrelid = c.conrelid AND a.attnum = key_attnum + WHERE a.attname = 'project_id' + ) OR NOT EXISTS ( + SELECT 1 FROM unnest(c.confkey) key_attnum + JOIN pg_attribute a ON a.attrelid = c.confrelid AND a.attnum = key_attnum + WHERE a.attname = 'project_id' + ) + ) + ORDER BY 1, 2 + `)) as unknown as Array<{ object_name: string; detail: string }>; + if (relationalGaps.length > 0) { + throw new Error( + `Project-owned keys or relationships are globally scoped: ${relationalGaps.map(({ object_name, detail }) => `${object_name}.${detail}`).join(", ")}`, + ); + } + } + + /* + FNXC:PostgresMigrationColumnCoverage 2026-07-14-13:17: + Apply SQLite schema parity independently of the original baseline and ownership migration. Existing partial cutover targets must gain all late source columns before the idempotent migration retry rebuilds its table plan. + */ + if (!sqliteSchemaParityAlreadyApplied) { + const sqliteSchemaParitySql = await readFile(SQLITE_SCHEMA_PARITY_MIGRATION_PATH, "utf8"); + await tx.execute(sql.raw(sqliteSchemaParitySql)); + await tx.execute( + sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${SQLITE_SCHEMA_PARITY_VERSION}) ON CONFLICT (version) DO NOTHING`, + ); + schemaChanged = true; + } + + /* + FNXC:PlannerOversight 2026-07-14-18:49: + Apply session_advisor_enabled independently of 0007 so databases that already + recorded SQLite schema parity still gain the per-task session-advisor column + before TaskStore/Drizzle SELECT paths run on boot. + */ + if (!sessionAdvisorEnabledAlreadyApplied) { + const sessionAdvisorEnabledSql = await readFile(SESSION_ADVISOR_ENABLED_MIGRATION_PATH, "utf8"); + await tx.execute(sql.raw(sessionAdvisorEnabledSql)); + await tx.execute( + sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${SESSION_ADVISOR_ENABLED_SCHEMA_VERSION}) ON CONFLICT (version) DO NOTHING`, + ); + schemaChanged = true; + } + + /* + FNXC:MissionFixIdempotency 2026-07-14-18:55: + Existing PostgreSQL databases receive the validator-run lineage uniqueness invariant independently of earlier schema versions. Duplicate historical rows fail the migration visibly instead of being silently discarded. + + FNXC:PostgresConflictResolution 2026-07-14-20:52: + Main assigned migration 0008 to session-advisor state before the cutover landed, so mission lineage uniqueness advances to 0009. Both migrations must run in order; sharing a bookkeeping version would silently skip one invariant. + */ + if (!missionFixIdempotencyAlreadyApplied) { + const migrationSql = await readFile(MISSION_FIX_IDEMPOTENCY_MIGRATION_PATH, "utf8"); + await tx.execute(sql.raw(migrationSql)); + await tx.execute( + sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${MISSION_FIX_IDEMPOTENCY_VERSION}) ON CONFLICT (version) DO NOTHING`, + ); + schemaChanged = true; + } + + /* + FNXC:GitHubImportTranslate 2026-07-15-09:30: + Create the import-translation cache table independently of earlier schema versions so existing databases gain it on boot before any Import Tasks translate/import read runs against it. + */ + if (!importTranslationCacheAlreadyApplied) { + const importTranslationCacheSql = await readFile(IMPORT_TRANSLATION_CACHE_MIGRATION_PATH, "utf8"); + await tx.execute(sql.raw(importTranslationCacheSql)); + await tx.execute( + sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${IMPORT_TRANSLATION_CACHE_VERSION}) ON CONFLICT (version) DO NOTHING`, + ); + schemaChanged = true; + } + + /* + FNXC:MultiProjectIsolation 2026-07-15-23:40: + Apply the owner_project_id domain/partition split independently of earlier + schema versions so existing databases gain the domain column (backfilled from + the previously conflated partition value) before any store read/write path + that now targets owner_project_id runs on boot. + */ + if (!ownerProjectIdSplitAlreadyApplied) { + const ownerProjectIdSplitSql = await readFile(OWNER_PROJECT_ID_SPLIT_MIGRATION_PATH, "utf8"); + await tx.execute(sql.raw(ownerProjectIdSplitSql)); + await tx.execute( + sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${OWNER_PROJECT_ID_SPLIT_VERSION}) ON CONFLICT (version) DO NOTHING`, + ); + schemaChanged = true; + } + + /* + FNXC:ChatPinned 2026-07-16-12:30: + Apply the pin timestamp separately from the baseline so all pre-existing + databases can safely read and write Direct chat pins after this rollout. + */ + if (!chatSessionPinsAlreadyApplied) { + const chatSessionPinsSql = await readFile(CHAT_SESSION_PINS_MIGRATION_PATH, "utf8"); + await tx.execute(sql.raw(chatSessionPinsSql)); + await tx.execute( + sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${CHAT_SESSION_PINS_VERSION}) ON CONFLICT (version) DO NOTHING`, + ); + schemaChanged = true; + } + + if (!executorToolFailureRetryAlreadyApplied) { + const executorToolFailureRetrySql = await readFile(EXECUTOR_TOOL_FAILURE_RETRY_MIGRATION_PATH, "utf8"); + await tx.execute(sql.raw(executorToolFailureRetrySql)); + await tx.execute( + sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${EXECUTOR_TOOL_FAILURE_RETRY_VERSION}) ON CONFLICT (version) DO NOTHING`, + ); + schemaChanged = true; + } + + if (!executorEscalationAttemptAlreadyApplied) { + const executorEscalationAttemptSql = await readFile(EXECUTOR_ESCALATION_ATTEMPT_MIGRATION_PATH, "utf8"); + await tx.execute(sql.raw(executorEscalationAttemptSql)); + await tx.execute( + sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${EXECUTOR_ESCALATION_ATTEMPT_VERSION}) ON CONFLICT (version) DO NOTHING`, + ); + schemaChanged = true; + } + + if (!globalRoutinesAlreadyApplied) { + const migrationSql = await readFile(GLOBAL_ROUTINES_MIGRATION_PATH, "utf8"); + await tx.execute(sql.raw(migrationSql)); + await tx.execute( + sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${GLOBAL_ROUTINES_SCHEMA_VERSION}) ON CONFLICT (version) DO NOTHING`, + ); + schemaChanged = true; + } + if (!taskMergerModelLaneAlreadyApplied) { + const migrationSql = await readFile(TASK_MERGER_MODEL_LANE_MIGRATION_PATH, "utf8"); + await tx.execute(sql.raw(migrationSql)); + await tx.execute( + sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${TASK_MERGER_MODEL_LANE_VERSION}) ON CONFLICT (version) DO NOTHING`, + ); + schemaChanged = true; + } + /* + FNXC:Lifecycle 2026-07-16-22:35: + FN-8141 bulk-completion-refusal taint marker. PR #2260 added the column to + the model + 0000 baseline but no forward migration, so existing clusters + (baseline marker already present) never gained it and crashed on the first + TaskStore SELECT. Apply it as a forward migration so those clusters recover. + */ + if (!bulkCompletionRefusalAtAlreadyApplied) { + const migrationSql = await readFile(BULK_COMPLETION_REFUSAL_AT_MIGRATION_PATH, "utf8"); + await tx.execute(sql.raw(migrationSql)); + await tx.execute( + sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${BULK_COMPLETION_REFUSAL_AT_VERSION}) ON CONFLICT (version) DO NOTHING`, + ); + schemaChanged = true; + } + + if (!taskProposalClaimAlreadyApplied) { + const migrationSql = await readFile(TASK_PROPOSAL_CLAIM_MIGRATION_PATH, "utf8"); + await tx.execute(sql.raw(migrationSql)); + await tx.execute(sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${TASK_PROPOSAL_CLAIM_VERSION}) ON CONFLICT (version) DO NOTHING`); + schemaChanged = true; + } + + /* + FNXC:GitHubImportTranslate 2026-07-16-23:30: + 0010's marker prevents its corrected fresh-install definition from running + on upgrades. Apply 0016 separately before runtime cache reads so existing + rows, RLS, and unbound compatibility stores share one partition contract. + */ + /* + FNXC:ConfigVersioning 2026-07-18-00:00: + Migrations are explicitly registered rather than discovered. Keep 0021's + bookkeeping check adjacent to its apply block so upgrades cannot silently + omit configuration history while fresh installs appear healthy. + */ + if (!configurationRevisionsAlreadyApplied) { + const migrationSql = await readFile(CONFIGURATION_REVISIONS_MIGRATION_PATH, "utf8"); + await tx.execute(sql.raw(migrationSql)); + await tx.execute(sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${CONFIGURATION_REVISIONS_VERSION}) ON CONFLICT (version) DO NOTHING`); + schemaChanged = true; + } + + /* + FNXC:Ideation 2026-07-30-15:30: + Register the ideation migration explicitly. New SQL files are never auto-discovered, + and upgrades must receive project-scoped session/candidate tables before the store opens. + */ + if (!ideationAlreadyApplied) { + const migrationSql = await readFile(IDEATION_MIGRATION_PATH, "utf8"); + await tx.execute(sql.raw(migrationSql)); + await tx.execute(sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${IDEATION_SCHEMA_VERSION}) ON CONFLICT (version) DO NOTHING`); + schemaChanged = true; + } + + if (!researchFeatureProvenanceAlreadyApplied) { + const migrationSql = await readFile(RESEARCH_FEATURE_PROVENANCE_MIGRATION_PATH, "utf8"); + await tx.execute(sql.raw(migrationSql)); + await tx.execute(sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${RESEARCH_FEATURE_PROVENANCE_VERSION}) ON CONFLICT (version) DO NOTHING`); + schemaChanged = true; + } + + if (!importTranslationCacheScopeFixAlreadyApplied) { + const migrationSql = await readFile(IMPORT_TRANSLATION_CACHE_SCOPE_FIX_MIGRATION_PATH, "utf8"); + await tx.execute(sql.raw(migrationSql)); + await tx.execute( + sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${IMPORT_TRANSLATION_CACHE_SCOPE_FIX_VERSION}) ON CONFLICT (version) DO NOTHING`, + ); + schemaChanged = true; + } + + if (!taskVerificationRequestAlreadyApplied) { + const migrationSql = await readFile(TASK_VERIFICATION_REQUEST_MIGRATION_PATH, "utf8"); + await tx.execute(sql.raw(migrationSql)); + await tx.execute(sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${TASK_VERIFICATION_REQUEST_VERSION}) ON CONFLICT (version) DO NOTHING`); + schemaChanged = true; + } + + /* + FNXC:SymbolLock 2026-07-30-14:10: + Migration files are manually registered. Apply 0025 independently so both + fresh and upgraded installations gain forced RLS after 0006 owns its setup. + */ + if (!symbolLocksAlreadyApplied) { + const migrationSql = await readFile(SYMBOL_LOCKS_MIGRATION_PATH, "utf8"); + await tx.execute(sql.raw(migrationSql)); + await tx.execute(sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${SYMBOL_LOCKS_SCHEMA_VERSION}) ON CONFLICT (version) DO NOTHING`); + schemaChanged = true; + } + + if (!importTranslationCacheLegacyPartitionBackfillAlreadyApplied) { + const migrationSql = await readFile(IMPORT_TRANSLATION_CACHE_LEGACY_PARTITION_BACKFILL_MIGRATION_PATH, "utf8"); + await tx.execute(sql.raw(migrationSql)); + await tx.execute( + sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${IMPORT_TRANSLATION_CACHE_LEGACY_PARTITION_BACKFILL_VERSION}) ON CONFLICT (version) DO NOTHING`, + ); + schemaChanged = true; + } + + /* + FNXC:PostgresBigintCounters 2026-07-18-21:45: + Existing embedded-PG clusters that were created before this migration still have + integer columns for unbounded token/usage counters, so the SQLite migrator fails on + values larger than int4. Apply the column type widening after ownership/parity and + record the version so fresh baselines already benefit from the bigint DDL. + */ + if (!bigintCountersAlreadyApplied) { + const bigintCountersSql = await readFile(BIGINT_COUNTERS_MIGRATION_PATH, "utf8"); + await tx.execute(sql.raw(bigintCountersSql)); + await tx.execute( + sql`INSERT INTO public.${sql.identifier(MIGRATION_BOOKKEEPING_TABLE)} (version) VALUES (${BIGINT_COUNTERS_VERSION}) ON CONFLICT (version) DO NOTHING`, + ); + schemaChanged = true; + } + return { applied: schemaChanged, pluginHooksRun: pluginHooks.length }; + }); +} diff --git a/packages/core/src/postgres/schema/project.ts b/packages/core/src/postgres/schema/project.ts index 262b449431..d2c4eec24c 100644 --- a/packages/core/src/postgres/schema/project.ts +++ b/packages/core/src/postgres/schema/project.ts @@ -154,11 +154,11 @@ export const tasks = projectSchema.table("tasks", { is not enough for Gate boot-smoke before health reconciliation runs. */ sessionAdvisorEnabled: integer("session_advisor_enabled"), - tokenUsageInputTokens: integer("token_usage_input_tokens"), - tokenUsageOutputTokens: integer("token_usage_output_tokens"), - tokenUsageCachedTokens: integer("token_usage_cached_tokens"), - tokenUsageCacheWriteTokens: integer("token_usage_cache_write_tokens"), - tokenUsageTotalTokens: integer("token_usage_total_tokens"), + tokenUsageInputTokens: bigint("token_usage_input_tokens", { mode: "number" }), + tokenUsageOutputTokens: bigint("token_usage_output_tokens", { mode: "number" }), + tokenUsageCachedTokens: bigint("token_usage_cached_tokens", { mode: "number" }), + tokenUsageCacheWriteTokens: bigint("token_usage_cache_write_tokens", { mode: "number" }), + tokenUsageTotalTokens: bigint("token_usage_total_tokens", { mode: "number" }), tokenUsageFirstUsedAt: text("token_usage_first_used_at"), tokenUsageLastUsedAt: text("token_usage_last_used_at"), tokenUsageModelProvider: text("token_usage_model_provider"), @@ -171,7 +171,7 @@ export const tasks = projectSchema.table("tasks", { updatedAt: text("updated_at").notNull(), columnMovedAt: text("column_moved_at"), firstExecutionAt: text("first_execution_at"), - cumulativeActiveMs: integer("cumulative_active_ms"), + cumulativeActiveMs: bigint("cumulative_active_ms", { mode: "number" }), /* FNXC:PostgresMigrationColumnCoverage 2026-07-14-13:17: Keep the task schema aligned with late SQLite lifecycle migrations. JSON lifecycle markers stay jsonb for native backend reads; retired board/question fields remain text so their legacy payloads round-trip byte-for-byte. @@ -248,7 +248,7 @@ export const tasks = projectSchema.table("tasks", { checkoutNodeId: text("checkout_node_id"), checkoutRunId: text("checkout_run_id"), checkoutLeaseRenewedAt: text("checkout_lease_renewed_at"), - checkoutLeaseEpoch: integer("checkout_lease_epoch").default(0), + checkoutLeaseEpoch: bigint("checkout_lease_epoch", { mode: "number" }).default(0), deletedAt: text("deleted_at"), allowResurrection: integer("allow_resurrection").default(0), transitionPending: text("transition_pending"), @@ -1928,11 +1928,11 @@ export const chatTokenUsage = projectSchema.table("chat_token_usage", { agentId: text("agent_id"), modelProvider: text("model_provider"), modelId: text("model_id"), - inputTokens: integer("input_tokens").notNull(), - outputTokens: integer("output_tokens").notNull(), - cachedTokens: integer("cached_tokens").notNull(), - cacheWriteTokens: integer("cache_write_tokens").notNull(), - totalTokens: integer("total_tokens").notNull(), + inputTokens: bigint("input_tokens", { mode: "number" }).notNull(), + outputTokens: bigint("output_tokens", { mode: "number" }).notNull(), + cachedTokens: bigint("cached_tokens", { mode: "number" }).notNull(), + cacheWriteTokens: bigint("cache_write_tokens", { mode: "number" }).notNull(), + totalTokens: bigint("total_tokens", { mode: "number" }).notNull(), createdAt: text("created_at").notNull(), }, (t) => [ index("idxChatTokenUsageCreatedAt").on(t.createdAt), diff --git a/packages/core/src/postgres/sqlite-migrator.ts b/packages/core/src/postgres/sqlite-migrator.ts index bd80550944..bfe341397c 100644 --- a/packages/core/src/postgres/sqlite-migrator.ts +++ b/packages/core/src/postgres/sqlite-migrator.ts @@ -70,6 +70,131 @@ import { getErrorMessage } from "../error-message.js"; const log = createLogger("sqlite-migrator"); +/** + * FNXC:PostgresMigrationSharedSingleton 2026-07-19-05:00: + * A small set of project-schema tables intentionally have no project_id and are + * shared across all projects in the PostgreSQL cluster. When two legacy SQLite + * files use the same primary-key value (e.g. distributed_task_id_state.prefix), + * a plain INSERT ... ON CONFLICT DO NOTHING discards the second file's row and + * the content-verification subset check fails because the target row differs from + * the source. These tables must instead merge conflicting rows using semantics + * appropriate to each column. Only distributed_task_id_state is currently known + * to need this treatment; any future shared singleton with overlapping keys can + * register here. + */ +function isSharedSingletonTable(pgSchema: string, pgTable: string): boolean { + return pgSchema === PROJECT_SCHEMA && pgTable === "distributed_task_id_state"; +} + +/** + * FNXC:PostgresMigrationSharedSingleton 2026-07-19-05:05: + * For distributed_task_id_state, the source SQLite files are per-project, but the + * PostgreSQL target is a cluster-wide singleton keyed by prefix. The allocator + * sequence floor (next_sequence) must be the maximum across all source files so + * newly-created task IDs never reuse an already-issued number. The committed-cluster + * counter and last-committed-task-id follow the row with the highest sequence. + * updated_at is kept as the latest ISO timestamp. This merge is idempotent. + */ +function buildSharedSingletonConflictClause(pgSchema: string, pgTable: string, cols: readonly ColumnMapping[]): ReturnType { + if (pgTable === "distributed_task_id_state") { + const colSet = new Set(cols.map((c) => c.pgName)); + const prefix = quoteIdent("prefix"); + const projectId = quoteIdent("project_id"); + const tableName = quoteIdent("distributed_task_id_state"); + const parts: string[] = []; + if (colSet.has("next_sequence")) { + parts.push(`${quoteIdent("next_sequence")} = GREATEST(${tableName}.${quoteIdent("next_sequence")}, EXCLUDED.${quoteIdent("next_sequence")})`); + } + if (colSet.has("committed_cluster_task_count")) { + parts.push(`${quoteIdent("committed_cluster_task_count")} = GREATEST(${tableName}.${quoteIdent("committed_cluster_task_count")}, EXCLUDED.${quoteIdent("committed_cluster_task_count")})`); + } + if (colSet.has("last_committed_task_id")) { + parts.push(`${quoteIdent("last_committed_task_id")} = CASE WHEN EXCLUDED.${quoteIdent("next_sequence")} > ${tableName}.${quoteIdent("next_sequence")} THEN EXCLUDED.${quoteIdent("last_committed_task_id")} ELSE ${tableName}.${quoteIdent("last_committed_task_id")} END`); + } + if (colSet.has("updated_at")) { + parts.push(`${quoteIdent("updated_at")} = GREATEST(${tableName}.${quoteIdent("updated_at")}, EXCLUDED.${quoteIdent("updated_at")})`); + } + return sql.raw(`ON CONFLICT (${projectId}, ${prefix}) DO UPDATE SET ${parts.join(", ")}`); + } + return sql.raw("ON CONFLICT DO NOTHING"); +} + +/** + * FNXC:PostgresMigrationSharedSingleton 2026-07-19-05:10: + * Verify a shared singleton table by dominance rather than exact row equality. + * For each source row, the target must contain a row with the same primary key + * and counter values that are at least as large as the source values. This lets + * a later project migrate safely even when an earlier project already populated + * the shared table with a higher sequence floor for the same prefix. + * + * FNXC:PostgresMigrationSharedSingleton 2026-07-19-08:40: + * sourceRows come from raw SQLite `SELECT *` (PRAGMA camelCase columns: + * nextSequence, committedClusterTaskCount). Target counters stay snake_case + * from PostgreSQL. Reading snake_case on the source always fell back to 0 and + * let a non-dominating target pass verification. + * + * FNXC:PostgresMigrationSharedSingleton 2026-07-19-09:50: + * PK is (project_id, prefix). Scope the target lookup by partitionProjectId so + * a higher next_sequence on another project's same prefix cannot make result[0] + * pass while the migrating project's row is missing or lower. Fail closed when + * the partition id is absent — prefix-only lookup is ambiguous under multi-project. + */ +async function verifySharedSingletonTable( + db: PostgresJsDatabase>, + pgSchema: string, + pgTable: string, + sourceRows: readonly Record[], + partitionProjectId: string | undefined, +): Promise { + if (pgTable !== "distributed_task_id_state") return false; + if (typeof partitionProjectId !== "string" || partitionProjectId.length === 0) { + log.warn( + `Shared singleton ${pgSchema}.${pgTable} verification requires partitionProjectId ` + + `(composite PK project_id, prefix); refusing prefix-only lookup`, + ); + return false; + } + const schemaQualifiedTable = `${quoteIdent(pgSchema)}.${quoteIdent(pgTable)}`; + for (const row of sourceRows) { + const prefix = row.prefix; + if (typeof prefix !== "string") continue; + const result = (await db.execute(sql` + SELECT + ${sql.raw(quoteIdent("next_sequence"))} AS next_sequence, + ${sql.raw(quoteIdent("committed_cluster_task_count"))} AS committed_cluster_task_count, + ${sql.raw(quoteIdent("updated_at"))} AS updated_at + FROM ${sql.raw(schemaQualifiedTable)} + WHERE ${sql.raw(quoteIdent("project_id"))} = ${partitionProjectId} + AND ${sql.raw(quoteIdent("prefix"))} = ${prefix} + `)) as unknown as Array<{ + next_sequence: number | string; + committed_cluster_task_count: number | string; + updated_at: string; + }>; + if (result.length === 0) { + log.warn( + `Shared singleton ${pgSchema}.${pgTable} missing project_id=${partitionProjectId} prefix=${prefix}`, + ); + return false; + } + const target = result[0]; + // Legacy SQLite schema uses camelCase; never prefer snake_case for source. + const sourceNext = Number(row.nextSequence ?? 0); + const targetNext = Number(target.next_sequence ?? 0); + const sourceCount = Number(row.committedClusterTaskCount ?? 0); + const targetCount = Number(target.committed_cluster_task_count ?? 0); + if (targetNext < sourceNext || targetCount < sourceCount) { + log.warn( + `Shared singleton ${pgSchema}.${pgTable} project_id=${partitionProjectId} prefix=${prefix} does not dominate source: ` + + `source=(nextSequence=${sourceNext}, count=${sourceCount}), target=(next_sequence=${targetNext}, count=${targetCount})`, + ); + return false; + } + } + return true; +} + + /** Batch size for streaming row inserts. */ const INSERT_BATCH_SIZE = 200; @@ -1667,9 +1792,21 @@ async function migrateTable( const targetCanonicalRows = await computeTargetCanonicalRows( db, plan.pgSchema, plan.pgTable, insertableCols, plan.partitionProjectId, ); - contentOk = verifiesSharedProjectTable - ? isCanonicalMultisetSubset(sourceCanonicalRows, targetCanonicalRows) - : checksumCanonicalRows(sourceCanonicalRows) === checksumCanonicalRows(targetCanonicalRows); + /* + FNXC:PostgresMigrationSharedSingleton 2026-07-19-09:50: + distributed_task_id_state always uses project-scoped dominance verification + (GREATEST merge + composite PK). Other shared project tables keep multiset + subset when unbound; partitioned tables keep exact checksums. + */ + contentOk = isSharedSingletonTable(plan.pgSchema, plan.pgTable) + ? await verifySharedSingletonTable( + db, plan.pgSchema, plan.pgTable, + sqlite.prepare(`SELECT * FROM ${quoteIdent(plan.table)}`).all() as Record[], + plan.partitionProjectId, + ) + : verifiesSharedProjectTable + ? isCanonicalMultisetSubset(sourceCanonicalRows, targetCanonicalRows) + : checksumCanonicalRows(sourceCanonicalRows) === checksumCanonicalRows(targetCanonicalRows); if (!contentOk) { log.warn( `Content checksum mismatch for ${plan.pgSchema}.${plan.pgTable}: ` + @@ -1824,12 +1961,15 @@ async function insertBatch( const replacesCentralSeed = plan.pgSchema === CENTRAL_SCHEMA && (plan.pgTable === "central_settings" || plan.pgTable === "global_concurrency"); + const isSharedSingleton = isSharedSingletonTable(plan.pgSchema, plan.pgTable); const conflictClause = replacesCentralSeed ? sql.raw(`ON CONFLICT (${quoteIdent("id")}) DO UPDATE SET ${cols .filter((column) => column.pgName !== "id") .map((column) => `${quoteIdent(column.pgName)} = EXCLUDED.${quoteIdent(column.pgName)}`) .join(", ")}`) - : sql`ON CONFLICT DO NOTHING`; + : isSharedSingleton + ? buildSharedSingletonConflictClause(plan.pgSchema, plan.pgTable, cols) + : sql`ON CONFLICT DO NOTHING`; /* FNXC:PostgresMigration 2026-07-13-21:05: