Files
fusion/packages/core/src/insight-run-executor.ts
gsxdsm c15c78feeb feat: migrate storage from SQLite to PostgreSQL (#1793)
# 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>
2026-07-13 19:07:58 -07:00

315 lines
10 KiB
TypeScript

import { setTimeout as delay } from "node:timers/promises";
import type {
InsightRun,
InsightRunCreateInput,
InsightRunFailureClass,
InsightRunOutputMetadata,
InsightRunTrigger,
InsightRunUpdateInput,
} from "./insight-types.js";
import { InsightLifecycleError, InsightStore } from "./insight-store.js";
import type { AsyncInsightStore } from "./async-insight-store.js";
/*
* FNXC:InsightStore 2026-06-28-10:00:
* The run executor must drive both backends: the sync SQLite `InsightStore` and
* the PostgreSQL-backed `AsyncInsightStore` (async). Both expose the SAME method
* names returning the SAME shapes, so the executor types `store` as the union and
* `await`s every store call — a sync method's awaited return is identical to its
* direct return, so lifecycle semantics are preserved across both backends.
*/
type InsightRunExecutorStore = InsightStore | AsyncInsightStore;
export interface InsightRunAttemptResult {
summary?: string | null;
insightsCreated: number;
insightsUpdated: number;
outputMetadata?: InsightRunOutputMetadata;
}
export interface InsightRunAttemptContext {
run: InsightRun;
attempt: number;
maxAttempts: number;
signal: AbortSignal;
}
export interface InsightRunExecutorOptions {
store: InsightRunExecutorStore;
projectId: string;
input: InsightRunCreateInput;
executeAttempt: (ctx: InsightRunAttemptContext) => Promise<InsightRunAttemptResult>;
timeoutMs?: number;
maxAttempts?: number;
retryDelayMs?: number;
signal?: AbortSignal;
}
export interface InsightRunExecutorErrorClassification {
failureClass: InsightRunFailureClass;
retryable: boolean;
terminalReason: "cancelled" | "failed" | "timed_out";
terminalCause: string;
}
function isAbortLike(error: unknown): boolean {
return error instanceof DOMException && error.name === "AbortError";
}
function asErrorMessage(error: unknown): string {
return error instanceof Error ? error.message : String(error);
}
export function classifyInsightRunError(error: unknown): InsightRunExecutorErrorClassification {
if (isAbortLike(error)) {
return {
failureClass: "cancelled",
retryable: false,
terminalReason: "cancelled",
terminalCause: asErrorMessage(error),
};
}
const message = asErrorMessage(error);
if (/timeout|timed out|deadline/i.test(message)) {
return {
failureClass: "timed_out",
retryable: true,
terminalReason: "timed_out",
terminalCause: message,
};
}
if (/ECONNRESET|ENOTFOUND|EAI_AGAIN|ETIMEDOUT|429|5\d\d/i.test(message)) {
return {
failureClass: "retryable_transient",
retryable: true,
terminalReason: "failed",
terminalCause: message,
};
}
return {
failureClass: "non_retryable",
retryable: false,
terminalReason: "failed",
terminalCause: message,
};
}
function composeSignal(timeoutMs: number | undefined, parent?: AbortSignal): { signal: AbortSignal; clear: () => void } {
const controller = new AbortController();
const timeoutId = timeoutMs && timeoutMs > 0
? setTimeout(() => controller.abort(new Error(`Insight run timed out after ${timeoutMs}ms`)), timeoutMs)
: undefined;
const onAbort = () => {
controller.abort(parent?.reason ?? new DOMException("Aborted", "AbortError"));
};
if (parent) {
if (parent.aborted) onAbort();
else parent.addEventListener("abort", onAbort, { once: true });
}
return {
signal: controller.signal,
clear: () => {
if (timeoutId) clearTimeout(timeoutId);
if (parent) parent.removeEventListener("abort", onAbort);
},
};
}
function patchForStatus(status: "completed" | "failed" | "cancelled", patch: InsightRunUpdateInput): InsightRunUpdateInput {
if (status === "cancelled") {
return {
...patch,
cancelledAt: patch.cancelledAt ?? new Date().toISOString(),
};
}
return patch;
}
async function executeExistingRun(
store: InsightRunExecutorStore,
run: InsightRun,
options: Omit<InsightRunExecutorOptions, "input" | "projectId"> & { maxAttempts: number; retryDelayMs: number },
): Promise<InsightRun> {
const started = await store.updateRun(run.id, {
status: "running",
startedAt: run.startedAt ?? new Date().toISOString(),
lifecycle: {
...run.lifecycle,
maxAttempts: options.maxAttempts,
attempt: run.lifecycle.attempt ?? 1,
},
});
let active = started ?? run;
await store.appendRunEvent(active.id, { type: "status_changed", status: "running", message: "Run started" });
for (let attempt = active.lifecycle.attempt ?? 1; attempt <= options.maxAttempts; attempt += 1) {
const { signal, clear } = composeSignal(options.timeoutMs, options.signal);
try {
if (signal.aborted) {
throw signal.reason instanceof Error ? signal.reason : new DOMException("Aborted", "AbortError");
}
await store.appendRunEvent(active.id, {
type: "info",
message: `Attempt ${attempt}/${options.maxAttempts}`,
metadata: { attempt, maxAttempts: options.maxAttempts },
});
const result = await options.executeAttempt({ run: active, attempt, maxAttempts: options.maxAttempts, signal });
const completed = await store.updateRun(active.id, {
status: "completed",
summary: result.summary ?? null,
insightsCreated: result.insightsCreated,
insightsUpdated: result.insightsUpdated,
outputMetadata: result.outputMetadata,
lifecycle: {
...active.lifecycle,
attempt,
maxAttempts: options.maxAttempts,
terminalReason: "completed",
retryable: false,
},
});
if (!completed) throw new Error(`Run disappeared while completing: ${active.id}`);
await store.appendRunEvent(completed.id, { type: "status_changed", status: "completed", message: "Run completed" });
return completed;
} catch (error) {
const classification = classifyInsightRunError(error);
const canRetry = classification.retryable && attempt < options.maxAttempts;
await store.appendRunEvent(active.id, {
type: canRetry ? "retry_scheduled" : "error",
status: canRetry ? "running" : classification.terminalReason === "cancelled" ? "cancelled" : "failed",
classification: classification.failureClass,
message: canRetry
? `Attempt ${attempt} failed (${classification.failureClass}); retrying`
: `Run failed (${classification.failureClass})`,
metadata: { attempt, maxAttempts: options.maxAttempts, error: asErrorMessage(error) },
});
if (canRetry) {
active = await store.updateRun(active.id, {
lifecycle: {
...active.lifecycle,
attempt: attempt + 1,
maxAttempts: options.maxAttempts,
failureClass: classification.failureClass,
retryable: true,
},
}) ?? active;
if (options.retryDelayMs > 0) {
await delay(options.retryDelayMs, undefined, { signal: options.signal });
}
continue;
}
const terminalStatus = classification.terminalReason === "cancelled" ? "cancelled" : "failed";
const terminal = await store.updateRun(active.id, patchForStatus(terminalStatus, {
status: terminalStatus,
error: asErrorMessage(error),
lifecycle: {
...active.lifecycle,
attempt,
maxAttempts: options.maxAttempts,
terminalReason: classification.terminalReason,
terminalCause: classification.terminalCause,
failureClass: classification.failureClass,
retryable: classification.failureClass === "retryable_transient",
timeoutAt: classification.failureClass === "timed_out" ? new Date().toISOString() : active.lifecycle.timeoutAt,
},
}));
if (!terminal) throw new Error(`Run disappeared while failing: ${active.id}`);
return terminal;
} finally {
clear();
}
}
const failed = await store.updateRun(active.id, {
status: "failed",
error: "Run exhausted attempts",
lifecycle: {
...active.lifecycle,
terminalReason: "failed",
terminalCause: "Run exhausted attempts",
failureClass: "non_retryable",
retryable: false,
attempt: options.maxAttempts,
maxAttempts: options.maxAttempts,
},
});
if (!failed) throw new Error(`Run disappeared after attempts exhausted: ${active.id}`);
return failed;
}
export async function executeInsightRunLifecycle(options: InsightRunExecutorOptions): Promise<InsightRun> {
const maxAttempts = Math.max(1, options.maxAttempts ?? 2);
const retryDelayMs = Math.max(0, options.retryDelayMs ?? 250);
let run: InsightRun;
try {
run = await options.store.createRunOrThrowConflict(options.projectId, {
...options.input,
lifecycle: {
...options.input.lifecycle,
attempt: options.input.lifecycle?.attempt ?? 1,
maxAttempts,
rootRunId: options.input.lifecycle?.rootRunId,
},
});
} catch (error) {
if (error instanceof InsightLifecycleError && error.code === "active_run_conflict") {
throw error;
}
throw error;
}
await options.store.appendRunEvent(run.id, {
type: "status_changed",
status: "pending",
message: "Run created",
});
return executeExistingRun(options.store, run, {
...options,
maxAttempts,
retryDelayMs,
});
}
export async function retryInsightRunLifecycle(
options: Omit<InsightRunExecutorOptions, "input" | "projectId"> & { runId: string; trigger?: InsightRunTrigger; inputMetadata?: InsightRunCreateInput["inputMetadata"] },
): Promise<{ run: InsightRun; retryOf: InsightRun }> {
const original = await options.store.getRun(options.runId);
if (!original) {
throw new Error(`Insight run not found: ${options.runId}`);
}
if (original.status !== "failed") {
throw new InsightLifecycleError(`Run ${original.id} must be failed to retry`, "not_retryable");
}
if (!original.lifecycle.retryable || original.lifecycle.failureClass !== "retryable_transient") {
throw new InsightLifecycleError(`Run ${original.id} is non-retryable`, "not_retryable");
}
const run = await executeInsightRunLifecycle({
...options,
projectId: original.projectId,
input: {
trigger: options.trigger ?? original.trigger,
inputMetadata: options.inputMetadata ?? original.inputMetadata,
lifecycle: {
retryOfRunId: original.id,
rootRunId: original.lifecycle.rootRunId ?? original.id,
attempt: 1,
},
},
});
return { run, retryOf: original };
}