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.
<!-- This is an auto-generated comment: release notes by coderabbit.ai
-->
## 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.
<!-- end of auto-generated comment: release notes by coderabbit.ai -->
---------
Co-authored-by: gsxdsm <gsxdsm@users.noreply.github.com>
This commit is contained in:
@@ -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,
|
||||
]);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -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
|
||||
);
|
||||
|
||||
|
||||
@@ -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
|
||||
$$;
|
||||
File diff suppressed because it is too large
Load Diff
@@ -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),
|
||||
|
||||
@@ -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<typeof sql.raw> {
|
||||
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<Record<string, never>>,
|
||||
pgSchema: string,
|
||||
pgTable: string,
|
||||
sourceRows: readonly Record<string, unknown>[],
|
||||
partitionProjectId: string | undefined,
|
||||
): Promise<boolean> {
|
||||
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<string, unknown>[],
|
||||
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:
|
||||
|
||||
Reference in New Issue
Block a user