# 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>
272 lines
9.4 KiB
TypeScript
272 lines
9.4 KiB
TypeScript
import { randomUUID } from "node:crypto";
|
|
|
|
import {
|
|
EXPERIMENT_RUN_OUTCOMES,
|
|
type ExperimentMetricDefinition,
|
|
type ExperimentRunOutcome,
|
|
type ExperimentRunRecordPayload,
|
|
type ExperimentSecondaryMetric,
|
|
type ExperimentSession,
|
|
type ExperimentSessionRecord,
|
|
type ExperimentSessionStore,
|
|
} from "@fusion/core";
|
|
|
|
import { AgentSemaphore } from "./concurrency.js";
|
|
import { runBenchmark as defaultRunBenchmark, type BenchmarkRunOptions } from "./experiment/benchmark-runner.js";
|
|
import { defaultGitOps, type GitOps } from "./experiment/git-ops.js";
|
|
import { commitKept, ExperimentRevertConflictError, revertDiscarded } from "./experiment/git-policy.js";
|
|
import { parseMetricLines } from "./experiment/metric-parser.js";
|
|
import { createLogger, formatError } from "./logger.js";
|
|
|
|
export class ExperimentMaxIterationsError extends Error {}
|
|
export class ExperimentGitNotConfiguredError extends Error {}
|
|
|
|
export interface ExperimentExecutorOptions {
|
|
store: ExperimentSessionStore;
|
|
git?: GitOps;
|
|
runBenchmark?: typeof defaultRunBenchmark;
|
|
maxConcurrentExperiments?: number;
|
|
logger?: ReturnType<typeof createLogger>;
|
|
}
|
|
|
|
export interface InitExperimentInput {
|
|
name: string;
|
|
metric: ExperimentMetricDefinition;
|
|
maxIterations?: number;
|
|
workingDir?: string;
|
|
rules?: string;
|
|
ideas?: string;
|
|
projectId?: string;
|
|
tags?: string[];
|
|
}
|
|
|
|
export interface RunExperimentInput {
|
|
sessionId: string;
|
|
command: string;
|
|
cwd: string;
|
|
timeoutMs?: number;
|
|
env?: NodeJS.ProcessEnv;
|
|
onProgress?: BenchmarkRunOptions["onProgress"];
|
|
}
|
|
|
|
export interface RunExperimentResult {
|
|
runHandle: string;
|
|
exitCode: number;
|
|
stdout: string;
|
|
stderr: string;
|
|
durationMs: number;
|
|
primaryMetric?: { name: string; value: number; unit?: string };
|
|
secondaryMetrics: ExperimentSecondaryMetric[];
|
|
parseWarnings: string[];
|
|
status: "pending" | "errored";
|
|
truncatedTempFile?: string;
|
|
}
|
|
|
|
export interface LogExperimentInput {
|
|
sessionId: string;
|
|
runResult: RunExperimentResult;
|
|
outcome: ExperimentRunOutcome;
|
|
description?: string;
|
|
asi?: Record<string, unknown>;
|
|
confidence?: number;
|
|
commitMessage?: string;
|
|
baselineCommit?: string;
|
|
}
|
|
|
|
export interface ExperimentExecutorStatus {
|
|
sessionId: string;
|
|
status: ExperimentSession["status"];
|
|
currentSegment: number;
|
|
runsInSegment: number;
|
|
activeHandles: string[];
|
|
maxIterations?: number;
|
|
}
|
|
|
|
export class ExperimentExecutor {
|
|
private readonly semaphore: AgentSemaphore;
|
|
private readonly activeRuns = new Map<string, { controller: AbortController; sessionId: string; startedAt: number }>();
|
|
private readonly runBenchmark;
|
|
private readonly logger;
|
|
|
|
constructor(private readonly options: ExperimentExecutorOptions) {
|
|
this.semaphore = new AgentSemaphore(options.maxConcurrentExperiments ?? 2);
|
|
this.runBenchmark = options.runBenchmark ?? defaultRunBenchmark;
|
|
this.logger = options.logger ?? createLogger("experiment-executor");
|
|
}
|
|
|
|
async initExperiment(input: InitExperimentInput): Promise<{ session: ExperimentSession; configRecord: ExperimentSessionRecord }> {
|
|
const configPayload = {
|
|
metric: input.metric,
|
|
maxIterations: input.maxIterations,
|
|
workingDir: input.workingDir,
|
|
rules: input.rules,
|
|
ideas: input.ideas,
|
|
};
|
|
|
|
const existing = (await this.options.store
|
|
.listSessions({ projectId: input.projectId }))
|
|
.find((session) => session.name === input.name && ["active", "finalizing"].includes(session.status));
|
|
|
|
if (existing) {
|
|
const result = await this.options.store.startNewSegment(existing.id, configPayload);
|
|
this.logger.log(`initExperiment: ${existing.id} mode=new-segment`);
|
|
return { session: result.session, configRecord: result.record };
|
|
}
|
|
|
|
const session = await this.options.store.createSession({
|
|
name: input.name,
|
|
projectId: input.projectId,
|
|
metric: input.metric,
|
|
maxIterations: input.maxIterations,
|
|
workingDir: input.workingDir,
|
|
tags: input.tags,
|
|
status: "active",
|
|
currentSegment: 1,
|
|
});
|
|
|
|
const configRecord = await this.options.store.appendRecord(session.id, {
|
|
type: "config",
|
|
payload: configPayload,
|
|
segment: session.currentSegment,
|
|
});
|
|
|
|
this.logger.log(`initExperiment: ${session.id} mode=created`);
|
|
return { session, configRecord };
|
|
}
|
|
|
|
async runExperiment(input: RunExperimentInput, opts?: { abortSignal?: AbortSignal }): Promise<RunExperimentResult> {
|
|
const session = await this.options.store.getSession(input.sessionId);
|
|
if (!session || session.status !== "active") throw new Error("Session not active");
|
|
|
|
const runsInSegment = (await this.options.store
|
|
.listRecords(input.sessionId, { segment: session.currentSegment, type: "run" }))
|
|
.length;
|
|
if (session.maxIterations !== undefined && runsInSegment >= session.maxIterations) {
|
|
throw new ExperimentMaxIterationsError(`Session ${input.sessionId} reached max iterations`);
|
|
}
|
|
|
|
await this.semaphore.acquire();
|
|
const controller = new AbortController();
|
|
const runHandle = randomUUID();
|
|
if (opts?.abortSignal) {
|
|
opts.abortSignal.addEventListener("abort", () => controller.abort(), { once: true });
|
|
}
|
|
this.activeRuns.set(runHandle, { controller, sessionId: input.sessionId, startedAt: Date.now() });
|
|
|
|
try {
|
|
const benchmark = await this.runBenchmark({
|
|
command: input.command,
|
|
cwd: input.cwd,
|
|
timeoutMs: input.timeoutMs,
|
|
env: input.env,
|
|
abortSignal: controller.signal,
|
|
onProgress: input.onProgress,
|
|
sessionId: input.sessionId,
|
|
});
|
|
const parsed = parseMetricLines(benchmark.stdout);
|
|
const status = benchmark.exitCode !== 0 || benchmark.timedOut || !parsed.primary ? "errored" : "pending";
|
|
return {
|
|
runHandle,
|
|
exitCode: benchmark.exitCode,
|
|
stdout: benchmark.stdout,
|
|
stderr: benchmark.stderr,
|
|
durationMs: benchmark.durationMs,
|
|
primaryMetric: parsed.primary,
|
|
secondaryMetrics: parsed.secondary,
|
|
parseWarnings: parsed.warnings,
|
|
status,
|
|
truncatedTempFile: benchmark.truncatedTempFile,
|
|
};
|
|
} catch (error) {
|
|
this.logger.error(`runExperiment failed: ${formatError(error)}`);
|
|
throw error;
|
|
} finally {
|
|
this.activeRuns.delete(runHandle);
|
|
this.semaphore.release();
|
|
}
|
|
}
|
|
|
|
async logExperiment(input: LogExperimentInput): Promise<{ runRecord: ExperimentSessionRecord; commit?: string; revertedTo?: string }> {
|
|
const session = await this.options.store.getSession(input.sessionId);
|
|
if (!session) throw new Error(`Experiment session not found: ${input.sessionId}`);
|
|
if (!EXPERIMENT_RUN_OUTCOMES.includes(input.outcome)) throw new Error(`Invalid outcome: ${input.outcome}`);
|
|
if (input.outcome === "keep" && !input.runResult.primaryMetric) throw new Error("keep outcome requires primary metric");
|
|
if (["discard", "checks_failed"].includes(input.outcome) && !input.baselineCommit) {
|
|
throw new Error("baselineCommit is required for discard/checks_failed");
|
|
}
|
|
if (input.outcome === "keep" && !this.options.git) {
|
|
throw new ExperimentGitNotConfiguredError("Git ops not configured");
|
|
}
|
|
|
|
const payload: ExperimentRunRecordPayload = {
|
|
commit: undefined,
|
|
primaryMetric: input.runResult.primaryMetric?.value ?? Number.NaN,
|
|
secondaryMetrics: input.runResult.secondaryMetrics,
|
|
status: input.outcome,
|
|
description: input.description,
|
|
confidence: input.confidence,
|
|
asi: input.asi,
|
|
durationMs: input.runResult.durationMs,
|
|
};
|
|
|
|
const runRecord = await this.options.store.appendRecord(input.sessionId, {
|
|
type: "run",
|
|
payload,
|
|
segment: session.currentSegment,
|
|
});
|
|
|
|
let commit: string | undefined;
|
|
let revertedTo: string | undefined;
|
|
|
|
if (input.outcome === "keep" && this.options.git) {
|
|
const result = await commitKept({
|
|
session,
|
|
runRecord,
|
|
runPayload: payload,
|
|
git: this.options.git,
|
|
commitMessage: input.commitMessage,
|
|
});
|
|
commit = result.commit;
|
|
await this.options.store.updateRecordPayload(runRecord.id, { commit });
|
|
await this.options.store.setBestRun(input.sessionId, runRecord.id);
|
|
await this.options.store.recordKept(input.sessionId, runRecord.id);
|
|
}
|
|
|
|
if (["discard", "checks_failed", "errored"].includes(input.outcome) && input.baselineCommit && this.options.git) {
|
|
const result = await revertDiscarded({ session, git: this.options.git, baselineCommit: input.baselineCommit });
|
|
revertedTo = result.revertedTo;
|
|
}
|
|
|
|
return { runRecord, commit, revertedTo };
|
|
}
|
|
|
|
async getStatus(sessionId: string): Promise<ExperimentExecutorStatus> {
|
|
const session = await this.options.store.getSession(sessionId);
|
|
if (!session) throw new Error(`Experiment session not found: ${sessionId}`);
|
|
const runsInSegment = (await this.options.store
|
|
.listRecords(sessionId, { segment: session.currentSegment, type: "run" }))
|
|
.length;
|
|
const activeHandles = [...this.activeRuns.entries()]
|
|
.filter(([, value]) => value.sessionId === sessionId)
|
|
.map(([handle]) => handle);
|
|
|
|
return {
|
|
sessionId,
|
|
status: session.status,
|
|
currentSegment: session.currentSegment,
|
|
runsInSegment,
|
|
activeHandles,
|
|
maxIterations: session.maxIterations,
|
|
};
|
|
}
|
|
|
|
cancel(runHandle: string): boolean {
|
|
const active = this.activeRuns.get(runHandle);
|
|
if (!active) return false;
|
|
active.controller.abort();
|
|
return true;
|
|
}
|
|
}
|
|
|
|
export { ExperimentRevertConflictError, defaultGitOps };
|