From b0ca17a1f78a5dd90bd45f8376311c1d1c09765b Mon Sep 17 00:00:00 2001 From: gsxdsm Date: Sat, 11 Jul 2026 19:35:07 -0700 Subject: [PATCH] fix(compound-engineering): claim plan handoffs atomically --- .../src/__tests__/session-store.test.ts | 18 ++++++- .../src/schema.ts | 10 ++++ .../src/session/orchestrator.ts | 21 +++----- .../src/session/session-store.ts | 52 +++++++++++++++++-- 4 files changed, 81 insertions(+), 20 deletions(-) diff --git a/plugins/fusion-plugin-compound-engineering/src/__tests__/session-store.test.ts b/plugins/fusion-plugin-compound-engineering/src/__tests__/session-store.test.ts index 09f2ff1784..98e6e6c1f4 100644 --- a/plugins/fusion-plugin-compound-engineering/src/__tests__/session-store.test.ts +++ b/plugins/fusion-plugin-compound-engineering/src/__tests__/session-store.test.ts @@ -1,5 +1,5 @@ import { afterEach, beforeEach, describe, expect, it } from "vitest"; -import { CeSessionStore, STALE_INTERVAL_MULTIPLE } from "../session/session-store.js"; +import { CeSessionStore, PlanHandoffClaimError, STALE_INTERVAL_MULTIPLE } from "../session/session-store.js"; import { ensureCeSchema } from "../schema.js"; import { makeHarness, type TestHarness } from "./_harness.js"; @@ -28,6 +28,10 @@ describe("ensureCeSchema", () => { "lastActivityAt", ]), ); + const claimsTable = h.db + .prepare("SELECT name FROM sqlite_master WHERE type = 'table' AND name = 'ce_plan_handoff_claims'") + .get() as { name: string } | undefined; + expect(claimsTable?.name).toBe("ce_plan_handoff_claims"); }); }); @@ -104,6 +108,18 @@ describe("multi-session independence + delete", () => { // Deleting a missing row reports false, no throw. expect(store.delete(b.id)).toBe(false); }); + + it("atomically claims a requirements artifact for one Plan session and releases it on discard", () => { + const store = new CeSessionStore(h.db); + const artifactPath = "docs/plans/2026-07-11-001-topic-plan.md"; + const first = store.createWithPlanHandoffClaim({ stage: "plan", projectId: "p1" }, artifactPath); + + expect(() => store.createWithPlanHandoffClaim({ stage: "plan", projectId: "p1" }, artifactPath)).toThrow(PlanHandoffClaimError); + expect(store.delete(first.id)).toBe(true); + + const retry = store.createWithPlanHandoffClaim({ stage: "plan", projectId: "p1" }, artifactPath); + expect(retry.artifactPath).toBe(artifactPath); + }); }); describe("interval-relative staleness (FN-4172 rubric)", () => { diff --git a/plugins/fusion-plugin-compound-engineering/src/schema.ts b/plugins/fusion-plugin-compound-engineering/src/schema.ts index 5d9d6f61b5..d6a8949fae 100644 --- a/plugins/fusion-plugin-compound-engineering/src/schema.ts +++ b/plugins/fusion-plugin-compound-engineering/src/schema.ts @@ -68,6 +68,16 @@ export function ensureCeSchema(db: Database): void { CREATE INDEX IF NOT EXISTS idxCeSessionsProject ON ce_sessions(projectId, updatedAt DESC, id); + CREATE TABLE IF NOT EXISTS ce_plan_handoff_claims ( + artifactPath TEXT PRIMARY KEY, + sessionId TEXT NOT NULL UNIQUE, + projectId TEXT, + createdAt TEXT NOT NULL + ); + + CREATE INDEX IF NOT EXISTS idxCePlanHandoffClaimsSession + ON ce_plan_handoff_claims(sessionId); + CREATE TABLE IF NOT EXISTS ce_pipeline_links ( id TEXT PRIMARY KEY, taskId TEXT NOT NULL, diff --git a/plugins/fusion-plugin-compound-engineering/src/session/orchestrator.ts b/plugins/fusion-plugin-compound-engineering/src/session/orchestrator.ts index 572448f9b1..3fa0e88554 100644 --- a/plugins/fusion-plugin-compound-engineering/src/session/orchestrator.ts +++ b/plugins/fusion-plugin-compound-engineering/src/session/orchestrator.ts @@ -465,29 +465,20 @@ export class CeOrchestrator { /* * FNXC:CompoundEngineeringPlanning 2026-07-10-22:52: - * Brainstorm creates the requirements-only unified plan. A same-project Plan session must carry the selected completed predecessor's safe docs/plans artifact path, accept it only while it remains requirements-only, and atomically replace it with valid implementation-ready output; absent a compatible handoff, legacy new-file behavior remains available. + * Brainstorm creates the requirements-only unified plan. A same-project Plan session must carry the selected completed predecessor's safe docs/plans artifact path, accept it only while it remains requirements-only, and atomically claim it with row creation so concurrent starts cannot enrich the same file; absent a compatible handoff, legacy new-file behavior remains available. */ const handoffArtifactPath = stageId === PLAN_STAGE_ID ? this.findBrainstormHandoffArtifact(opts.projectId ?? null, opts.sourceSessionId) : null; - if (handoffArtifactPath) { - const competingPlan = this.store.list({ stage: PLAN_STAGE_ID }).find((candidate) => ( - candidate.projectId === (opts.projectId ?? null) - && candidate.artifactPath === handoffArtifactPath - && candidate.status !== "completed" - && candidate.status !== "error" - && candidate.status !== "interrupted" - )); - if (competingPlan) { - throw new Error(`Plan session ${competingPlan.id} is already enriching ${handoffArtifactPath}`); - } - } - const session = this.store.create({ + const sessionInput = { stage: stageId, projectId: opts.projectId ?? null, artifactPath: handoffArtifactPath, turnIntervalMs: this.turnTimeoutMs, - }); + }; + const session = handoffArtifactPath + ? this.store.createWithPlanHandoffClaim(sessionInput, handoffArtifactPath) + : this.store.create(sessionInput); this.store.appendHistory(session.id, { role: "user", text: opts.openingMessage, at: new Date().toISOString() }); const turn = this.runOpeningTurn(session.id, stage, opts.openingMessage); diff --git a/plugins/fusion-plugin-compound-engineering/src/session/session-store.ts b/plugins/fusion-plugin-compound-engineering/src/session/session-store.ts index 5f0f5d8309..ab786f7ce5 100644 --- a/plugins/fusion-plugin-compound-engineering/src/session/session-store.ts +++ b/plugins/fusion-plugin-compound-engineering/src/session/session-store.ts @@ -93,6 +93,16 @@ export interface CreateCeSessionInput { 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"; + } +} + /** * Default multiple of the turn interval beyond which a non-terminal session is * considered stale. Mirrors the FN-4172 rubric (`> 3× interval`), interval- @@ -172,8 +182,34 @@ export class CeSessionStore { } 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 }); + return this.db.transactionImmediate(() => { + const existing = this.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); + this.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(); - const session: CeSession = { + return { id: input.id ?? randomUUID(), stage: input.stage, status: "launching", @@ -187,6 +223,9 @@ export class CeSessionStore { createdAt: now, updatedAt: now, }; + } + + private insert(session: CeSession): void { this.db .prepare( `INSERT INTO ce_sessions @@ -207,7 +246,6 @@ export class CeSessionStore { session.createdAt, session.updatedAt, ); - return session; } get(id: string): CeSession | undefined { @@ -281,8 +319,14 @@ export class CeSessionStore { /** Delete a session row. Returns true when a row was removed. */ delete(id: string): boolean { - const result = this.db.prepare(`DELETE FROM ce_sessions WHERE id = ?`).run(id); - return Number(result.changes ?? 0) > 0; + return this.db.transactionImmediate(() => { + const result = this.db.prepare(`DELETE FROM ce_sessions WHERE id = ?`).run(id); + if (Number(result.changes ?? 0) > 0) { + this.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). */