# Migrate storage from SQLite to PostgreSQL — full dashboard cutover Migrates Fusion's storage layer to the embedded PostgreSQL `AsyncDataLayer` (the default backend) and **completes the satellite-store + feature cutover** so every dashboard and Command Center surface works in PG mode. ## Status — every surface works in embedded-PG mode Verified live against a running embedded-Postgres dashboard (all **200**, zero 5xx) and gate-tested (**23 files / 99 tests** on embedded PG, plus engine-core 294 and ci-shape 63 in the blocking merge gate; core/engine/cli/dashboard typecheck clean). | Area | Surfaces | State | |---|---|---| | Satellite stores | workflows, todos, insights, research, missions, goals, mailbox | ✅ | | Views | artifacts, documents, evals | ✅ | | Command Center | activity, productivity, team, tokens, tools, **workflows**, **github**, **signals**, **plugin-activations**, **live** (all 10) | ✅ | | Run execution | insight generation, research run execution | ✅ (store-path; AI step needs a provider) | | Live updates | SSE push for mission/research/insight events | ✅ | | Workflow editing | create / update / delete / select (+ id counter) | ✅ | | Engine | mission autopilot, incident-signal ingestion, regression storm-guard, agent wake-on-message | ✅ | | Core | tasks, agents, secrets, automations, memory, chat, usage, PRs, git | ✅ | ## Approach Each satellite store gets an `Async<Store>` wrapper exposing the sync store's method names over the existing `async-*-store.ts` helpers; `get<Store>Store()` returns a `Sync | Async` union; consumers `await` (harmless on sync), and engine/CLI paths that can't convert use `instanceof Sync` graceful fallback. Analytics aggregators branch on `"ping" in dbOrLayer` to run schema-qualified raw SQL over `project.*` (snake_case) in PG. Executors/orchestrators/autopilot are await-converted to drive the union store; the async store wrappers extend `EventEmitter` so SSE live-push fires in both backends. Not-yet-ported capabilities degrade gracefully (never 500) and are individually called out in commits. ## Sync with main The branch is kept continuously merged with `main` (currently through FN-7845, 2026-07-12); the earlier "final rebase deferred" note no longer applies. Use **Create a merge commit** (or squash) to land it — GitHub's rebase-merge cannot replay a merge-maintained branch. ## Residual Review Findings Multi-agent code review of the PostgreSQL satellite-store ports (U1–U5) applied 3 safe fixes (see `fix(review): apply autofix feedback`). The following are **real but gated** — recorded here as follow-up work rather than auto-applied. All are SQLite→PostgreSQL **concurrency/atomicity regressions**: the sync stores were immune only by SQLite's single-writer, single-threaded-handler execution; the async ports open multi-await read-modify-write windows. **Reachability is low today** because the execution engines that generate concurrent same-run mutations (insight run executor, research orchestrator/dispatcher) are `instanceof`-gated to sync mode in PG. No process-crash class survived (all engine fallbacks correctly guard the sync store). - **[P1] Research `appendResearchEvent` dual-write is non-atomic** (`packages/core/src/async-research-store.ts`, corroborated: adversarial + reliability). The `research_run_events` insert (own transaction) and the `run.events` jsonb update are separate writes — a crash between them, or two concurrent appends, splits the table count from the jsonb array. **Fix:** perform the seq-insert and the jsonb update in one `layer.transactionImmediate`. - **[P1] Research run terminal-reversion via stale full-row persist** (`async-research-store.ts` `persistResearchRun`/`updateResearchStatus`). Concurrent `PATCH /runs/:id/status` + `POST /runs/:id/events` can revert a terminal run to `running` by overwriting the whole row, bypassing the transition guard. **Fix:** scoped column `UPDATE`s with a `WHERE status …` guard, or optimistic version column. - **[P2] `updateResearchRun`/`updateInsightRun` read-then-write TOCTOU** — concurrent PATCHes last-writer-wins on the lifecycle merge. **Fix:** `SELECT … FOR UPDATE` / enclosing transaction. - **[P2] `upsertRun`/`createRunOrThrowConflict` check-then-create race** (`async-insight-store.ts`) — two callers can each create an "active" run. **Fix:** partial unique index on `(projectId, trigger) WHERE status IN ('pending','running')`. - **[P3] `createResearchRetryRun` return-value divergence** — sync returns the pre-update `queued` snapshot; async returns the reloaded `retry_waiting` run (persisted state is identical). Pick one side for cross-backend parity. - **[P2/perf] Mission `getMissionWithHierarchy`/`getMissionHealth` N+1 fan-out** — O(milestones×slices) sequential round-trips hold one pool slot per request; can starve the pool for large hierarchies. **Fix:** batched/joined reads. - **Testing gaps:** no PG-mode concurrency tests (interleaved status/event mutations), no sync↔async parity assertion for the lifecycle-error codes, and no mission status/health rollup parity test vs the sync `MissionStore`. ~~Out of scope (deferred): AI run *execution* (insight/research) + mission autopilot + live SSE mission events remain sync-gated/degraded in PG mode.~~ **Since ported** — insight/research run execution, mission autopilot, and SSE live push all run on the async layer now, which also makes the concurrency findings above genuinely reachable; they remain open follow-ups. --- ## Update — 2026-07-12: production-readiness hardening & live acceptance Everything below landed on this branch since the description above was written: **Production blockers from review — fixed** - `recoverStaleTransitionPending` ported to the async layer (backend moves write + clear the crash-safe marker; startup/maintenance sweeps no longer throw). - Lost-update class fixed: `atomicWriteTaskJson`/`WithAudit` write changed columns only (full-row upserts silently resurrected stale fields across concurrent store instances — the "task stuck unplanned forever" bug). - First-boot **auto-migration**: booting the PG backend over a project with a legacy `fusion.db` migrates it automatically (loud failure, SQLite kept as backup), and the dashboard shows a one-time **"your data was migrated" banner** with the backup paths and a Need-help Discord link. - `pg_dump`/`pg_restore` discovered from common install locations for embedded-mode backups. - The PG suite is part of the blocking merge gate (`test:pg-gate`). **Multi-project isolation (PR #2007, merged into this branch)** - `project_id` partition key on tasks / archived tasks / config, `taskProjectScope` threaded through every scan/claim/count, per-project config rows, layer bound to the project at startup. - Review P1 follow-up: the shared cold-storage `archive.archived_tasks` table is also partitioned and all archived-board reads/counts/searches are scoped. - Schema drift self-heal generalized to schema-qualified columns so existing databases upgrade in place. **Other changes** - Node settings sync **removed** in PG mode (409 `settings-sync-disabled-postgres`) — nodes share state by connecting to the same database; auth sync kept (per-machine file). - Perf (review findings): `listTasks` pushes column filter + ORDER BY + LIMIT/OFFSET into SQL; `getConversation` capped to the most recent 200 messages. - Fixed a false "operator action required" pause-abort log fired on every successfully auto-merged task. **Live acceptance — PASSED (2026-07-12)** A sandboxed instance (isolated HOME, embedded PG, real Opus executor) ran a task through the complete cycle: create → triage (AI spec) → execute → in-review → AI squash-merge landed on the project's `main` → done. A write+read sweep of every data surface (settings, comments, documents, attachments + artifact bridge + artifact edit, chat with real generation, goals, missions, agent mail, secrets, workflows, memory, CC analytics) was green on embedded PG. **Known remaining work** - The per-project `config` PK re-key has no upgrade path for pre-isolation embedded-PG databases (needs a real `DROP CONSTRAINT`/re-key migration; fresh databases are fine). - `pg_dump`/`pg_restore` binaries are not yet bundled in release artifacts (PATH/common-location discovery only). - The satellite-store concurrency findings listed above. --------- Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com> Co-authored-by: Phil Larson <hello@phillarson.xyz> Co-authored-by: fusion-merge <fusion-merge@local>
195 lines
7.6 KiB
TypeScript
195 lines
7.6 KiB
TypeScript
/**
|
|
* Async task-ID integrity detector for PostgreSQL (U8).
|
|
*
|
|
* FNXC:TaskIdIntegrity 2026-06-24-15:00:
|
|
* PostgreSQL-backed equivalent of the SQLite `detectTaskIdIntegrityAnomalies`
|
|
* in `task-id-integrity.ts`. The detector is preserved per the feature
|
|
* description ("Preserve task-ID-integrity detector") and surfaces the same
|
|
* anomaly kinds via the same `TaskIdIntegrityReport` shape so the dashboard
|
|
* banner and `/api/health` payload remain compatible.
|
|
*
|
|
* The detector checks for (VAL-HEALTH-003):
|
|
* - duplicate task IDs inside `tasks`
|
|
* - task IDs that exist in both `tasks` and `archived_tasks` (cross-table
|
|
* collision)
|
|
* - `distributed_task_id_state.next_sequence` values that point at or below
|
|
* an already-used numeric suffix (sequence drift)
|
|
* - active task rows whose prefix falls outside the prefixes declared in
|
|
* `distributed_task_id_state`
|
|
*
|
|
* All queries use the Drizzle query builder against the project schema so the
|
|
* detector works identically against embedded or external PostgreSQL.
|
|
*/
|
|
|
|
import { sql } from "drizzle-orm";
|
|
import type { DrizzleDb } from "./data-layer.js";
|
|
import type {
|
|
TaskIdIntegrityAnomaly,
|
|
TaskIdIntegrityReport,
|
|
} from "../task-id-integrity.js";
|
|
import { PROJECT_SCHEMA } from "./schema/_shared.js";
|
|
|
|
const TASK_ID_PATTERN = /^([A-Z][A-Z0-9]*)-(\d+)$/;
|
|
|
|
function parseTaskId(taskId: string): { prefix: string; sequence: number } | null {
|
|
const match = taskId.trim().toUpperCase().match(TASK_ID_PATTERN);
|
|
if (!match) return null;
|
|
const sequence = Number.parseInt(match[2], 10);
|
|
if (!Number.isFinite(sequence)) return null;
|
|
return { prefix: match[1], sequence };
|
|
}
|
|
|
|
function uniqueSorted(values: Iterable<string>): string[] {
|
|
return Array.from(new Set(values)).sort((a, b) => a.localeCompare(b));
|
|
}
|
|
|
|
function buildReport(checkedAt: string, anomalies: TaskIdIntegrityAnomaly[]): TaskIdIntegrityReport {
|
|
return {
|
|
status: anomalies.length > 0 ? "anomaly" : "ok",
|
|
checkedAt,
|
|
anomalies,
|
|
};
|
|
}
|
|
|
|
/**
|
|
* FNXC:TaskIdIntegrity 2026-06-24-15:05:
|
|
* Detect task-ID integrity anomalies against a PostgreSQL backend via Drizzle.
|
|
* This is the async PostgreSQL equivalent of the sync SQLite
|
|
* `detectTaskIdIntegrityAnomalies(db)`.
|
|
*
|
|
* The detector intentionally does NOT filter on `deletedAt` for the `tasks`
|
|
* table — soft-deleted IDs must remain visible to integrity checks (FN-5105).
|
|
*
|
|
* @param db The runtime Drizzle instance.
|
|
* @returns The integrity report with the same shape as the SQLite version.
|
|
*/
|
|
export async function detectTaskIdIntegrityAnomaliesAsync(db: DrizzleDb): Promise<TaskIdIntegrityReport> {
|
|
const checkedAt = new Date().toISOString();
|
|
|
|
try {
|
|
const anomalies: TaskIdIntegrityAnomaly[] = [];
|
|
|
|
// Read all active and archived task IDs. We intentionally do not filter
|
|
// deletedAt on tasks (FN-5105). Use raw SQL for direct column access
|
|
// without needing full Drizzle row-type mapping.
|
|
const activeRows = (await db.execute(
|
|
sql.raw(`SELECT id FROM ${PROJECT_SCHEMA}.tasks`),
|
|
)) as unknown as Array<{ id: string }>;
|
|
const archivedRows = (await db.execute(
|
|
sql.raw(`SELECT id FROM ${PROJECT_SCHEMA}.archived_tasks`),
|
|
)) as unknown as Array<{ id: string }>;
|
|
|
|
const activeIds = activeRows.map((r) => String(r.id ?? ""));
|
|
const archivedIds = archivedRows.map((r) => String(r.id ?? ""));
|
|
const allIds = [...activeIds, ...archivedIds];
|
|
|
|
// 1. Duplicate active IDs.
|
|
const idCounts = new Map<string, number>();
|
|
for (const id of activeIds) {
|
|
idCounts.set(id, (idCounts.get(id) ?? 0) + 1);
|
|
}
|
|
for (const [id, count] of idCounts) {
|
|
if (count > 1) {
|
|
const parsed = parseTaskId(id);
|
|
anomalies.push({
|
|
kind: "duplicate_active_id",
|
|
prefix: parsed?.prefix ?? "unknown",
|
|
affectedIds: [id],
|
|
details: `Active tasks contains ${count} rows for ${id}.`,
|
|
});
|
|
}
|
|
}
|
|
|
|
// 2. IDs in both active and archived (cross-table collision).
|
|
const archivedIdSet = new Set(archivedIds);
|
|
const activeAndArchived = uniqueSorted(activeIds.filter((id) => archivedIdSet.has(id)));
|
|
if (activeAndArchived.length > 0) {
|
|
const byPrefix = new Map<string, string[]>();
|
|
for (const taskId of activeAndArchived) {
|
|
const prefix = parseTaskId(taskId)?.prefix ?? "unknown";
|
|
byPrefix.set(prefix, [...(byPrefix.get(prefix) ?? []), taskId]);
|
|
}
|
|
for (const [prefix, affectedIds] of byPrefix) {
|
|
anomalies.push({
|
|
kind: "id_in_active_and_archived",
|
|
prefix,
|
|
affectedIds,
|
|
details: `Task IDs exist in both active and archived storage for prefix ${prefix}.`,
|
|
});
|
|
}
|
|
}
|
|
|
|
// 3. Compute max used sequence per prefix.
|
|
const maxUsedSequenceByPrefix = new Map<string, { maxSequence: number; taskIds: string[] }>();
|
|
for (const taskId of allIds) {
|
|
const parsed = parseTaskId(taskId);
|
|
if (!parsed) continue;
|
|
const existing = maxUsedSequenceByPrefix.get(parsed.prefix);
|
|
if (!existing || parsed.sequence > existing.maxSequence) {
|
|
maxUsedSequenceByPrefix.set(parsed.prefix, { maxSequence: parsed.sequence, taskIds: [taskId] });
|
|
continue;
|
|
}
|
|
if (parsed.sequence === existing.maxSequence) {
|
|
existing.taskIds.push(taskId);
|
|
}
|
|
}
|
|
|
|
// Read allocator state rows.
|
|
const stateRows = (await db.execute(
|
|
sql.raw(`SELECT prefix, next_sequence FROM ${PROJECT_SCHEMA}.distributed_task_id_state`),
|
|
)) as unknown as Array<{ prefix: string; next_sequence: string | number }>;
|
|
|
|
// 4. Sequence drift: next_sequence at or below a used suffix.
|
|
for (const stateRow of stateRows) {
|
|
const prefix = String(stateRow.prefix).trim().toUpperCase();
|
|
const nextSequence = Number(stateRow.next_sequence);
|
|
const maxUsed = maxUsedSequenceByPrefix.get(prefix);
|
|
if (!maxUsed) continue;
|
|
if (nextSequence <= maxUsed.maxSequence) {
|
|
anomalies.push({
|
|
kind: "next_sequence_at_or_below_used",
|
|
prefix,
|
|
affectedIds: uniqueSorted(maxUsed.taskIds),
|
|
details: `distributed_task_id_state.next_sequence=${nextSequence} is at or below existing sequence ${maxUsed.maxSequence} for prefix ${prefix}.`,
|
|
});
|
|
}
|
|
}
|
|
|
|
// 5. Active task rows with a prefix outside known allocator prefixes.
|
|
if (stateRows.length > 0) {
|
|
const knownPrefixes = new Set(
|
|
stateRows
|
|
.map((row) => String(row.prefix).trim().toUpperCase())
|
|
.filter((prefix) => prefix.length > 0),
|
|
);
|
|
if (knownPrefixes.size > 0) {
|
|
const outsideKnownPrefix = new Map<string, string[]>();
|
|
for (const taskId of activeIds) {
|
|
const parsed = parseTaskId(taskId);
|
|
const prefix = parsed?.prefix ?? "unknown";
|
|
if (!parsed || !knownPrefixes.has(prefix)) {
|
|
outsideKnownPrefix.set(prefix, [...(outsideKnownPrefix.get(prefix) ?? []), taskId]);
|
|
}
|
|
}
|
|
for (const [prefix, affectedIds] of outsideKnownPrefix) {
|
|
anomalies.push({
|
|
kind: "task_row_outside_known_prefix",
|
|
prefix,
|
|
affectedIds: uniqueSorted(affectedIds),
|
|
details:
|
|
prefix === "unknown"
|
|
? "Active task rows contain IDs that do not match the expected PREFIX-123 format."
|
|
: `Active task rows use prefix ${prefix}, which is not declared in distributed_task_id_state.`,
|
|
});
|
|
}
|
|
}
|
|
}
|
|
|
|
return buildReport(checkedAt, anomalies);
|
|
} catch {
|
|
// On any query failure, return an ok report (fail-open). The separate
|
|
// health check will surface connectivity issues via the corruption banner.
|
|
return buildReport(checkedAt, []);
|
|
}
|
|
}
|