fix(compound-engineering): claim plan handoffs atomically

This commit is contained in:
gsxdsm
2026-07-11 19:35:07 -07:00
parent 337bdcd081
commit b0ca17a1f7
4 changed files with 81 additions and 20 deletions

View File

@@ -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)", () => {

View File

@@ -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,

View File

@@ -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);

View File

@@ -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). */