# Remove dead SQLite dual-path code; keep migration-only readers ## Summary PostgreSQL cutover left hundreds of production dual-path branches (`backendMode ? PG : SQLite/store.db`) whose SQLite arms only hit throwing `Database`/`ArchiveDatabase`/`CentralDatabase` stubs. This change mechanically collapses those unreachable arms so production authority is AsyncDataLayer/PostgreSQL only, while preserving the six authorized read-only migration/recovery `DatabaseSync` seams. ## Dual-path mass removed | Metric | Before | After | |---|---|---| | `if (…backendMode)` (non-test) | ~328 | ~70 | | `store.db` / `this.db` refs in core (non-test) | ~570+ | ~375 (mostly pure legacy MissionStore/eval/insight SQLite classes + thin getters) | | Net diff | — | **~6.7k lines removed** across 41 files | Remaining `backendMode` checks are intentional (incomplete-PG sync safe-defaults, settings-sync disabled-on-PG, symbol-lock PG-only gates, “requires PostgreSQL” config versioning throws), not live SQLite authority. ## Subsystems cleaned - **Core TaskStore / task-store/***: collapsed if/else and early-return dual-path across reads, moves, lifecycle, mutations, workflow, archive, branch/PR, artifacts, comments, audit, project ops, etc. `initImpl` is PostgreSQL-only (SQLite startup tail deleted). - **Satellite stores**: automation, agent, routine, plugin, secrets, approval-request, central-core dual-path arms collapsed. - **Plugins**: reports async methods, compound-engineering pipeline + session stores, CLI Printing Press store — SQLite fallbacks removed; PG required. - **Engine**: no functional dual-path change beyond whitespace (settings-sync / peer-exchange PG-disabled behavior kept). ## Six migration-only readers retained (allowlist unchanged) 1. `packages/core/src/postgres/sqlite-migrator.ts` 2. `packages/core/src/project-identity.ts` 3. `packages/core/src/sqlite-validation.ts` 4. `packages/core/src/postgres/startup-factory.ts` 5. `packages/cli/src/commands/db.ts` 6. `scripts/lib/start-local-project.mjs` Plus low-level `sqlite-adapter` and migrator/startup-import tests. Inventory ratchet still requires exactly these six `new DatabaseSync(` production sites, all `readOnly: true`. ## Not treated as SQLite - `.fusion/project.json`, `task.json`, `agent-log.jsonl` file storage - AsyncDataLayer / Drizzle PG paths - Incomplete-PG sync safe-default stubs (still return empty/false/null under backend without consulting SQLite) ## Verification - `sqlite-production-reader-inventory.test.ts` — 15/15 pass - `incomplete-pg-ports.pg.test.ts` — 6/6 pass - Targeted PG tests (create-task, move, handoff, runtime-persistence, agent, mission, insight, central-core) — green - `tsc --noEmit` for `@fusion/core`, `@fusion/engine`, `@fusion/dashboard` — green - `scripts/check-no-getdatabase.mjs` — clean <!-- This is an auto-generated comment: release notes by coderabbit.ai --> ## Summary by CodeRabbit * **Improvements** * Improved end-to-end consistency by making PostgreSQL/async persistence the standard across core task/workflow, automation, agents, plugins, routines, secrets, approvals, central operations, and session storage. * Unified scheduling, settings, configuration revision writes, run/workflow selection, queues/leases/transitions, and audit/lifecycle updates around consistent async transaction behavior. * **Bug Fixes** * Fixed edge cases for archived/deleted reads, unarchive/recovery flows, not-found handling, and task/artifact/document/log/comment operations, including more reliable emissions and hydration across search/list and lifecycle operations. <!-- end of auto-generated comment: release notes by coderabbit.ai -->
551 lines
25 KiB
TypeScript
551 lines
25 KiB
TypeScript
import { randomUUID } from "node:crypto";
|
||
import { type AsyncDataLayer, type Database, type PlanningQuestion, type PluginContext } from "@fusion/core";
|
||
import { sql } from "drizzle-orm";
|
||
/* FNXC:CompoundEngineeringPostgres 2026-07-13-23:42: Import SQL construction from Drizzle directly because the CLI's bundled-plugin @fusion/core shim does not expose database query builders. */
|
||
import { ensureCeSchema } from "../schema.js";
|
||
|
||
/**
|
||
* CE session lifecycle states (mirrors the plan's state machine):
|
||
* launching → active → awaiting_input ↔ active → completed | error | interrupted;
|
||
* interrupted/error → active on resume/retry.
|
||
*/
|
||
export const CE_SESSION_STATUSES = [
|
||
"launching",
|
||
"active",
|
||
"awaiting_input",
|
||
"completed",
|
||
"error",
|
||
"interrupted",
|
||
] as const;
|
||
|
||
export type CeSessionStatus = (typeof CE_SESSION_STATUSES)[number];
|
||
|
||
/** Narrow an arbitrary string (e.g. a query param) to a valid status, else undefined. */
|
||
export function asCeSessionStatus(value: string | undefined): CeSessionStatus | undefined {
|
||
return value && (CE_SESSION_STATUSES as readonly string[]).includes(value)
|
||
? (value as CeSessionStatus)
|
||
: undefined;
|
||
}
|
||
|
||
/** A single recorded turn in the conversation history (for resume). */
|
||
export interface CeConversationTurn {
|
||
role: "user" | "agent";
|
||
/** Free text, or a serialized question/answer marker. */
|
||
text: string;
|
||
at: string;
|
||
}
|
||
|
||
/**
|
||
* One line of live agent activity (mid-turn working output): an accumulated
|
||
* thinking/text block or a discrete tool execution marker.
|
||
*/
|
||
export interface CeActivityTurn {
|
||
kind: "thinking" | "text" | "tool";
|
||
text: string;
|
||
at: string;
|
||
/** Tool turns: execution finished. */
|
||
done?: boolean;
|
||
/** Tool turns: execution finished with an error. */
|
||
isError?: boolean;
|
||
}
|
||
|
||
export interface CeSession {
|
||
id: string;
|
||
stage: string;
|
||
status: CeSessionStatus;
|
||
currentQuestion: PlanningQuestion | null;
|
||
conversationHistory: CeConversationTurn[];
|
||
/**
|
||
* TRANSIENT: in-flight working output for the current turn, attached by the
|
||
* GET-session route from the orchestrator's in-memory buffer. Never persisted
|
||
* to the row; absent when no turn is running (or in another process).
|
||
*/
|
||
liveActivity?: CeActivityTurn[];
|
||
projectId: string | null;
|
||
artifactPath: string | null;
|
||
error: string | null;
|
||
/** Expected per-turn interval (ms); drives interval-relative staleness. */
|
||
turnIntervalMs: number;
|
||
/** Epoch millis of the last produced event (liveness anchor). */
|
||
lastActivityAt: number;
|
||
createdAt: string;
|
||
updatedAt: string;
|
||
}
|
||
|
||
interface CeSessionRow {
|
||
id: string;
|
||
stage: string;
|
||
status: CeSessionStatus;
|
||
currentQuestion: string | null;
|
||
conversationHistory: string;
|
||
projectId: string | null;
|
||
artifactPath: string | null;
|
||
error: string | null;
|
||
turnIntervalMs: number;
|
||
lastActivityAt: number;
|
||
createdAt: string;
|
||
updatedAt: string;
|
||
}
|
||
|
||
export interface CreateCeSessionInput {
|
||
stage: string;
|
||
projectId?: string | null;
|
||
artifactPath?: string | null;
|
||
turnIntervalMs?: number;
|
||
id?: string;
|
||
}
|
||
|
||
export class PlanHandoffClaimError extends Error {
|
||
constructor(
|
||
readonly artifactPath: string,
|
||
readonly sessionId: string,
|
||
) {
|
||
super(`Plan session ${sessionId} is already enriching ${artifactPath}`);
|
||
this.name = "PlanHandoffClaimError";
|
||
}
|
||
}
|
||
|
||
function isUniqueConstraintError(error: unknown): boolean {
|
||
let current: unknown = error;
|
||
const seen = new Set<unknown>();
|
||
while (current && typeof current === "object" && !seen.has(current)) {
|
||
seen.add(current);
|
||
const candidate = current as { code?: unknown; message?: unknown; cause?: unknown };
|
||
if (candidate.code === "23505") return true;
|
||
if (typeof candidate.message === "string" && /duplicate key|unique constraint/i.test(candidate.message)) return true;
|
||
current = candidate.cause;
|
||
}
|
||
return false;
|
||
}
|
||
|
||
/**
|
||
* Default multiple of the turn interval beyond which a non-terminal session is
|
||
* considered stale. Mirrors the FN-4172 rubric (`> 3× interval`), interval-
|
||
* relative rather than a raw last-event age.
|
||
*/
|
||
export const STALE_INTERVAL_MULTIPLE = 3;
|
||
|
||
const DEFAULT_TURN_INTERVAL_MS = 120000;
|
||
|
||
/**
|
||
* Parse a JSON column, falling back to `fallback` when it is missing, fails to
|
||
* parse (syntax error), OR parses to the wrong shape. Shape validation matters:
|
||
* a column holding `'null'` or `'{}'` parses fine but would yield a non-array
|
||
* `conversationHistory` that later crashes `appendHistory`'s spread — so a
|
||
* semantically-corrupt value is treated exactly like a syntactically-corrupt one.
|
||
*/
|
||
function safeParse<T>(raw: string | null, fallback: T, isValid: (value: unknown) => value is T): T {
|
||
if (!raw) return fallback;
|
||
try {
|
||
const parsed: unknown = JSON.parse(raw);
|
||
return isValid(parsed) ? parsed : fallback;
|
||
} catch {
|
||
// A corrupted JSON column must not crash reads of an otherwise-valid row
|
||
// (and must not destroy the rest of the session). Degrade to the fallback;
|
||
// the row's status/error still surface the session's real state.
|
||
return fallback;
|
||
}
|
||
}
|
||
|
||
function isConversationHistory(value: unknown): value is CeConversationTurn[] {
|
||
return (
|
||
Array.isArray(value)
|
||
&& value.every((turn) => {
|
||
if (typeof turn !== "object" || turn === null) return false;
|
||
const t = turn as Record<string, unknown>;
|
||
return (t.role === "user" || t.role === "agent") && typeof t.text === "string" && typeof t.at === "string";
|
||
})
|
||
);
|
||
}
|
||
|
||
function isPlanningQuestionOrNull(value: unknown): value is PlanningQuestion | null {
|
||
if (value === null) return true;
|
||
if (typeof value !== "object") return false;
|
||
const q = value as Record<string, unknown>;
|
||
return typeof q.id === "string" && typeof q.type === "string" && typeof q.question === "string";
|
||
}
|
||
|
||
function rowToSession(row: CeSessionRow): CeSession {
|
||
return {
|
||
id: row.id,
|
||
stage: row.stage,
|
||
status: row.status,
|
||
currentQuestion: safeParse<PlanningQuestion | null>(row.currentQuestion, null, isPlanningQuestionOrNull),
|
||
conversationHistory: safeParse<CeConversationTurn[]>(row.conversationHistory, [], isConversationHistory),
|
||
projectId: row.projectId,
|
||
artifactPath: row.artifactPath,
|
||
error: row.error,
|
||
turnIntervalMs: row.turnIntervalMs,
|
||
lastActivityAt: Number(row.lastActivityAt),
|
||
createdAt: row.createdAt,
|
||
updatedAt: row.updatedAt,
|
||
};
|
||
}
|
||
|
||
/**
|
||
* Plugin-local persistence for CE interactive sessions. Reaches the DB the same
|
||
* way reports does (via `ctx.taskStore.getAsyncLayer()`), and ensures its schema
|
||
* defensively on construction so a store created before `onSchemaInit` ran (or
|
||
* in a test) still works.
|
||
*/
|
||
export class CeSessionStore {
|
||
// FNXC:RuntimeSatelliteAsync 2026-06-24-22:45:
|
||
// db is null in backend mode (PostgreSQL). Store methods that use sync
|
||
// SQLite will throw in backend mode until the async path is implemented.
|
||
private readonly db: Database | null;
|
||
|
||
constructor(db: Database | null, private readonly asyncLayer: AsyncDataLayer | null = null) {
|
||
this.db = db;
|
||
if (db) ensureCeSchema(db);
|
||
}
|
||
|
||
/** Asserts sync db is available (throws in backend mode). */
|
||
private syncDb(): Database {
|
||
if (!this.db) throw new Error("CeSessionStore: sync Database is null (backend mode)");
|
||
return this.db;
|
||
}
|
||
|
||
create(input: CreateCeSessionInput): CeSession {
|
||
const session = this.newSession(input);
|
||
this.insert(session);
|
||
return session;
|
||
}
|
||
|
||
/*
|
||
* FNXC:CompoundEngineeringPlanning 2026-07-11-00:18:
|
||
* A requirements artifact can have exactly one Plan owner until that session is discarded. Claim it in the same immediate transaction as the Plan row so concurrent dashboard/API requests cannot both enrich the same document; discarding that session intentionally releases the claim for a retry.
|
||
*/
|
||
createWithPlanHandoffClaim(input: CreateCeSessionInput, artifactPath: string): CeSession {
|
||
const session = this.newSession({ ...input, artifactPath });
|
||
const db = this.syncDb();
|
||
return db.transactionImmediate(() => {
|
||
const existing = db
|
||
.prepare("SELECT sessionId FROM ce_plan_handoff_claims WHERE artifactPath = ?")
|
||
.get(artifactPath) as { sessionId: string } | undefined;
|
||
if (existing) throw new PlanHandoffClaimError(artifactPath, existing.sessionId);
|
||
|
||
this.insert(session);
|
||
db
|
||
.prepare("INSERT INTO ce_plan_handoff_claims (artifactPath, sessionId, projectId, createdAt) VALUES (?, ?, ?, ?)")
|
||
.run(artifactPath, session.id, session.projectId, session.createdAt);
|
||
return session;
|
||
});
|
||
}
|
||
|
||
private newSession(input: CreateCeSessionInput): CeSession {
|
||
const now = new Date().toISOString();
|
||
return {
|
||
id: input.id ?? randomUUID(),
|
||
stage: input.stage,
|
||
status: "launching",
|
||
currentQuestion: null,
|
||
conversationHistory: [],
|
||
projectId: input.projectId ?? null,
|
||
artifactPath: input.artifactPath ?? null,
|
||
error: null,
|
||
turnIntervalMs: input.turnIntervalMs ?? DEFAULT_TURN_INTERVAL_MS,
|
||
lastActivityAt: Date.now(),
|
||
createdAt: now,
|
||
updatedAt: now,
|
||
};
|
||
}
|
||
|
||
private insert(session: CeSession): void {
|
||
this.syncDb()
|
||
.prepare(
|
||
`INSERT INTO ce_sessions
|
||
(id, stage, status, currentQuestion, conversationHistory, projectId, artifactPath, error, turnIntervalMs, lastActivityAt, createdAt, updatedAt)
|
||
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`,
|
||
)
|
||
.run(
|
||
session.id,
|
||
session.stage,
|
||
session.status,
|
||
null,
|
||
JSON.stringify(session.conversationHistory),
|
||
session.projectId,
|
||
session.artifactPath,
|
||
null,
|
||
session.turnIntervalMs,
|
||
session.lastActivityAt,
|
||
session.createdAt,
|
||
session.updatedAt,
|
||
);
|
||
}
|
||
|
||
get(id: string): CeSession | undefined {
|
||
const row = this.syncDb().prepare(`SELECT * FROM ce_sessions WHERE id = ?`).get(id) as CeSessionRow | undefined;
|
||
return row ? rowToSession(row) : undefined;
|
||
}
|
||
|
||
list(filter: { status?: CeSessionStatus; stage?: string; projectId?: string } = {}): CeSession[] {
|
||
const clauses: string[] = [];
|
||
const params: unknown[] = [];
|
||
if (filter.status) {
|
||
clauses.push("status = ?");
|
||
params.push(filter.status);
|
||
}
|
||
if (filter.stage) {
|
||
clauses.push("stage = ?");
|
||
params.push(filter.stage);
|
||
}
|
||
if (filter.projectId) {
|
||
clauses.push("projectId = ?");
|
||
params.push(filter.projectId);
|
||
}
|
||
const where = clauses.length > 0 ? `WHERE ${clauses.join(" AND ")}` : "";
|
||
const rows = this.syncDb()
|
||
.prepare(`SELECT * FROM ce_sessions ${where} ORDER BY updatedAt DESC, id`)
|
||
.all(...params) as CeSessionRow[];
|
||
return rows.map(rowToSession);
|
||
}
|
||
|
||
/**
|
||
* Patch a session. Always bumps `updatedAt`; bumps `lastActivityAt` unless the
|
||
* caller explicitly overrides it (used by liveness tests to simulate age).
|
||
*/
|
||
update(
|
||
id: string,
|
||
patch: Partial<
|
||
Pick<
|
||
CeSession,
|
||
"status" | "currentQuestion" | "conversationHistory" | "artifactPath" | "error" | "lastActivityAt" | "projectId"
|
||
>
|
||
>,
|
||
): CeSession | undefined {
|
||
const existing = this.get(id);
|
||
if (!existing) return undefined;
|
||
const next: CeSession = {
|
||
...existing,
|
||
...patch,
|
||
lastActivityAt: patch.lastActivityAt ?? Date.now(),
|
||
updatedAt: new Date().toISOString(),
|
||
};
|
||
this.syncDb()
|
||
.prepare(
|
||
`UPDATE ce_sessions SET
|
||
status = ?, currentQuestion = ?, conversationHistory = ?, projectId = ?,
|
||
artifactPath = ?, error = ?, lastActivityAt = ?, updatedAt = ?
|
||
WHERE id = ?`,
|
||
)
|
||
.run(
|
||
next.status,
|
||
next.currentQuestion ? JSON.stringify(next.currentQuestion) : null,
|
||
JSON.stringify(next.conversationHistory),
|
||
next.projectId,
|
||
next.artifactPath,
|
||
next.error,
|
||
next.lastActivityAt,
|
||
next.updatedAt,
|
||
id,
|
||
);
|
||
return next;
|
||
}
|
||
|
||
/** Delete a session row. Returns true when a row was removed. */
|
||
delete(id: string): boolean {
|
||
const db = this.syncDb();
|
||
return db.transactionImmediate(() => {
|
||
const result = db.prepare(`DELETE FROM ce_sessions WHERE id = ?`).run(id);
|
||
if (Number(result.changes ?? 0) > 0) {
|
||
db.prepare("DELETE FROM ce_plan_handoff_claims WHERE sessionId = ?").run(id);
|
||
return true;
|
||
}
|
||
return false;
|
||
});
|
||
}
|
||
|
||
/** Append a turn to the conversation history (no other field touched). */
|
||
appendHistory(id: string, turn: CeConversationTurn): CeSession | undefined {
|
||
const existing = this.get(id);
|
||
if (!existing) return undefined;
|
||
return this.update(id, { conversationHistory: [...existing.conversationHistory, turn] });
|
||
}
|
||
|
||
/**
|
||
* Interval-relative staleness: a non-terminal session is stale only when its
|
||
* last activity is older than `multiple × turnIntervalMs`. A healthy-but-slow
|
||
* session (within the interval band) is NOT stale. Terminal sessions
|
||
* (completed/error/interrupted) are never "stale" — they are already settled.
|
||
*/
|
||
isStale(session: CeSession, now = Date.now(), multiple = STALE_INTERVAL_MULTIPLE): boolean {
|
||
if (session.status === "completed" || session.status === "error" || session.status === "interrupted") {
|
||
return false;
|
||
}
|
||
return now - session.lastActivityAt > multiple * session.turnIntervalMs;
|
||
}
|
||
|
||
/**
|
||
* Recover sessions left non-terminal by a crash/restart. A session with a
|
||
* persisted `currentQuestion` is restored to `awaiting_input` (resumable);
|
||
* one without is marked `interrupted` with its progress preserved — never
|
||
* silently dropped. Returns the ids transitioned.
|
||
*/
|
||
recoverStaleSessions(now = Date.now(), multiple = STALE_INTERVAL_MULTIPLE): string[] {
|
||
// Only IN-FLIGHT agent turns are subject to the interval-staleness rubric:
|
||
// `active`/`launching` mean an agent turn should be progressing, so exceeding
|
||
// the interval band signals a crashed/abandoned turn worth recovering.
|
||
//
|
||
// `awaiting_input` is DELIBERATELY excluded: a session waiting on human input
|
||
// is not a crashed turn — human response time is unbounded, and the interval
|
||
// rubric measures agent turns, not human waits. Flagging it stale would
|
||
// misclassify a legitimately-paused session. It is already in its resumable
|
||
// state, so no recovery action is needed.
|
||
const candidates = this.list().filter(
|
||
(s) => (s.status === "active" || s.status === "launching") && this.isStale(s, now, multiple),
|
||
);
|
||
const recovered: string[] = [];
|
||
for (const s of candidates) {
|
||
if (s.currentQuestion) {
|
||
this.update(s.id, { status: "awaiting_input" });
|
||
} else {
|
||
this.update(s.id, { status: "interrupted", error: s.error ?? "Session interrupted — progress preserved, resume to continue" });
|
||
}
|
||
recovered.push(s.id);
|
||
}
|
||
return recovered;
|
||
}
|
||
|
||
/**
|
||
* FNXC:CompoundEngineeringPostgresPersistence 2026-07-13-22:37:
|
||
* Session orchestration uses these async siblings so backend mode persists through the project-bound AsyncDataLayer while SQLite callers retain their established synchronous API. PostgreSQL reads ignore caller-supplied cross-project filters and always enforce the layer's project ID.
|
||
*/
|
||
async createAsync(input: CreateCeSessionInput): Promise<CeSession> {
|
||
const projectId = this.requireProjectId();
|
||
const session = this.newSession({ ...input, projectId });
|
||
await this.insertAsync(session);
|
||
return session;
|
||
}
|
||
|
||
async createWithPlanHandoffClaimAsync(input: CreateCeSessionInput, artifactPath: string): Promise<CeSession> {
|
||
const projectId = this.requireProjectId();
|
||
const session = this.newSession({ ...input, projectId, artifactPath });
|
||
try {
|
||
await this.requireLayer().transactionImmediate(async (tx) => {
|
||
await tx.execute(sql`INSERT INTO project.ce_sessions
|
||
(id, stage, status, current_question, conversation_history, project_id, artifact_path, error, turn_interval_ms, last_activity_at, created_at, updated_at)
|
||
VALUES(${session.id}, ${session.stage}, ${session.status}, NULL, '[]', ${projectId}, ${artifactPath}, NULL, ${session.turnIntervalMs}, ${session.lastActivityAt}, ${session.createdAt}, ${session.updatedAt})`);
|
||
await tx.execute(sql`INSERT INTO project.ce_plan_handoff_claims(project_id, artifact_path, session_id, created_at)
|
||
VALUES(${projectId}, ${artifactPath}, ${session.id}, ${session.createdAt})`);
|
||
});
|
||
} catch (error) {
|
||
if (isUniqueConstraintError(error)) {
|
||
const rows = await this.requireLayer().db.execute(sql`SELECT session_id FROM project.ce_plan_handoff_claims WHERE project_id=${projectId} AND artifact_path=${artifactPath} LIMIT 1`) as unknown as Array<{ session_id: string }>;
|
||
throw new PlanHandoffClaimError(artifactPath, rows[0]?.session_id ?? "unknown");
|
||
}
|
||
throw error;
|
||
}
|
||
return session;
|
||
}
|
||
|
||
/*
|
||
FNXC:SqliteDualPathCleanup 2026-07-26-15:20:
|
||
Async session methods require AsyncDataLayer after SQLite dual-path removal; assert for TS null narrowing.
|
||
*/
|
||
private requireLayer() {
|
||
if (!this.asyncLayer) throw new Error("CE PostgreSQL persistence requires a project-bound data layer");
|
||
return this.asyncLayer;
|
||
}
|
||
|
||
private requireProjectId(): string {
|
||
const projectId = this.requireLayer().projectId?.trim();
|
||
if (!projectId) throw new Error("CE PostgreSQL persistence requires a project-bound data layer");
|
||
return projectId;
|
||
}
|
||
|
||
private async insertAsync(session: CeSession): Promise<void> {
|
||
await this.requireLayer().db.execute(sql`INSERT INTO project.ce_sessions
|
||
(id, stage, status, current_question, conversation_history, project_id, artifact_path, error, turn_interval_ms, last_activity_at, created_at, updated_at)
|
||
VALUES(${session.id}, ${session.stage}, ${session.status}, NULL, ${JSON.stringify(session.conversationHistory)}, ${session.projectId}, ${session.artifactPath}, NULL, ${session.turnIntervalMs}, ${session.lastActivityAt}, ${session.createdAt}, ${session.updatedAt})`);
|
||
}
|
||
|
||
async getAsync(id: string): Promise<CeSession | undefined> {
|
||
const rows = await this.requireLayer().db.execute(sql`SELECT id, stage, status, current_question AS "currentQuestion", conversation_history AS "conversationHistory", project_id AS "projectId", artifact_path AS "artifactPath", error, turn_interval_ms AS "turnIntervalMs", last_activity_at AS "lastActivityAt", created_at AS "createdAt", updated_at AS "updatedAt" FROM project.ce_sessions WHERE project_id=${this.requireProjectId()} AND id=${id} LIMIT 1`) as unknown as CeSessionRow[];
|
||
return rows[0] ? rowToSession(rows[0]) : undefined;
|
||
}
|
||
|
||
async listAsync(filter: { status?: CeSessionStatus; stage?: string; projectId?: string } = {}): Promise<CeSession[]> {
|
||
const projectId = this.requireProjectId();
|
||
const rows = await this.requireLayer().db.execute(sql`SELECT id, stage, status, current_question AS "currentQuestion", conversation_history AS "conversationHistory", project_id AS "projectId", artifact_path AS "artifactPath", error, turn_interval_ms AS "turnIntervalMs", last_activity_at AS "lastActivityAt", created_at AS "createdAt", updated_at AS "updatedAt" FROM project.ce_sessions WHERE project_id=${projectId} AND (${filter.status ?? null}::text IS NULL OR status=${filter.status ?? null}) AND (${filter.stage ?? null}::text IS NULL OR stage=${filter.stage ?? null}) ORDER BY updated_at DESC, id`) as unknown as CeSessionRow[];
|
||
return rows.map(rowToSession);
|
||
}
|
||
|
||
async updateAsync(id: string, patch: Partial<Pick<CeSession, "status" | "currentQuestion" | "conversationHistory" | "artifactPath" | "error" | "lastActivityAt" | "projectId">>): Promise<CeSession | undefined> {
|
||
const projectId = this.requireProjectId();
|
||
const updatedAt = new Date().toISOString();
|
||
const lastActivityAt = patch.lastActivityAt ?? Date.now();
|
||
const has = (key: keyof typeof patch): boolean => Object.prototype.hasOwnProperty.call(patch, key);
|
||
/*
|
||
* FNXC:CompoundEngineeringConcurrency 2026-07-14-00:18:
|
||
* PostgreSQL session patches must update only the fields named by the caller. A prior read-modify-write rewrote the entire row, allowing a delayed heartbeat to restore stale history or a pre-terminal status after a question/completion had committed.
|
||
*/
|
||
const rows = await this.requireLayer().db.execute(sql`
|
||
UPDATE project.ce_sessions SET
|
||
status = CASE WHEN ${has("status")} THEN ${patch.status ?? null}::text ELSE status END,
|
||
current_question = CASE WHEN ${has("currentQuestion")} THEN ${patch.currentQuestion ? JSON.stringify(patch.currentQuestion) : null}::text ELSE current_question END,
|
||
conversation_history = CASE WHEN ${has("conversationHistory")} THEN ${patch.conversationHistory ? JSON.stringify(patch.conversationHistory) : null}::text ELSE conversation_history END,
|
||
artifact_path = CASE WHEN ${has("artifactPath")} THEN ${patch.artifactPath ?? null}::text ELSE artifact_path END,
|
||
error = CASE WHEN ${has("error")} THEN ${patch.error ?? null}::text ELSE error END,
|
||
last_activity_at = ${lastActivityAt},
|
||
updated_at = ${updatedAt}
|
||
WHERE project_id=${projectId} AND id=${id}
|
||
RETURNING id, stage, status, current_question AS "currentQuestion", conversation_history AS "conversationHistory", project_id AS "projectId", artifact_path AS "artifactPath", error, turn_interval_ms AS "turnIntervalMs", last_activity_at AS "lastActivityAt", created_at AS "createdAt", updated_at AS "updatedAt"
|
||
`) as unknown as CeSessionRow[];
|
||
return rows[0] ? rowToSession(rows[0]) : undefined;
|
||
}
|
||
|
||
/** Atomically append one history turn so simultaneous progress/terminal writes cannot drop either turn. */
|
||
async appendHistoryAsync(id: string, turn: CeConversationTurn): Promise<CeSession | undefined> {
|
||
const projectId = this.requireProjectId();
|
||
const now = Date.now();
|
||
const updatedAt = new Date(now).toISOString();
|
||
const rows = await this.requireLayer().db.execute(sql`
|
||
UPDATE project.ce_sessions SET
|
||
conversation_history = ((conversation_history::jsonb || jsonb_build_array(${JSON.stringify(turn)}::jsonb))::text),
|
||
last_activity_at = ${now},
|
||
updated_at = ${updatedAt}
|
||
WHERE project_id=${projectId} AND id=${id}
|
||
RETURNING id, stage, status, current_question AS "currentQuestion", conversation_history AS "conversationHistory", project_id AS "projectId", artifact_path AS "artifactPath", error, turn_interval_ms AS "turnIntervalMs", last_activity_at AS "lastActivityAt", created_at AS "createdAt", updated_at AS "updatedAt"
|
||
`) as unknown as CeSessionRow[];
|
||
return rows[0] ? rowToSession(rows[0]) : undefined;
|
||
}
|
||
|
||
/** Atomically touch only liveness columns; safe to run beside terminal/history mutations. */
|
||
async touchActivityAsync(id: string, at = Date.now()): Promise<boolean> {
|
||
if (!this.asyncLayer) return Boolean(this.update(id, { lastActivityAt: at }));
|
||
const rows = await this.requireLayer().db.execute(sql`
|
||
UPDATE project.ce_sessions
|
||
SET last_activity_at=${at}, updated_at=${new Date(at).toISOString()}
|
||
WHERE project_id=${this.requireProjectId()} AND id=${id}
|
||
RETURNING id
|
||
`) as unknown as Array<{ id: string }>;
|
||
return rows.length > 0;
|
||
}
|
||
async deleteAsync(id: string): Promise<boolean> { const existing = await this.getAsync(id); if (!existing) return false; await this.requireLayer().db.execute(sql`DELETE FROM project.ce_sessions WHERE project_id=${this.requireProjectId()} AND id=${id}`); return true; }
|
||
async recoverStaleSessionsAsync(now = Date.now(), multiple = STALE_INTERVAL_MULTIPLE): Promise<string[]> {
|
||
const rows = await this.requireLayer().db.execute(sql`SELECT id, stage, status, current_question AS "currentQuestion", conversation_history AS "conversationHistory", project_id AS "projectId", artifact_path AS "artifactPath", error, turn_interval_ms AS "turnIntervalMs", last_activity_at AS "lastActivityAt", created_at AS "createdAt", updated_at AS "updatedAt" FROM project.ce_sessions WHERE project_id=${this.requireProjectId()} AND status IN ('active', 'launching') AND last_activity_at < (${now}::bigint - (${multiple}::bigint * turn_interval_ms::bigint))`) as unknown as CeSessionRow[];
|
||
const candidates = rows.map(rowToSession);
|
||
await Promise.all(candidates.map((session) => this.updateAsync(
|
||
session.id,
|
||
session.currentQuestion
|
||
? { status: "awaiting_input" }
|
||
: { status: "interrupted", error: session.error ?? "Session interrupted — progress preserved, resume to continue" },
|
||
)));
|
||
return candidates.map((session) => session.id);
|
||
}
|
||
}
|
||
|
||
const storeCache = new WeakMap<object, CeSessionStore>();
|
||
|
||
/** WeakMap-cached store keyed by the TaskStore instance (mirrors reports). */
|
||
export function getCeSessionStore(ctx: PluginContext): CeSessionStore {
|
||
const key = ctx.taskStore as object;
|
||
const cached = storeCache.get(key);
|
||
if (cached) return cached;
|
||
const layer = ctx.taskStore.getAsyncLayer();
|
||
if (!layer) throw new Error("Compound Engineering session store requires the project PostgreSQL AsyncDataLayer");
|
||
const store = new CeSessionStore(null, layer);
|
||
storeCache.set(key, store);
|
||
return store;
|
||
}
|