Files
fusion/plugins/fusion-plugin-compound-engineering/src/sync/reconciler.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

254 lines
10 KiB
TypeScript

import type { PluginContext, Task } from "@fusion/core";
import { listPipelineStages } from "../session/stage-registry.js";
import { createCeTaskWithLink } from "./ce-task.js";
import {
getCePipelineStore,
type CePipelineLink,
type CePipelineState,
type CePipelineStore,
} from "./pipeline-store.js";
/**
* BIDIRECTIONAL SYNC RECONCILER (U8 / FN-5719 pattern).
*
* Two SEPARATE state machines are kept in sync, never merged (KTD4):
* - Board-task ownership → the task's `column` (board is authoritative).
* - CE-pipeline ownership → `ce_pipeline_state.{currentStage,status}` (CE flow
* is authoritative for artifact/pipeline content).
*
* INBOUND (board → pipeline): the lifecycle hooks (`onTaskMoved`/`onTaskCompleted`
* in index.ts) do the MINIMUM under the 5s hook budget — resolve the link and
* `enqueueSync(...)`, then return. They do NOT advance the pipeline inline.
*
* RECONCILE (the convergence guarantee): `reconcile()` is a single on-demand
* sweep — NOT a tight interval poll (per docs/performance/dashboard-load.md).
* It (1) drains the queue and (2) INDEPENDENTLY re-derives transitions by
* comparing live board state (`ctx.taskStore`) against pipeline state. Step (2)
* is why a DROPPED or never-enqueued hook event still converges: the queue is an
* optimization; the board↔state comparison is the source of truth.
*
* OUTBOUND (pipeline → board): when a pipeline advances to a stage that produces
* board work, the reconciler creates the next-stage board task via
* `ctx.taskStore.createTask` and links it — propagating the CE-flow change onto
* the board.
*
* TRIGGER MODEL (honest about the host seam): there is NO host scheduler wired to
* call this on a timer. In production the sweep is invoked (a) right after the
* hooks enqueue (a cheap drain on the same board mutation that triggered the
* hook), and (b) on demand from a route (U9 settings/refresh surface) or on a
* dashboard session-change. Because step (2) re-derives from board truth, any
* single missed trigger is recovered on the NEXT sweep — no continuous poll loop
* is needed for correctness.
*/
/** Columns that mean "this stage's board work is finished" → advance the pipeline. */
const TERMINAL_COLUMNS = new Set(["in-review", "done"]);
export interface ReconcileResult {
/** Queue entries drained this sweep. */
drained: number;
/** Pipelines whose state advanced this sweep. */
advanced: number;
/** Board tasks created outbound this sweep (next-stage propagation). */
tasksCreated: number;
/** Pipelines inspected. */
inspected: number;
}
/**
* The linear CE stage order. The pipeline advances along this sequence.
* Manual utility stages are excluded; remaining stages are sorted by explicit
* `order`, so runtime stages inserted mid-pipeline slot into the right place.
*/
function stageOrder(): string[] {
return listPipelineStages().map((s) => s.stageId);
}
/** The stage AFTER `stageId` in the pipeline, or `undefined` if it's terminal. */
export function nextStageAfter(stageId: string): string | undefined {
const order = stageOrder();
const idx = order.indexOf(stageId);
if (idx < 0 || idx >= order.length - 1) return undefined;
return order[idx + 1];
}
export class CeReconciler {
private readonly ctx: PluginContext;
private readonly store: CePipelineStore;
constructor(ctx: PluginContext) {
this.ctx = ctx;
this.store = getCePipelineStore(ctx);
}
/**
* Drain the queue AND re-derive missed transitions from live board state, then
* apply any pipeline advancement (with outbound board propagation). Idempotent:
* running it twice is a no-op once everything has converged.
*/
async reconcile(): Promise<ReconcileResult> {
const result: ReconcileResult = { drained: 0, advanced: 0, tasksCreated: 0, inspected: 0 };
// (1) Drain the queue. Draining is just an audit/ack — the actual decision is
// re-derived from board truth below, so a queue entry for an already-handled
// transition is harmless.
const pending = await this.store.listPendingSyncAsync();
for (const entry of pending) {
await this.store.markSyncProcessedAsync(entry.id);
result.drained++;
}
// (2) Convergence sweep: inspect EVERY pipeline that has state, not only the
// ones with queued entries. This is what recovers a dropped/never-enqueued
// hook event — board truth is compared against pipeline state regardless of
// whether a queue row exists.
const states = await this.store.listAllStateAsync();
for (const state of states) {
result.inspected++;
const advanced = await this.reconcileOne(state);
if (advanced) {
result.advanced++;
if (advanced.created) result.tasksCreated++;
}
}
return result;
}
/**
* Re-derive whether ONE pipeline should advance by reading the live board
* column of its current-stage task(s). Board is authoritative for task state;
* we never write the task column from here for the current stage.
*/
private async reconcileOne(
state: CePipelineState,
): Promise<{ created: boolean } | undefined> {
if (state.status === "completed") return undefined;
// All links for this pipeline (fetched once, reused by advance's idempotency
// check). The board tasks whose completion gates advancement are this
// pipeline's links AT its current stage.
const links = await this.store.listByPipelineAsync(state.cePipelineId);
const currentStageLinks = links.filter((l) => l.ceStageId === state.currentStage);
if (currentStageLinks.length === 0) return undefined;
const tasks = await this.loadTasks(currentStageLinks);
if (tasks.length === 0) return undefined;
// A deleted/missing task yields `undefined` (treated as ABSENT, not
// terminal AND not blocking). Compute terminality over the tasks that still
// EXIST so one deleted current-stage task cannot wedge the pipeline forever.
const existing = tasks.filter((t): t is Task => t != null);
if (existing.length === 0) {
// Every current-stage task was deleted — there is nothing left on the
// board to gate advancement, but also no completion signal to act on.
// Safest non-wedging behavior: leave the pipeline state unchanged (do not
// advance off a vanished stage, do not crash). A later sweep with a real
// board task re-derives the transition.
return undefined;
}
// Advancement rule: every EXISTING current-stage board task has reached a
// terminal column (board-authoritative read). Partial completion keeps it
// running; deleted tasks are excluded above rather than counted as blocking.
const allTerminal = existing.every((t) => TERMINAL_COLUMNS.has(t.column));
if (!allTerminal) {
// Still running on the board — make sure our status reflects that and stop.
if (state.status !== "running") {
await this.store.transitionStateAsync(state.cePipelineId, { status: "running" });
}
return undefined;
}
const next = nextStageAfter(state.currentStage);
if (!next) {
// Terminal stage finished → pipeline completed. No outbound task.
await this.store.transitionStateAsync(state.cePipelineId, { status: "completed" });
this.ctx.emitEvent("compound-engineering:pipeline-completed", {
cePipelineId: state.cePipelineId,
stage: state.currentStage,
});
return { created: false };
}
// CONFLICT POLICY (explicit): board is authoritative for the task columns we
// just READ (we never rewrote them); CE flow is authoritative for the
// pipeline content we WRITE (currentStage, artifact, the next-stage task).
// Advancing only moves the CE-owned fields + creates a NEW board task; it
// never mutates the already-terminal board tasks, so the two writers never
// contend over the same cell.
const created = await this.advance(state, next, links);
return { created };
}
/**
* Advance the pipeline to `nextStage` (CE-owned write) and propagate OUTBOUND
* by creating the next-stage board task (board-owned write on a NEW row).
* Idempotent: if a link for the next stage already exists, we don't duplicate.
*/
private async advance(
state: CePipelineState,
nextStage: string,
links: CePipelineLink[],
): Promise<boolean> {
// Idempotency guard: if we already advanced (a next-stage link exists), just
// ensure state is consistent and skip the outbound create.
const already = links.some((l) => l.ceStageId === nextStage);
if (already) {
await this.store.transitionStateAsync(state.cePipelineId, {
currentStage: nextStage,
status: "running",
});
return false;
}
// Shared contract: create the CE-tagged next-stage board task AND its
// authoritative pipeline-link row (FN-5719) in one place.
const task = await createCeTaskWithLink(this.ctx.taskStore, this.store, {
title: `CE ${nextStage}: continue pipeline`,
description: `Continue the compound-engineering pipeline at the "${nextStage}" stage.`,
cePipelineId: state.cePipelineId,
ceStageId: nextStage,
ceArtifactPath: state.lastArtifactPath,
});
// Single state write on the create path: advance to the next stage and mark
// the pipeline as waiting on the freshly-created board task.
await this.store.transitionStateAsync(state.cePipelineId, {
currentStage: nextStage,
status: "awaiting_board",
});
this.ctx.emitEvent("compound-engineering:pipeline-advanced", {
cePipelineId: state.cePipelineId,
fromStage: state.currentStage,
toStage: nextStage,
taskId: task.id,
});
return true;
}
/** Load the live board tasks for a set of links (board-authoritative read). */
private async loadTasks(links: CePipelineLink[]): Promise<Array<Task | undefined>> {
const out: Array<Task | undefined> = [];
for (const link of links) {
try {
const task = await this.ctx.taskStore.getTask(link.taskId);
out.push(task ?? undefined);
} catch {
// A deleted/missing task is treated as absent, not terminal.
out.push(undefined);
}
}
return out;
}
}
/**
* Convenience: build a reconciler and run one sweep. This is the entry point a
* route handler or post-hook drain calls.
*/
export async function reconcileCePipelines(ctx: PluginContext): Promise<ReconcileResult> {
return new CeReconciler(ctx).reconcile();
}