# 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>
241 lines
8.1 KiB
TypeScript
241 lines
8.1 KiB
TypeScript
import { sql } from "drizzle-orm";
|
|
import type { AgentStore } from "./agent-store.js";
|
|
import type { Database } from "./db.js";
|
|
import type { DrizzleDb } from "./postgres/data-layer.js";
|
|
import { PROJECT_SCHEMA } from "./postgres/schema/_shared.js";
|
|
import type { TaskStore } from "./store.js";
|
|
import type { AgentRole } from "./types.js";
|
|
|
|
export interface AgentTokenUsageWindowSummary {
|
|
totalInputTokens: number;
|
|
totalCachedTokens: number;
|
|
totalCacheWriteTokens: number;
|
|
totalOutputTokens: number;
|
|
nTasks: number;
|
|
hitRatio: number;
|
|
}
|
|
|
|
export interface AgentTokenUsageSummary {
|
|
agentId: string;
|
|
role: AgentRole;
|
|
last24h: AgentTokenUsageWindowSummary;
|
|
last7d: AgentTokenUsageWindowSummary;
|
|
allTime: AgentTokenUsageWindowSummary;
|
|
}
|
|
|
|
export interface AgentTaskTokenTotals {
|
|
inputTokens: number;
|
|
cachedTokens: number;
|
|
cacheWriteTokens: number;
|
|
outputTokens: number;
|
|
totalTokens: number;
|
|
nTasks: number;
|
|
}
|
|
|
|
interface TaskTokenLinkRow {
|
|
taskId: string;
|
|
assignedAgentId: string | null;
|
|
sourceAgentId: string | null;
|
|
checkedOutBy: string | null;
|
|
inputTokens: number | null;
|
|
cachedTokens: number | null;
|
|
cacheWriteTokens: number | null;
|
|
outputTokens: number | null;
|
|
totalTokens: number | null;
|
|
}
|
|
|
|
/**
|
|
* FNXC:PostgresCutover 2026-07-04:
|
|
* Shared accumulation core for task-derived token totals by agent link.
|
|
* Both the sync SQLite path and the async PostgreSQL path funnel the same
|
|
* TaskTokenLinkRow[] through this so the assigned/source/checkout attribution
|
|
* stays identical across backends (every linked agent is credited, no double
|
|
* counting when a task links the same agent via multiple fields).
|
|
*/
|
|
function accumulateTaskTokenLinkRows(
|
|
rows: TaskTokenLinkRow[],
|
|
totalsByAgentId: Map<string, AgentTaskTokenTotals>,
|
|
): void {
|
|
for (const row of rows) {
|
|
const agentIds = new Set([row.assignedAgentId, row.sourceAgentId, row.checkedOutBy].filter((value): value is string => Boolean(value)));
|
|
for (const agentId of agentIds) {
|
|
const existing = totalsByAgentId.get(agentId) ?? createTaskTokenTotals();
|
|
existing.inputTokens += row.inputTokens ?? 0;
|
|
existing.cachedTokens += row.cachedTokens ?? 0;
|
|
existing.cacheWriteTokens += row.cacheWriteTokens ?? 0;
|
|
existing.outputTokens += row.outputTokens ?? 0;
|
|
existing.totalTokens += row.totalTokens ?? (row.inputTokens ?? 0) + (row.cachedTokens ?? 0) + (row.cacheWriteTokens ?? 0) + (row.outputTokens ?? 0);
|
|
existing.nTasks += 1;
|
|
totalsByAgentId.set(agentId, existing);
|
|
}
|
|
}
|
|
}
|
|
|
|
export function aggregateTaskTokenTotalsByAgentLink(db: Database): Map<string, AgentTaskTokenTotals> {
|
|
/*
|
|
FNXC:AgentTokenUsage 2026-06-27-23:06:
|
|
List-row token totals must use the same assigned/source/checkout attribution as Agent Detail so ephemeral task-worker agents do not report zero when they only sourced or checked out a task.
|
|
*/
|
|
const rows = db.prepare(`
|
|
SELECT
|
|
id AS taskId,
|
|
assignedAgentId,
|
|
sourceAgentId,
|
|
checkedOutBy,
|
|
tokenUsageInputTokens AS inputTokens,
|
|
tokenUsageCachedTokens AS cachedTokens,
|
|
tokenUsageCacheWriteTokens AS cacheWriteTokens,
|
|
tokenUsageOutputTokens AS outputTokens,
|
|
tokenUsageTotalTokens AS totalTokens
|
|
FROM tasks
|
|
WHERE tokenUsageInputTokens IS NOT NULL
|
|
OR tokenUsageCachedTokens IS NOT NULL
|
|
OR tokenUsageCacheWriteTokens IS NOT NULL
|
|
OR tokenUsageOutputTokens IS NOT NULL
|
|
OR tokenUsageTotalTokens IS NOT NULL
|
|
`).all() as TaskTokenLinkRow[];
|
|
|
|
const totalsByAgentId = new Map<string, AgentTaskTokenTotals>();
|
|
accumulateTaskTokenLinkRows(rows, totalsByAgentId);
|
|
return totalsByAgentId;
|
|
}
|
|
|
|
/**
|
|
* FNXC:PostgresCutover 2026-07-04:
|
|
* Async PostgreSQL equivalent of {@link aggregateTaskTokenTotalsByAgentLink}.
|
|
* Reads the same task-link token columns from the schema-qualified
|
|
* project.tasks table (snake_case) and runs them through the shared
|
|
* accumulator so the Agents listing surfaces real task-derived totals in
|
|
* backend (PostgreSQL) mode instead of silently degrading to zero. The route
|
|
* dispatches to this when the scoped store carries an AsyncDataLayer.
|
|
*/
|
|
export async function aggregateTaskTokenTotalsByAgentLinkAsync(
|
|
db: DrizzleDb,
|
|
): Promise<Map<string, AgentTaskTokenTotals>> {
|
|
const rawRows = (await db.execute(
|
|
sql.raw(`
|
|
SELECT
|
|
id AS "taskId",
|
|
assigned_agent_id AS "assignedAgentId",
|
|
source_agent_id AS "sourceAgentId",
|
|
checked_out_by AS "checkedOutBy",
|
|
token_usage_input_tokens AS "inputTokens",
|
|
token_usage_cached_tokens AS "cachedTokens",
|
|
token_usage_cache_write_tokens AS "cacheWriteTokens",
|
|
token_usage_output_tokens AS "outputTokens",
|
|
token_usage_total_tokens AS "totalTokens"
|
|
FROM ${PROJECT_SCHEMA}.tasks
|
|
WHERE token_usage_input_tokens IS NOT NULL
|
|
OR token_usage_cached_tokens IS NOT NULL
|
|
OR token_usage_cache_write_tokens IS NOT NULL
|
|
OR token_usage_output_tokens IS NOT NULL
|
|
OR token_usage_total_tokens IS NOT NULL
|
|
`),
|
|
)) as unknown as TaskTokenLinkRow[];
|
|
|
|
const totalsByAgentId = new Map<string, AgentTaskTokenTotals>();
|
|
accumulateTaskTokenLinkRows(rawRows, totalsByAgentId);
|
|
return totalsByAgentId;
|
|
}
|
|
|
|
function createTaskTokenTotals(): AgentTaskTokenTotals {
|
|
return {
|
|
inputTokens: 0,
|
|
cachedTokens: 0,
|
|
cacheWriteTokens: 0,
|
|
outputTokens: 0,
|
|
totalTokens: 0,
|
|
nTasks: 0,
|
|
};
|
|
}
|
|
|
|
export async function aggregateAgentTokenUsage({
|
|
taskStore,
|
|
agentStore,
|
|
agentId,
|
|
now = new Date(),
|
|
}: {
|
|
taskStore: TaskStore;
|
|
agentStore: AgentStore;
|
|
agentId: string;
|
|
now?: Date;
|
|
}): Promise<AgentTokenUsageSummary | null> {
|
|
const agent = await agentStore.getAgent(agentId);
|
|
if (!agent) {
|
|
return null;
|
|
}
|
|
|
|
/*
|
|
FNXC:AgentTokenUsage 2026-06-27-19:10:
|
|
Ephemeral/task-worker agents must surface task-derived token usage because their cumulative agent token fields are never accumulated by the durable-agent heartbeat path.
|
|
*/
|
|
const tasks = await taskStore.listTasks({ slim: true, includeArchived: true });
|
|
const nowMs = now.getTime();
|
|
const last24hMs = nowMs - (24 * 60 * 60 * 1000);
|
|
const last7dMs = nowMs - (7 * 24 * 60 * 60 * 1000);
|
|
|
|
const allTime = createWindowSummary();
|
|
const last24h = createWindowSummary();
|
|
const last7d = createWindowSummary();
|
|
|
|
for (const task of tasks) {
|
|
if (!task.tokenUsage) continue;
|
|
const matchesAgent = task.assignedAgentId === agentId || task.sourceAgentId === agentId || task.checkedOutBy === agentId;
|
|
if (!matchesAgent) continue;
|
|
|
|
const usage = task.tokenUsage;
|
|
applyTaskUsage(allTime, usage.inputTokens ?? 0, usage.cachedTokens ?? 0, usage.outputTokens ?? 0, usage.cacheWriteTokens ?? 0);
|
|
|
|
const lastUsedAtMs = Date.parse(usage.lastUsedAt ?? "");
|
|
if (!Number.isFinite(lastUsedAtMs)) continue;
|
|
|
|
if (lastUsedAtMs >= last24hMs) {
|
|
applyTaskUsage(last24h, usage.inputTokens ?? 0, usage.cachedTokens ?? 0, usage.outputTokens ?? 0, usage.cacheWriteTokens ?? 0);
|
|
}
|
|
if (lastUsedAtMs >= last7dMs) {
|
|
applyTaskUsage(last7d, usage.inputTokens ?? 0, usage.cachedTokens ?? 0, usage.outputTokens ?? 0, usage.cacheWriteTokens ?? 0);
|
|
}
|
|
}
|
|
|
|
return {
|
|
agentId,
|
|
role: agent.role as AgentRole,
|
|
last24h: finalizeWindowSummary(last24h),
|
|
last7d: finalizeWindowSummary(last7d),
|
|
allTime: finalizeWindowSummary(allTime),
|
|
};
|
|
}
|
|
|
|
function createWindowSummary(): AgentTokenUsageWindowSummary {
|
|
return {
|
|
totalInputTokens: 0,
|
|
totalCachedTokens: 0,
|
|
totalCacheWriteTokens: 0,
|
|
totalOutputTokens: 0,
|
|
nTasks: 0,
|
|
hitRatio: 0,
|
|
};
|
|
}
|
|
|
|
function applyTaskUsage(
|
|
summary: AgentTokenUsageWindowSummary,
|
|
inputTokens: number,
|
|
cachedTokens: number,
|
|
outputTokens: number,
|
|
cacheWriteTokens: number,
|
|
): void {
|
|
summary.totalInputTokens += inputTokens;
|
|
summary.totalCachedTokens += cachedTokens;
|
|
summary.totalCacheWriteTokens += cacheWriteTokens;
|
|
summary.totalOutputTokens += outputTokens;
|
|
summary.nTasks += 1;
|
|
}
|
|
|
|
function finalizeWindowSummary(summary: AgentTokenUsageWindowSummary): AgentTokenUsageWindowSummary {
|
|
const denominator = summary.totalInputTokens + summary.totalCachedTokens;
|
|
return {
|
|
...summary,
|
|
hitRatio: denominator > 0 ? summary.totalCachedTokens / denominator : 0,
|
|
};
|
|
}
|