Files
fusion/packages/core/src/secrets-store.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

427 lines
15 KiB
TypeScript

import { randomUUID } from "node:crypto";
import type { Database as ProjectDatabase } from "./db.js";
import type { CentralDatabase } from "./central-db.js";
import { createSecretCipher, SecretCryptoError, type MasterKeyProvider } from "./secrets-crypto.js";
import type { AsyncDataLayer } from "./postgres/data-layer.js";
import * as asyncSecretsStore from "./async-secrets-store.js";
export type SecretScope = "project" | "global";
export function isSecretScope(value: unknown): value is SecretScope {
return value === "project" || value === "global";
}
export type SecretAccessPolicy = "auto" | "prompt" | "deny";
export interface SecretRecord {
id: string;
key: string;
scope: SecretScope;
description: string | null;
accessPolicy: SecretAccessPolicy;
envExportable: boolean;
envExportKey: string | null;
createdAt: string;
updatedAt: string;
lastReadAt: string | null;
lastReadBy: string | null;
}
export interface EnvExportableSecret {
id: string;
key: string;
exportKey: string;
scope: SecretScope;
plaintextValue: string;
}
interface SecretRow {
id: string;
key: string;
description: string | null;
access_policy: SecretAccessPolicy;
env_exportable: number;
env_export_key: string | null;
created_at: string;
updated_at: string;
last_read_at: string | null;
last_read_by: string | null;
}
interface SecretCipherRow extends SecretRow {
value_ciphertext: Buffer;
nonce: Buffer;
}
type SecretsDb = Pick<ProjectDatabase, "prepare" | "bumpLastModified"> | Pick<CentralDatabase, "prepare" | "bumpLastModified">;
type SecretsStoreAuditEvent = {
mutationType: "secret:create" | "secret:update" | "secret:delete" | "secret:read";
scope: SecretScope;
secretId: string;
key: string;
actor?: { agentId?: string | null; userId?: string | null };
};
export interface SecretsStoreOptions {
/** Optional non-blocking audit emitter. Errors are swallowed/warned so CRUD paths continue. */
auditEmitter?: (event: SecretsStoreAuditEvent) => void;
/**
* FNXC:SecretsStore 2026-06-24-21:00:
* When provided, the store enters backend (PostgreSQL) mode and delegates all
* data access to the async helpers in async-secrets-store.ts. The sync SQLite
* databases (projectDb/centralDb) are ignored in this mode. This is the
* dual-path pattern: the same class serves both SQLite (CLI/desktop) and
* PostgreSQL (backend) deployments.
*/
asyncLayer?: AsyncDataLayer | null;
}
export class SecretsStoreError extends Error {
readonly code: "duplicate-key" | "not-found" | "invalid-policy" | "invalid-key" | "decrypt-failed";
constructor(params: {
code: "duplicate-key" | "not-found" | "invalid-policy" | "invalid-key" | "decrypt-failed";
message: string;
}) {
super(params.message);
this.name = "SecretsStoreError";
this.code = params.code;
}
}
function tableForScope(scope: SecretScope): "secrets" | "secrets_global" {
return scope === "project" ? "secrets" : "secrets_global";
}
function isSqliteUniqueError(error: unknown): boolean {
return error instanceof Error && /UNIQUE constraint failed/u.test(error.message);
}
function isAccessPolicy(value: string): value is SecretAccessPolicy {
return value === "auto" || value === "prompt" || value === "deny";
}
export class SecretsStore {
private readonly cipher: ReturnType<typeof createSecretCipher>;
/**
* FNXC:SecretsStore 2026-06-24-21:05:
* When non-null, the store is in backend (PostgreSQL) mode and all data
* access delegates to the async helpers. The sync projectDb/centralDb are
* not used in this mode.
*/
private readonly asyncLayer: AsyncDataLayer | null;
constructor(
private readonly projectDb: Pick<ProjectDatabase, "prepare" | "bumpLastModified">,
private readonly centralDb: Pick<CentralDatabase, "prepare" | "bumpLastModified">,
masterKeyProvider: MasterKeyProvider,
private readonly options: SecretsStoreOptions = {},
) {
this.cipher = createSecretCipher(masterKeyProvider);
this.asyncLayer = options.asyncLayer ?? null;
}
/** True when the store is backed by PostgreSQL (AsyncDataLayer present). */
private get backendMode(): boolean {
return this.asyncLayer !== null;
}
private emitAudit(event: SecretsStoreAuditEvent): void {
if (!this.options.auditEmitter) return;
try {
this.options.auditEmitter(event);
} catch (error) {
console.warn("[secrets-store] audit emitter failed", error);
}
}
private dbForScope(scope: SecretScope): SecretsDb {
return scope === "project" ? this.projectDb : this.centralDb;
}
private rowToRecord(row: SecretRow, scope: SecretScope): SecretRecord {
return {
id: row.id,
key: row.key,
scope,
description: row.description,
accessPolicy: row.access_policy,
envExportable: row.env_exportable === 1,
envExportKey: row.env_export_key,
createdAt: row.created_at,
updatedAt: row.updated_at,
lastReadAt: row.last_read_at,
lastReadBy: row.last_read_by,
};
}
async listSecrets(scope?: SecretScope): Promise<SecretRecord[]> {
if (this.backendMode) {
return asyncSecretsStore.listSecrets(this.asyncLayer!.db, scope);
}
if (scope) {
const db = this.dbForScope(scope);
const table = tableForScope(scope);
const rows = db.prepare(`SELECT id, key, description, access_policy, env_exportable, env_export_key, created_at, updated_at, last_read_at, last_read_by FROM ${table} ORDER BY key COLLATE NOCASE ASC`).all() as SecretRow[];
return rows.map((row) => this.rowToRecord(row, scope));
}
const [project, global] = await Promise.all([
this.listSecrets("project"),
this.listSecrets("global"),
]);
return [...project, ...global];
}
async listEnvExportable(opts?: { keyPrefix?: string }): Promise<EnvExportableSecret[]> {
const keyPrefix = opts?.keyPrefix;
const projectRows = await this.listSecrets("project");
const globalRows = await this.listSecrets("global");
const exported = new Map<string, EnvExportableSecret>();
const collect = async (row: SecretRecord): Promise<void> => {
if (!row.envExportable) return;
if (keyPrefix && !row.key.startsWith(keyPrefix)) return;
const exportKey = row.envExportKey?.trim() || row.key;
if (exported.has(exportKey)) {
if (row.scope === "global") {
console.debug(`[secrets-store] dropping global env export key due to project override: ${exportKey}`);
}
return;
}
try {
const revealed = await this.revealSecret(row.id, row.scope, {
agentId: null,
userId: "fusion:secrets-env-writer",
});
exported.set(exportKey, {
id: row.id,
key: row.key,
exportKey,
scope: row.scope,
plaintextValue: revealed.plaintextValue,
});
} catch (error) {
const message = error instanceof Error ? error.message : String(error);
console.warn(`[secrets-store] failed to reveal env exportable secret ${row.scope}:${row.key}: ${message}`);
}
};
for (const row of projectRows) {
await collect(row);
}
for (const row of globalRows) {
await collect(row);
}
return [...exported.values()];
}
async getSecretMetadata(id: string, scope: SecretScope): Promise<SecretRecord | null> {
if (this.backendMode) {
return asyncSecretsStore.getSecretMetadata(this.asyncLayer!.db, id, scope);
}
const db = this.dbForScope(scope);
const table = tableForScope(scope);
const row = db.prepare(`SELECT id, key, description, access_policy, env_exportable, env_export_key, created_at, updated_at, last_read_at, last_read_by FROM ${table} WHERE id = ?`).get(id) as SecretRow | undefined;
return row ? this.rowToRecord(row, scope) : null;
}
async createSecret(input: {
scope: SecretScope;
key: string;
plaintextValue: string;
description?: string | null;
accessPolicy?: SecretAccessPolicy;
envExportable?: boolean;
envExportKey?: string | null;
}): Promise<SecretRecord> {
const key = input.key.trim();
if (!key) {
throw new SecretsStoreError({ code: "invalid-key", message: "Secret key is required" });
}
if (input.accessPolicy && !isAccessPolicy(input.accessPolicy)) {
throw new SecretsStoreError({ code: "invalid-policy", message: "Invalid access policy" });
}
if (this.backendMode) {
const created = await asyncSecretsStore.createSecret(this.asyncLayer!.db, this.cipher, input);
this.emitAudit({ mutationType: "secret:create", scope: input.scope, secretId: created.id, key: created.key });
return created;
}
const now = new Date().toISOString();
const id = randomUUID();
const encrypted = await this.cipher.encrypt(input.plaintextValue);
const scope = input.scope;
const db = this.dbForScope(scope);
const table = tableForScope(scope);
try {
db.prepare(`INSERT INTO ${table} (id, key, value_ciphertext, nonce, description, access_policy, env_exportable, env_export_key, created_at, updated_at, last_read_at, last_read_by) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, NULL, NULL)`)
.run(
id,
key,
encrypted.ciphertext,
encrypted.nonce,
input.description ?? null,
input.accessPolicy ?? "auto",
input.envExportable ? 1 : 0,
input.envExportKey ?? null,
now,
now,
);
db.bumpLastModified();
} catch (error) {
if (isSqliteUniqueError(error)) {
throw new SecretsStoreError({ code: "duplicate-key", message: "Secret key already exists" });
}
throw error;
}
const created = (await this.getSecretMetadata(id, scope))!;
this.emitAudit({ mutationType: "secret:create", scope, secretId: created.id, key: created.key });
return created;
}
async updateSecret(id: string, scope: SecretScope, patch: {
key?: string;
plaintextValue?: string;
description?: string | null;
accessPolicy?: SecretAccessPolicy;
envExportable?: boolean;
envExportKey?: string | null;
}): Promise<SecretRecord> {
if (this.backendMode) {
const updated = await asyncSecretsStore.updateSecret(this.asyncLayer!.db, this.cipher, id, scope, patch);
this.emitAudit({ mutationType: "secret:update", scope, secretId: updated.id, key: updated.key });
return updated;
}
const existing = await this.getSecretMetadata(id, scope);
if (!existing) {
throw new SecretsStoreError({ code: "not-found", message: "Secret not found" });
}
const updates: string[] = ["updated_at = ?"];
const params: Array<string | number | Buffer | null> = [new Date().toISOString()];
if (patch.key !== undefined) {
const key = patch.key.trim();
if (!key) {
throw new SecretsStoreError({ code: "invalid-key", message: "Secret key is required" });
}
updates.push("key = ?");
params.push(key);
}
if (patch.description !== undefined) {
updates.push("description = ?");
params.push(patch.description ?? null);
}
if (patch.accessPolicy !== undefined) {
if (!isAccessPolicy(patch.accessPolicy)) {
throw new SecretsStoreError({ code: "invalid-policy", message: "Invalid access policy" });
}
updates.push("access_policy = ?");
params.push(patch.accessPolicy);
}
if (patch.envExportable !== undefined) {
updates.push("env_exportable = ?");
params.push(patch.envExportable ? 1 : 0);
}
if (patch.envExportKey !== undefined) {
updates.push("env_export_key = ?");
params.push(patch.envExportKey ?? null);
}
if (patch.plaintextValue !== undefined) {
const encrypted = await this.cipher.encrypt(patch.plaintextValue);
updates.push("value_ciphertext = ?", "nonce = ?");
params.push(encrypted.ciphertext, encrypted.nonce);
}
const db = this.dbForScope(scope);
const table = tableForScope(scope);
try {
params.push(id);
db.prepare(`UPDATE ${table} SET ${updates.join(", ")} WHERE id = ?`).run(...params);
db.bumpLastModified();
} catch (error) {
if (isSqliteUniqueError(error)) {
throw new SecretsStoreError({ code: "duplicate-key", message: "Secret key already exists" });
}
throw error;
}
const updated = (await this.getSecretMetadata(id, scope))!;
this.emitAudit({ mutationType: "secret:update", scope, secretId: updated.id, key: updated.key });
return updated;
}
async deleteSecret(id: string, scope: SecretScope): Promise<void> {
if (this.backendMode) {
const existing = await this.getSecretMetadata(id, scope);
if (!existing) {
throw new SecretsStoreError({ code: "not-found", message: "Secret not found" });
}
await asyncSecretsStore.deleteSecret(this.asyncLayer!.db, id, scope);
this.emitAudit({ mutationType: "secret:delete", scope, secretId: id, key: existing.key });
return;
}
const existing = await this.getSecretMetadata(id, scope);
if (!existing) {
throw new SecretsStoreError({ code: "not-found", message: "Secret not found" });
}
const db = this.dbForScope(scope);
const table = tableForScope(scope);
db.prepare(`DELETE FROM ${table} WHERE id = ?`).run(id);
db.bumpLastModified();
this.emitAudit({ mutationType: "secret:delete", scope, secretId: id, key: existing.key });
}
async revealSecret(
id: string,
scope: SecretScope,
reader: { agentId?: string | null; userId?: string | null },
): Promise<{ key: string; plaintextValue: string }> {
if (this.backendMode) {
const revealed = await asyncSecretsStore.revealSecret(this.asyncLayer!.db, this.cipher, id, scope, reader);
this.emitAudit({ mutationType: "secret:read", scope, secretId: id, key: revealed.key, actor: reader });
return revealed;
}
const db = this.dbForScope(scope);
const table = tableForScope(scope);
const row = db.prepare(`SELECT id, key, value_ciphertext, nonce, description, access_policy, env_exportable, env_export_key, created_at, updated_at, last_read_at, last_read_by FROM ${table} WHERE id = ?`).get(id) as SecretCipherRow | undefined;
if (!row) {
throw new SecretsStoreError({ code: "not-found", message: "Secret not found" });
}
let plaintextValue: string;
try {
plaintextValue = await this.cipher.decrypt({ ciphertext: row.value_ciphertext, nonce: row.nonce });
} catch (error) {
if (error instanceof SecretCryptoError && error.code === "decryption-failed") {
throw new SecretsStoreError({ code: "decrypt-failed", message: "Secret decryption failed" });
}
throw new SecretsStoreError({ code: "decrypt-failed", message: "Secret decryption failed" });
}
const now = new Date().toISOString();
const lastReadBy = reader.userId ?? reader.agentId ?? null;
db.prepare(`UPDATE ${table} SET last_read_at = ?, last_read_by = ?, updated_at = ? WHERE id = ?`).run(now, lastReadBy, now, id);
db.bumpLastModified();
this.emitAudit({ mutationType: "secret:read", scope, secretId: id, key: row.key, actor: reader });
return { key: row.key, plaintextValue };
}
}