fix(dashboard,cli): un-dead-end deleted plan tasks; harden fn task plan per review
Reported bug (screenshot): deleting the task created from a plan left the
session permanently stuck on PLANNING_CREATED_TASK_MISSING — Retry create
replayed the same 409 forever. A linked task absent from the
include-archived scan (task-row authority; a successful scan proves
deletion, not a flaky read) now clears the stale linkage and creates a
fresh task, in both the create-task route and createTaskFromPlanSession;
a still-listed-but-unreadable task keeps failing closed.
Multi-agent review of fdd120232 (correctness/adversarial/reliability):
- P1: CLI planning sessions were memory-only — setAiSessionStore only ran
in the dashboard server, so --resume could never find a session across
invocations. New ensureDurablePlanningSessionStore wires the durable
AiSessionStore over the board store's public asyncLayer in runTaskPlan.
- P1: resume failures now THROW instead of process.exit (fn_task_plan
runs inside the pi host — an exit killed the whole agent session), and
a no-question resume requires an explicit refine focus (the provided
description) so merely resuming never rotates the epoch.
- P2: claim and finalize CAS gained the same expected-epoch WHERE guard
as reconcile, so a stale-epoch creator can no longer finalize an
old-epoch task onto a rotated session.
- Side-effect failures (documents, logEntry, validate, reconcile) are now
logged instead of swallowed; post-insert failures no longer mislabel
the just-created task alreadyCreated:true; the keep-refining readline
closes on thrown prompts and a failed refine after creation returns the
created task id with a resume hint; cross-process generating guard
added to createTaskFromPlanSession.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
7
.changeset/planning-deleted-task-recreate.md
Normal file
7
.changeset/planning-deleted-task-recreate.md
Normal file
@@ -0,0 +1,7 @@
|
|||||||
|
---
|
||||||
|
"@runfusion/fusion": patch
|
||||||
|
---
|
||||||
|
|
||||||
|
summary: Deleting a task created from a plan no longer dead-ends the plan — Proceed creates a fresh task.
|
||||||
|
category: fix
|
||||||
|
dev: `PLANNING_CREATED_TASK_MISSING` now only fires when the linked task is still listed but unreadable (transient read); a task absent from the include-archived scan clears the stale linkage in both the create-task route and `createTaskFromPlanSession`. CLI/agent create side-effect failures are now logged; keep-refining readline closes on thrown prompts.
|
||||||
@@ -30,6 +30,7 @@ vi.mock("@fusion/dashboard/planning", () => ({
|
|||||||
submitResponse: vi.fn(),
|
submitResponse: vi.fn(),
|
||||||
validateSession: vi.fn(),
|
validateSession: vi.fn(),
|
||||||
getSession: vi.fn(),
|
getSession: vi.fn(),
|
||||||
|
ensureDurablePlanningSessionStore: vi.fn(async () => true),
|
||||||
/*
|
/*
|
||||||
FNXC:PlanningMultiTask 2026-07-24-02:30:
|
FNXC:PlanningMultiTask 2026-07-24-02:30:
|
||||||
The CLI now creates through the claim-aware shared path (idempotency + session linkage +
|
The CLI now creates through the claim-aware shared path (idempotency + session linkage +
|
||||||
@@ -510,18 +511,55 @@ describe("runTaskPlan", () => {
|
|||||||
data: { title: "Second task", description: "d2", suggestedSize: "S", suggestedDependencies: [], keyDeliverables: ["X"] },
|
data: { title: "Second task", description: "d2", suggestedSize: "S", suggestedDependencies: [], keyDeliverables: ["X"] },
|
||||||
});
|
});
|
||||||
mockQuestion
|
mockQuestion
|
||||||
.mockResolvedValueOnce("tighten the scope")
|
.mockResolvedValueOnce("split the rollout")
|
||||||
.mockResolvedValueOnce("DONE");
|
.mockResolvedValueOnce("DONE");
|
||||||
|
|
||||||
const taskId = await runTaskPlan(undefined, true, undefined, undefined, "resume-session-9");
|
const taskId = await runTaskPlan("tighten the scope", true, undefined, undefined, "resume-session-9");
|
||||||
|
|
||||||
expect(createSession).not.toHaveBeenCalled();
|
expect(createSession).not.toHaveBeenCalled();
|
||||||
// The no-question resume issues a refine turn to regenerate the interview.
|
// The no-question resume issues a refine turn carrying the provided description as focus.
|
||||||
expect(submitResponse).toHaveBeenNthCalledWith(1, "resume-session-9", { refine: true }, "/test/project", undefined, expect.anything());
|
expect(submitResponse).toHaveBeenNthCalledWith(1, "resume-session-9", { refine: true, focus: "tighten the scope" }, "/test/project", undefined, expect.anything());
|
||||||
expect(createTaskFromPlanSession).toHaveBeenCalledWith("resume-session-9", expect.anything(), { baseBranch: undefined });
|
expect(createTaskFromPlanSession).toHaveBeenCalledWith("resume-session-9", expect.anything(), { baseBranch: undefined });
|
||||||
expect(taskId).toBe("FN-042");
|
expect(taskId).toBe("FN-042");
|
||||||
});
|
});
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:PlanningMultiTask 2026-07-24-03:40:
|
||||||
|
Review findings: resume failures must THROW (fn_task_plan runs in the pi host — process.exit
|
||||||
|
killed the whole agent session), and a no-question resume without a refine focus must not
|
||||||
|
silently rotate the epoch.
|
||||||
|
*/
|
||||||
|
it("throws instead of exiting when the resumed session is missing", async () => {
|
||||||
|
setupTaskStoreMock();
|
||||||
|
(getSession as unknown as ReturnType<typeof vi.fn>).mockResolvedValueOnce(undefined);
|
||||||
|
const exitSpy = vi.spyOn(process, "exit").mockImplementation(() => {
|
||||||
|
throw new Error("Process.exit called");
|
||||||
|
});
|
||||||
|
|
||||||
|
await expect(runTaskPlan(undefined, true, undefined, undefined, "missing-session"))
|
||||||
|
.rejects.toThrow(/not found/);
|
||||||
|
expect(exitSpy).not.toHaveBeenCalled();
|
||||||
|
exitSpy.mockRestore();
|
||||||
|
});
|
||||||
|
|
||||||
|
it("throws when resuming a no-question session without a refine focus", async () => {
|
||||||
|
setupTaskStoreMock();
|
||||||
|
(getSession as unknown as ReturnType<typeof vi.fn>).mockResolvedValueOnce({
|
||||||
|
id: "resume-session-10",
|
||||||
|
currentQuestion: null,
|
||||||
|
summary: { title: "Existing plan", description: "d", suggestedSize: "M", suggestedDependencies: [], keyDeliverables: [] },
|
||||||
|
});
|
||||||
|
const exitSpy = vi.spyOn(process, "exit").mockImplementation(() => {
|
||||||
|
throw new Error("Process.exit called");
|
||||||
|
});
|
||||||
|
|
||||||
|
await expect(runTaskPlan(undefined, true, undefined, undefined, "resume-session-10"))
|
||||||
|
.rejects.toThrow(/refine/i);
|
||||||
|
expect(submitResponse).not.toHaveBeenCalled();
|
||||||
|
expect(exitSpy).not.toHaveBeenCalled();
|
||||||
|
exitSpy.mockRestore();
|
||||||
|
});
|
||||||
|
|
||||||
it("handles RateLimitError with proper message", async () => {
|
it("handles RateLimitError with proper message", async () => {
|
||||||
setupTaskStoreMock();
|
setupTaskStoreMock();
|
||||||
|
|
||||||
|
|||||||
@@ -2,7 +2,7 @@ import { TaskStore, COLUMNS, COLUMN_LABELS, CentralCore, buildAutoPauseClearPatc
|
|||||||
import { isInReviewMissingWorktreeSessionStartFailure, runAiMerge, landWorkspaceTask, installBaselineArchiveWorktreeDisposer } from "@fusion/engine";
|
import { isInReviewMissingWorktreeSessionStartFailure, runAiMerge, landWorkspaceTask, installBaselineArchiveWorktreeDisposer } from "@fusion/engine";
|
||||||
import { createInterface } from "node:readline/promises";
|
import { createInterface } from "node:readline/promises";
|
||||||
import type { PlanningQuestion, PlanningSummary } from "@fusion/core";
|
import type { PlanningQuestion, PlanningSummary } from "@fusion/core";
|
||||||
import { createSession, createTaskFromPlanSession, getSession as getPlanningSession, submitResponse, validateSession, RateLimitError, SessionNotFoundError, InvalidSessionStateError } from "@fusion/dashboard/planning";
|
import { createSession, createTaskFromPlanSession, ensureDurablePlanningSessionStore, getSession as getPlanningSession, submitResponse, validateSession, RateLimitError, SessionNotFoundError, InvalidSessionStateError } from "@fusion/dashboard/planning";
|
||||||
import { watchFile, unwatchFile, statSync, existsSync, readFileSync } from "node:fs";
|
import { watchFile, unwatchFile, statSync, existsSync, readFileSync } from "node:fs";
|
||||||
import { basename, join } from "node:path";
|
import { basename, join } from "node:path";
|
||||||
import * as dashboard from "@fusion/dashboard";
|
import * as dashboard from "@fusion/dashboard";
|
||||||
@@ -2207,6 +2207,16 @@ export async function runTaskPlan(
|
|||||||
let sessionId: string;
|
let sessionId: string;
|
||||||
let firstQuestion: PlanningQuestion;
|
let firstQuestion: PlanningQuestion;
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:PlanningMultiTask 2026-07-24-03:40:
|
||||||
|
Review P1: CLI planning sessions were memory-only (setAiSessionStore only ran in the
|
||||||
|
dashboard server), so `--resume` could never find a session across invocations and CLI
|
||||||
|
claim/linkage state was invisible to the dashboard. Wire the durable AiSessionStore over
|
||||||
|
the resolved board store before any session work; legacy stores without an async layer
|
||||||
|
stay in-memory and resume is honestly reported as unavailable.
|
||||||
|
*/
|
||||||
|
const durablePlanningStore = await ensureDurablePlanningSessionStore(store).catch(() => false);
|
||||||
|
|
||||||
try {
|
try {
|
||||||
showThinking();
|
showThinking();
|
||||||
const projectPath = context.projectPath;
|
const projectPath = context.projectPath;
|
||||||
@@ -2218,24 +2228,40 @@ export async function runTaskPlan(
|
|||||||
awaiting input, a refine turn regenerates one — the server reopens the session and
|
awaiting input, a refine turn regenerates one — the server reopens the session and
|
||||||
rotates the creation epoch when its current epoch already produced a task, so a later
|
rotates the creation epoch when its current epoch already produced a task, so a later
|
||||||
/validate creates a NEW task instead of replaying the old one.
|
/validate creates a NEW task instead of replaying the old one.
|
||||||
|
|
||||||
|
FNXC:PlanningMultiTask 2026-07-24-03:40:
|
||||||
|
Review findings: resume failures THROW instead of process.exit (fn_task_plan runs inside
|
||||||
|
the pi host process — an exit killed the whole agent session), and the regenerating
|
||||||
|
refine turn only fires WITH expressed intent (the provided description becomes the
|
||||||
|
refine focus) so merely resuming never rotates the epoch or un-validates the plan.
|
||||||
*/
|
*/
|
||||||
const existing = await getPlanningSession(resumeSessionId);
|
const existing = await getPlanningSession(resumeSessionId);
|
||||||
if (!existing) {
|
if (!existing) {
|
||||||
clearThinking();
|
clearThinking();
|
||||||
console.error(`\n Planning session ${resumeSessionId} not found or expired.\n`);
|
await closeProjectStore(context).catch(() => {});
|
||||||
await closeBoardContextAndExit(context, 1);
|
throw new SessionNotFoundError(
|
||||||
return undefined;
|
durablePlanningStore
|
||||||
|
? `Planning session ${resumeSessionId} not found or expired.`
|
||||||
|
: `Planning session ${resumeSessionId} not found — this store does not persist planning sessions, so resume only works within the process that created the session.`,
|
||||||
|
);
|
||||||
}
|
}
|
||||||
sessionId = resumeSessionId;
|
sessionId = resumeSessionId;
|
||||||
if (existing.currentQuestion) {
|
if (existing.currentQuestion) {
|
||||||
firstQuestion = existing.currentQuestion;
|
firstQuestion = existing.currentQuestion;
|
||||||
} else {
|
} else {
|
||||||
const regenerated = await submitResponse(sessionId, { refine: true }, projectPath, undefined, store);
|
const focus = initialPlan?.trim();
|
||||||
|
if (!focus) {
|
||||||
|
clearThinking();
|
||||||
|
await closeProjectStore(context).catch(() => {});
|
||||||
|
throw new InvalidSessionStateError(
|
||||||
|
"This plan has no question awaiting input. Provide a description of what to refine next (it becomes the refine focus) to continue the interview.",
|
||||||
|
);
|
||||||
|
}
|
||||||
|
const regenerated = await submitResponse(sessionId, { refine: true, focus }, projectPath, undefined, store);
|
||||||
if (regenerated.type !== "question") {
|
if (regenerated.type !== "question") {
|
||||||
clearThinking();
|
clearThinking();
|
||||||
console.error("\n Could not resume the interview for this session.\n");
|
await closeProjectStore(context).catch(() => {});
|
||||||
await closeBoardContextAndExit(context, 1);
|
throw new InvalidSessionStateError("Could not resume the interview for this session.");
|
||||||
return undefined;
|
|
||||||
}
|
}
|
||||||
firstQuestion = regenerated.data;
|
firstQuestion = regenerated.data;
|
||||||
}
|
}
|
||||||
@@ -2249,6 +2275,10 @@ export async function runTaskPlan(
|
|||||||
} catch (err) {
|
} catch (err) {
|
||||||
clearThinking();
|
clearThinking();
|
||||||
|
|
||||||
|
if (err instanceof SessionNotFoundError || err instanceof InvalidSessionStateError) {
|
||||||
|
// Resume-branch failures propagate to the caller (bin prints; fn_task_plan returns a tool error).
|
||||||
|
throw err;
|
||||||
|
}
|
||||||
if (err instanceof RateLimitError) {
|
if (err instanceof RateLimitError) {
|
||||||
console.error("\n Rate limit exceeded. Maximum 1000 planning sessions per hour.\n");
|
console.error("\n Rate limit exceeded. Maximum 1000 planning sessions per hour.\n");
|
||||||
await closeBoardContextAndExit(context, 1);
|
await closeBoardContextAndExit(context, 1);
|
||||||
@@ -2399,23 +2429,41 @@ export async function runTaskPlan(
|
|||||||
after one task and can continue later via `fn task plan --resume <sessionId>`.
|
after one task and can continue later via `fn task plan --resume <sessionId>`.
|
||||||
*/
|
*/
|
||||||
if (!yesFlag) {
|
if (!yesFlag) {
|
||||||
|
// FNXC:PlanningMultiTask 2026-07-24-03:20: close the interface on every path, including thrown prompts (review finding — a leaked readline keeps the process alive).
|
||||||
const rlContinue = createInterface({ input: process.stdin, output: process.stdout });
|
const rlContinue = createInterface({ input: process.stdin, output: process.stdout });
|
||||||
const continueAnswer = await rlContinue.question(" Keep refining this plan to create another task? [y/N]: ");
|
let wantsMore = false;
|
||||||
const wantsMore = ["y", "yes"].includes(continueAnswer.trim().toLowerCase());
|
let focus = "";
|
||||||
if (!wantsMore) {
|
try {
|
||||||
|
const continueAnswer = await rlContinue.question(" Keep refining this plan to create another task? [y/N]: ");
|
||||||
|
wantsMore = ["y", "yes"].includes(continueAnswer.trim().toLowerCase());
|
||||||
|
if (wantsMore) {
|
||||||
|
focus = (await rlContinue.question(" What should the next refinement focus on? ")).trim();
|
||||||
|
}
|
||||||
|
} finally {
|
||||||
rlContinue.close();
|
rlContinue.close();
|
||||||
|
}
|
||||||
|
if (!wantsMore) {
|
||||||
return task.id;
|
return task.id;
|
||||||
}
|
}
|
||||||
const focus = (await rlContinue.question(" What should the next refinement focus on? ")).trim();
|
/*
|
||||||
rlContinue.close();
|
FNXC:PlanningMultiTask 2026-07-24-03:40:
|
||||||
showThinking();
|
Review finding: a provider error on this refine turn used to escape the loop's
|
||||||
const refined = await submitResponse(sessionId, { refine: true, ...(focus ? { focus } : {}) }, context.projectPath, undefined, store);
|
error handling AFTER the task was created, discarding the task id. The task
|
||||||
clearThinking();
|
exists — report the refine failure and return it.
|
||||||
if (refined.type === "question") {
|
*/
|
||||||
currentQuestion = refined.data;
|
try {
|
||||||
continue;
|
showThinking();
|
||||||
|
const refined = await submitResponse(sessionId, { refine: true, ...(focus ? { focus } : {}) }, context.projectPath, undefined, store);
|
||||||
|
clearThinking();
|
||||||
|
if (refined.type === "question") {
|
||||||
|
currentQuestion = refined.data;
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
console.log("\n Could not continue the interview; the created task is ready.\n");
|
||||||
|
} catch (refineErr) {
|
||||||
|
clearThinking();
|
||||||
|
console.error(`\n Refine failed (${refineErr instanceof Error ? refineErr.message : String(refineErr)}); the created task is ready. Resume later with: fn task plan --resume ${sessionId}\n`);
|
||||||
}
|
}
|
||||||
console.log("\n Could not continue the interview; the created task is ready.\n");
|
|
||||||
}
|
}
|
||||||
|
|
||||||
return task.id;
|
return task.id;
|
||||||
|
|||||||
@@ -117,4 +117,30 @@ pgTest("planning session claim lifecycle (multi-task epochs)", () => {
|
|||||||
expect(afterReconcile.taskCreationEpoch).toBe(2);
|
expect(afterReconcile.taskCreationEpoch).toBe(2);
|
||||||
expect(afterReconcile.createdTaskIds).toEqual(["FN-1", "FN-2"]);
|
expect(afterReconcile.createdTaskIds).toEqual(["FN-1", "FN-2"]);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:PlanningMultiTask 2026-07-24-03:40:
|
||||||
|
Review finding: claim and finalize need the same epoch guard as reconcile, or a stale-epoch
|
||||||
|
creator could finalize an old-epoch task onto a session a concurrent edit already rotated.
|
||||||
|
*/
|
||||||
|
it("claim and finalize with a stale expected epoch are no-ops", async () => {
|
||||||
|
const db = h.layer().db;
|
||||||
|
const sessionId = "planning-claim-epoch-guarded-writes";
|
||||||
|
await upsertAiSession(db, planningRow(sessionId, { taskCreationEpoch: 2, createdTaskIds: ["FN-1"] }));
|
||||||
|
|
||||||
|
// Stale-epoch claim loses the CAS entirely.
|
||||||
|
expect(await claimPlanningSessionTaskCreation(db, sessionId, "stale-token", new Date().toISOString(), 1)).toBeNull();
|
||||||
|
expect(payloadOf(await getAiSession(db, sessionId)).createClaimStatus).toBeUndefined();
|
||||||
|
|
||||||
|
// Matching-epoch claim succeeds; a finalize whose expected epoch went stale mid-flight is a no-op.
|
||||||
|
const claimed = await claimPlanningSessionTaskCreation(db, sessionId, "live-token", new Date().toISOString(), 2);
|
||||||
|
expect(claimed).not.toBeNull();
|
||||||
|
expect(await finalizePlanningSessionTaskCreation(db, sessionId, "live-token", "FN-STALE", 1)).toBeNull();
|
||||||
|
const afterStaleFinalize = payloadOf(await getAiSession(db, sessionId));
|
||||||
|
expect(afterStaleFinalize.createClaimStatus).toBe("creating");
|
||||||
|
expect(afterStaleFinalize.createdTaskId).toBeUndefined();
|
||||||
|
|
||||||
|
const finalized = await finalizePlanningSessionTaskCreation(db, sessionId, "live-token", "FN-4", 2);
|
||||||
|
expect(payloadOf(finalized).createdTaskId).toBe("FN-4");
|
||||||
|
});
|
||||||
});
|
});
|
||||||
|
|||||||
@@ -310,11 +310,25 @@ survive. The WHERE guards still evaluate against the CURRENT row at update time.
|
|||||||
const CLAIM_KEYS_PATCH = (patch: Record<string, string>, removeKeys: string[]) =>
|
const CLAIM_KEYS_PATCH = (patch: Record<string, string>, removeKeys: string[]) =>
|
||||||
sql`(${schema.project.aiSessions.inputPayload} || ${JSON.stringify(patch)}::jsonb)${sql.raw(removeKeys.map((key) => ` - '${key.replace(/'/g, "''")}'`).join(""))}`;
|
sql`(${schema.project.aiSessions.inputPayload} || ${JSON.stringify(patch)}::jsonb)${sql.raw(removeKeys.map((key) => ` - '${key.replace(/'/g, "''")}'`).join(""))}`;
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:PlanningMultiTask 2026-07-24-03:40:
|
||||||
|
Review finding: without epoch guards a stale-epoch creator could claim/finalize an old-epoch
|
||||||
|
task onto a session that a concurrent edit turn had already rotated, re-linking an archived
|
||||||
|
plan's task to the NEW epoch. When the caller supplies expectedTaskCreationEpoch, the
|
||||||
|
conditional update only fires while the row's epoch still matches — a lost race is a no-op
|
||||||
|
and the caller recovers through the task-row proposalClaimId authority.
|
||||||
|
*/
|
||||||
|
const epochGuard = (expectedTaskCreationEpoch: number | undefined) =>
|
||||||
|
expectedTaskCreationEpoch === undefined
|
||||||
|
? []
|
||||||
|
: [sql`coalesce((${schema.project.aiSessions.inputPayload}->>'taskCreationEpoch')::int, 0) = ${expectedTaskCreationEpoch}`];
|
||||||
|
|
||||||
export async function claimPlanningSessionTaskCreation(
|
export async function claimPlanningSessionTaskCreation(
|
||||||
handle: QueryHandle,
|
handle: QueryHandle,
|
||||||
sessionId: string,
|
sessionId: string,
|
||||||
claimOwnerToken: string,
|
claimOwnerToken: string,
|
||||||
claimStartedAt: string,
|
claimStartedAt: string,
|
||||||
|
expectedTaskCreationEpoch?: number,
|
||||||
): Promise<AiSessionRow | null> {
|
): Promise<AiSessionRow | null> {
|
||||||
const rows = await handle.update(schema.project.aiSessions)
|
const rows = await handle.update(schema.project.aiSessions)
|
||||||
.set({
|
.set({
|
||||||
@@ -325,6 +339,7 @@ export async function claimPlanningSessionTaskCreation(
|
|||||||
eq(schema.project.aiSessions.id, sessionId),
|
eq(schema.project.aiSessions.id, sessionId),
|
||||||
eq(schema.project.aiSessions.type, "planning"),
|
eq(schema.project.aiSessions.type, "planning"),
|
||||||
sql`coalesce(${schema.project.aiSessions.inputPayload}->>'createClaimStatus', 'none') = 'none'`,
|
sql`coalesce(${schema.project.aiSessions.inputPayload}->>'createClaimStatus', 'none') = 'none'`,
|
||||||
|
...epochGuard(expectedTaskCreationEpoch),
|
||||||
))
|
))
|
||||||
.returning();
|
.returning();
|
||||||
return rows[0] ? rowToSession(rows[0]) : null;
|
return rows[0] ? rowToSession(rows[0]) : null;
|
||||||
@@ -336,13 +351,18 @@ export async function finalizePlanningSessionTaskCreation(
|
|||||||
sessionId: string,
|
sessionId: string,
|
||||||
claimOwnerToken: string,
|
claimOwnerToken: string,
|
||||||
createdTaskId: string,
|
createdTaskId: string,
|
||||||
|
expectedTaskCreationEpoch?: number,
|
||||||
): Promise<AiSessionRow | null> {
|
): Promise<AiSessionRow | null> {
|
||||||
const rows = await handle.update(schema.project.aiSessions)
|
const rows = await handle.update(schema.project.aiSessions)
|
||||||
.set({
|
.set({
|
||||||
inputPayload: CLAIM_KEYS_PATCH({ createClaimStatus: "created", createdTaskId }, ["claimOwnerToken", "claimStartedAt"]),
|
inputPayload: CLAIM_KEYS_PATCH({ createClaimStatus: "created", createdTaskId }, ["claimOwnerToken", "claimStartedAt"]),
|
||||||
updatedAt: new Date().toISOString(),
|
updatedAt: new Date().toISOString(),
|
||||||
})
|
})
|
||||||
.where(and(eq(schema.project.aiSessions.id, sessionId), sql`${schema.project.aiSessions.inputPayload}->>'claimOwnerToken' = ${claimOwnerToken}`))
|
.where(and(
|
||||||
|
eq(schema.project.aiSessions.id, sessionId),
|
||||||
|
sql`${schema.project.aiSessions.inputPayload}->>'claimOwnerToken' = ${claimOwnerToken}`,
|
||||||
|
...epochGuard(expectedTaskCreationEpoch),
|
||||||
|
))
|
||||||
.returning();
|
.returning();
|
||||||
return rows[0] ? rowToSession(rows[0]) : null;
|
return rows[0] ? rowToSession(rows[0]) : null;
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -238,9 +238,11 @@ describe("planning question regeneration instead of no-active-question errors",
|
|||||||
const { sessionId } = await startSessionAwaitingInput("10.2.0.13");
|
const { sessionId } = await startSessionAwaitingInput("10.2.0.13");
|
||||||
|
|
||||||
const tasks: Array<{ id: string; title: string; description: string; column: string; dependencies: string[]; proposalClaimId?: string }> = [];
|
const tasks: Array<{ id: string; title: string; description: string; column: string; dependencies: string[]; proposalClaimId?: string }> = [];
|
||||||
|
let taskSequence = 0;
|
||||||
const createTask = vi.fn(async (input: { title: string; description: string; dependencies?: string[]; proposalClaimId?: string }) => {
|
const createTask = vi.fn(async (input: { title: string; description: string; dependencies?: string[]; proposalClaimId?: string }) => {
|
||||||
|
taskSequence += 1;
|
||||||
const task = {
|
const task = {
|
||||||
id: `FN-CLI-${tasks.length + 1}`,
|
id: `FN-CLI-${taskSequence}`,
|
||||||
title: input.title,
|
title: input.title,
|
||||||
description: input.description,
|
description: input.description,
|
||||||
column: "triage",
|
column: "triage",
|
||||||
@@ -280,6 +282,21 @@ describe("planning question regeneration instead of no-active-question errors",
|
|||||||
expect(second.task.id).not.toBe(first.task.id);
|
expect(second.task.id).not.toBe(first.task.id);
|
||||||
expect(createTask).toHaveBeenCalledTimes(2);
|
expect(createTask).toHaveBeenCalledTimes(2);
|
||||||
expect(createTask.mock.calls[1][0].proposalClaimId).toBe(`planning-session:${sessionId}#1`);
|
expect(createTask.mock.calls[1][0].proposalClaimId).toBe(`planning-session:${sessionId}#1`);
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:PlanningMultiTask 2026-07-24-03:20:
|
||||||
|
Reported bug: deleting the created task dead-ended the plan forever
|
||||||
|
(PLANNING_CREATED_TASK_MISSING on every retry). A deleted linked task (absent from the
|
||||||
|
include-archived scan) must clear the stale linkage and create a fresh task.
|
||||||
|
*/
|
||||||
|
const deletedIndex = tasks.findIndex((task) => task.id === second.task.id);
|
||||||
|
tasks.splice(deletedIndex, 1);
|
||||||
|
const afterDeletion = await createTaskFromPlanSession(sessionId, taskStore);
|
||||||
|
expect(afterDeletion.alreadyCreated).toBe(false);
|
||||||
|
expect(afterDeletion.task.id).not.toBe(second.task.id);
|
||||||
|
expect(createTask).toHaveBeenCalledTimes(3);
|
||||||
|
// The replacement task reuses the current epoch's claim key (its unique-index row died with the deleted task).
|
||||||
|
expect(createTask.mock.calls[2][0].proposalClaimId).toBe(`planning-session:${sessionId}#1`);
|
||||||
});
|
});
|
||||||
|
|
||||||
/*
|
/*
|
||||||
|
|||||||
@@ -2967,6 +2967,100 @@ describe("Planning Mode Routes", () => {
|
|||||||
expect(store.createTask).not.toHaveBeenCalled();
|
expect(store.createTask).not.toHaveBeenCalled();
|
||||||
});
|
});
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:PlanningMultiTask 2026-07-24-03:20:
|
||||||
|
Reported bug: deleting the task created from a plan left the session permanently stuck on
|
||||||
|
PLANNING_CREATED_TASK_MISSING — Retry create replayed the same 409 forever. A linked task
|
||||||
|
that is absent from the include-archived scan clears the stale linkage and creates a
|
||||||
|
fresh task; a linked task still LISTED but unreadable keeps failing closed (never fork on
|
||||||
|
a flaky read).
|
||||||
|
*/
|
||||||
|
const buildLinkedGoneRow = (sessionId: string) => ({
|
||||||
|
id: sessionId,
|
||||||
|
type: "planning",
|
||||||
|
status: "complete",
|
||||||
|
title: "Linked task deleted",
|
||||||
|
inputPayload: JSON.stringify({ initialPlan: "Build a thing", validated: true, createdTaskId: "FN-GONE", createClaimStatus: "created" }),
|
||||||
|
conversationHistory: "[]",
|
||||||
|
currentQuestion: null,
|
||||||
|
result: JSON.stringify({
|
||||||
|
title: "Replacement plan",
|
||||||
|
description: "Recreate after the linked task was deleted",
|
||||||
|
suggestedSize: "M",
|
||||||
|
suggestedDependencies: [],
|
||||||
|
keyDeliverables: ["Implementation"],
|
||||||
|
}),
|
||||||
|
thinkingOutput: "",
|
||||||
|
error: null,
|
||||||
|
projectId: null,
|
||||||
|
createdAt: "2026-01-01T00:00:00.000Z",
|
||||||
|
updatedAt: "2026-01-01T00:00:00.000Z",
|
||||||
|
});
|
||||||
|
|
||||||
|
it("creates a fresh task when the linked task was deleted instead of dead-ending", async () => {
|
||||||
|
const sessionId = "planning-linked-task-deleted";
|
||||||
|
const mockStore = new MockAiSessionStore();
|
||||||
|
await mockStore.upsert(buildLinkedGoneRow(sessionId) as never);
|
||||||
|
setAiSessionStore(mockStore as unknown as Parameters<typeof setAiSessionStore>[0]);
|
||||||
|
|
||||||
|
(store.listTasks as ReturnType<typeof vi.fn>).mockResolvedValue([]);
|
||||||
|
(store.getTask as ReturnType<typeof vi.fn>).mockRejectedValue(new Error("Task FN-GONE not found"));
|
||||||
|
(store.createTask as ReturnType<typeof vi.fn>).mockResolvedValue({
|
||||||
|
id: "FN-REBORN",
|
||||||
|
description: "Recreated task",
|
||||||
|
column: "triage",
|
||||||
|
dependencies: [],
|
||||||
|
createdAt: "2026-01-01T00:00:00.000Z",
|
||||||
|
updatedAt: "2026-01-01T00:00:00.000Z",
|
||||||
|
});
|
||||||
|
(store.updateTask as ReturnType<typeof vi.fn>).mockResolvedValue({});
|
||||||
|
(store.logEntry as ReturnType<typeof vi.fn>).mockResolvedValue(undefined);
|
||||||
|
|
||||||
|
const appWithAiSessionStore = express();
|
||||||
|
appWithAiSessionStore.use(express.json());
|
||||||
|
appWithAiSessionStore.use("/api", createApiRoutes(store, { aiSessionStore: mockStore as any }));
|
||||||
|
|
||||||
|
const res = await REQUEST(
|
||||||
|
appWithAiSessionStore,
|
||||||
|
"POST",
|
||||||
|
"/api/planning/create-task",
|
||||||
|
JSON.stringify({ sessionId }),
|
||||||
|
{ "Content-Type": "application/json" },
|
||||||
|
);
|
||||||
|
|
||||||
|
expect(res.status).toBe(201);
|
||||||
|
expect(res.body.alreadyCreated).toBe(false);
|
||||||
|
expect(res.body.task.id).toBe("FN-REBORN");
|
||||||
|
expect(store.createTask).toHaveBeenCalledTimes(1);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("keeps failing closed when the linked task is still listed but unreadable", async () => {
|
||||||
|
const sessionId = "planning-linked-task-unreadable";
|
||||||
|
const mockStore = new MockAiSessionStore();
|
||||||
|
await mockStore.upsert(buildLinkedGoneRow(sessionId) as never);
|
||||||
|
setAiSessionStore(mockStore as unknown as Parameters<typeof setAiSessionStore>[0]);
|
||||||
|
|
||||||
|
(store.listTasks as ReturnType<typeof vi.fn>).mockResolvedValue([
|
||||||
|
{ id: "FN-GONE", proposalClaimId: undefined },
|
||||||
|
]);
|
||||||
|
(store.getTask as ReturnType<typeof vi.fn>).mockRejectedValue(new Error("transient store failure"));
|
||||||
|
|
||||||
|
const appWithAiSessionStore = express();
|
||||||
|
appWithAiSessionStore.use(express.json());
|
||||||
|
appWithAiSessionStore.use("/api", createApiRoutes(store, { aiSessionStore: mockStore as any }));
|
||||||
|
|
||||||
|
const res = await REQUEST(
|
||||||
|
appWithAiSessionStore,
|
||||||
|
"POST",
|
||||||
|
"/api/planning/create-task",
|
||||||
|
JSON.stringify({ sessionId }),
|
||||||
|
{ "Content-Type": "application/json" },
|
||||||
|
);
|
||||||
|
|
||||||
|
expect(res.status).toBe(409);
|
||||||
|
expect(store.createTask).not.toHaveBeenCalled();
|
||||||
|
});
|
||||||
|
|
||||||
it.each([
|
it.each([
|
||||||
{
|
{
|
||||||
sessionSource: "live",
|
sessionSource: "live",
|
||||||
|
|||||||
@@ -156,14 +156,15 @@ export class AiSessionStore extends EventEmitter<AiSessionStoreEvents> {
|
|||||||
Planning creation claims must use one conditional database update, not a read then an
|
Planning creation claims must use one conditional database update, not a read then an
|
||||||
upsert, so competing dashboard processes cannot both become the creator.
|
upsert, so competing dashboard processes cannot both become the creator.
|
||||||
*/
|
*/
|
||||||
async claimPlanningTaskCreation(sessionId: string, ownerToken: string, startedAt: string): Promise<AiSessionRow | null> {
|
// FNXC:PlanningMultiTask 2026-07-24-03:40: expectedTaskCreationEpoch guards claim/finalize against a concurrent epoch rotation — see the core CAS functions.
|
||||||
const row = await claimPlanningSessionTaskCreation(this.dbAsync, sessionId, ownerToken, startedAt) as AiSessionRow | null;
|
async claimPlanningTaskCreation(sessionId: string, ownerToken: string, startedAt: string, expectedTaskCreationEpoch?: number): Promise<AiSessionRow | null> {
|
||||||
|
const row = await claimPlanningSessionTaskCreation(this.dbAsync, sessionId, ownerToken, startedAt, expectedTaskCreationEpoch) as AiSessionRow | null;
|
||||||
if (row) this.emit("ai_session:updated", toSummary(row, row.updatedAt));
|
if (row) this.emit("ai_session:updated", toSummary(row, row.updatedAt));
|
||||||
return row;
|
return row;
|
||||||
}
|
}
|
||||||
|
|
||||||
async finalizePlanningTaskCreation(sessionId: string, ownerToken: string, taskId: string): Promise<AiSessionRow | null> {
|
async finalizePlanningTaskCreation(sessionId: string, ownerToken: string, taskId: string, expectedTaskCreationEpoch?: number): Promise<AiSessionRow | null> {
|
||||||
const row = await finalizePlanningSessionTaskCreation(this.dbAsync, sessionId, ownerToken, taskId) as AiSessionRow | null;
|
const row = await finalizePlanningSessionTaskCreation(this.dbAsync, sessionId, ownerToken, taskId, expectedTaskCreationEpoch) as AiSessionRow | null;
|
||||||
if (row) this.emit("ai_session:updated", toSummary(row, row.updatedAt));
|
if (row) this.emit("ai_session:updated", toSummary(row, row.updatedAt));
|
||||||
return row;
|
return row;
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -692,6 +692,25 @@ export function setAiSessionStore(store: AiSessionStore): void {
|
|||||||
_aiSessionStore.on("ai_session:deleted", _aiSessionDeletedListener);
|
_aiSessionStore.on("ai_session:deleted", _aiSessionDeletedListener);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:PlanningMultiTask 2026-07-24-03:40:
|
||||||
|
Review P1: `fn task plan --resume` / fn_task_plan resumeSessionId advertised cross-invocation
|
||||||
|
resume, but setAiSessionStore only ever ran inside the dashboard server — CLI planning
|
||||||
|
sessions were never persisted, getSession always missed in a fresh process, and dashboard
|
||||||
|
sessions were unreachable from the CLI. Non-server callers wire the same durable
|
||||||
|
AiSessionStore over their resolved board store's public asyncLayer here. Returns false when
|
||||||
|
no async layer exists (legacy SQLite mode): planning then stays in-memory and resume is
|
||||||
|
single-process only — callers must say so instead of failing silently.
|
||||||
|
*/
|
||||||
|
export async function ensureDurablePlanningSessionStore(store: TaskStore): Promise<boolean> {
|
||||||
|
if (_aiSessionStore) return true;
|
||||||
|
const layer = (store as { asyncLayer?: unknown }).asyncLayer;
|
||||||
|
if (!layer) return false;
|
||||||
|
const { AiSessionStore: DurableAiSessionStore } = await import("./ai-session-store.js");
|
||||||
|
setAiSessionStore(new DurableAiSessionStore(layer as ConstructorParameters<typeof DurableAiSessionStore>[0]));
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
|
||||||
function cleanupInMemorySession(sessionId: string): boolean {
|
function cleanupInMemorySession(sessionId: string): boolean {
|
||||||
const session = sessions.get(sessionId);
|
const session = sessions.get(sessionId);
|
||||||
if (!session) {
|
if (!session) {
|
||||||
@@ -3973,6 +3992,13 @@ export async function createTaskFromPlanSession(
|
|||||||
if (isPlanningTurnActive(sessionId) || planningStreamManager.hasPendingInitialTurn(sessionId)) {
|
if (isPlanningTurnActive(sessionId) || planningStreamManager.hasPendingInitialTurn(sessionId)) {
|
||||||
throw new GenerationInProgressError("Plan is still generating — wait for the current turn to finish, then create the task.");
|
throw new GenerationInProgressError("Plan is still generating — wait for the current turn to finish, then create the task.");
|
||||||
}
|
}
|
||||||
|
// FNXC:PlanningMultiTask 2026-07-24-03:40: cross-process guard (review finding) — mirror the route's durable status check so a turn generating in ANOTHER process cannot have its linkage torn by this creation.
|
||||||
|
if (_aiSessionStore) {
|
||||||
|
const liveRow = await _aiSessionStore.get(sessionId).catch(() => null);
|
||||||
|
if (liveRow?.type === "planning" && liveRow.status === "generating") {
|
||||||
|
throw new GenerationInProgressError("Plan is still generating — wait for the current turn to finish, then create the task.");
|
||||||
|
}
|
||||||
|
}
|
||||||
const summary = session.summary ?? buildRunningSummary(session.initialPlan, session.history);
|
const summary = session.summary ?? buildRunningSummary(session.initialPlan, session.history);
|
||||||
if (!summary) throw new InvalidSessionStateError("Planning session has no plan to create a task from");
|
if (!summary) throw new InvalidSessionStateError("Planning session has no plan to create a task from");
|
||||||
|
|
||||||
@@ -3982,13 +4008,44 @@ export async function createTaskFromPlanSession(
|
|||||||
(await store.listTasks({ includeArchived: true })).find((candidate) => candidate.proposalClaimId === proposalClaimId);
|
(await store.listTasks({ includeArchived: true })).find((candidate) => candidate.proposalClaimId === proposalClaimId);
|
||||||
const markSessionComplete = async (): Promise<void> => {
|
const markSessionComplete = async (): Promise<void> => {
|
||||||
const current = await getSession(sessionId);
|
const current = await getSession(sessionId);
|
||||||
if (current && !current.validated) await validateSession(sessionId).catch(() => undefined);
|
if (current && !current.validated) {
|
||||||
|
await validateSession(sessionId).catch((err) => {
|
||||||
|
diagnostics.warn("Planning create-task session completion failed", { sessionId, message: err instanceof Error ? err.message : String(err), operation: "create-task-session" });
|
||||||
|
});
|
||||||
|
}
|
||||||
};
|
};
|
||||||
const returnExisting = async (task: Task): Promise<{ task: Task; alreadyCreated: true }> => {
|
const returnExisting = async (task: Task): Promise<{ task: Task; alreadyCreated: true }> => {
|
||||||
await reconcilePlanningTaskCreation(sessionId, task.id, claimEpoch).catch(() => undefined);
|
await reconcilePlanningTaskCreation(sessionId, task.id, claimEpoch).catch((err) => {
|
||||||
|
diagnostics.warn("Planning create-task linkage reconcile failed", { sessionId, taskId: task.id, message: err instanceof Error ? err.message : String(err), operation: "create-task-session" });
|
||||||
|
});
|
||||||
await markSessionComplete();
|
await markSessionComplete();
|
||||||
return { task, alreadyCreated: true };
|
return { task, alreadyCreated: true };
|
||||||
};
|
};
|
||||||
|
/*
|
||||||
|
FNXC:PlanningMultiTask 2026-07-24-03:20:
|
||||||
|
Reported bug (dashboard surface, same contract here): deleting the task created from a plan
|
||||||
|
dead-ended the session forever. When the linked task is absent from the include-archived
|
||||||
|
task list (task-row authority — a successful scan proves deletion, not a flaky read), clear
|
||||||
|
the stale linkage so this attempt creates a fresh task; a transient read failure keeps
|
||||||
|
failing closed so we never fork on a hiccup.
|
||||||
|
*/
|
||||||
|
const clearStaleLinkedTask = async (staleTaskId: string): Promise<boolean> => {
|
||||||
|
const allTasks = await store.listTasks({ includeArchived: true }).catch(() => null);
|
||||||
|
if (allTasks === null || allTasks.some((candidate) => candidate.id === staleTaskId)) return false;
|
||||||
|
diagnostics.warn("Planning session linked task no longer exists; clearing stale linkage", {
|
||||||
|
sessionId,
|
||||||
|
staleTaskId,
|
||||||
|
operation: "create-task-session",
|
||||||
|
});
|
||||||
|
await updatePlanningCreateClaim(sessionId, { createClaimStatus: "none", createdTaskId: undefined, claimOwnerToken: undefined, claimStartedAt: undefined }).catch(() => undefined);
|
||||||
|
if (session) {
|
||||||
|
session.createdTaskId = undefined;
|
||||||
|
session.createClaimStatus = "none";
|
||||||
|
session.claimOwnerToken = undefined;
|
||||||
|
session.claimStartedAt = undefined;
|
||||||
|
}
|
||||||
|
return true;
|
||||||
|
};
|
||||||
|
|
||||||
// The task row under this epoch's key is the crash-window authority.
|
// The task row under this epoch's key is the crash-window authority.
|
||||||
const existingTask = await findCreatedTask();
|
const existingTask = await findCreatedTask();
|
||||||
@@ -3999,10 +4056,13 @@ export async function createTaskFromPlanSession(
|
|||||||
await markSessionComplete();
|
await markSessionComplete();
|
||||||
return { task: linked, alreadyCreated: true };
|
return { task: linked, alreadyCreated: true };
|
||||||
}
|
}
|
||||||
|
if (!(await clearStaleLinkedTask(session.createdTaskId))) {
|
||||||
|
throw new InvalidSessionStateError("PLANNING_CREATED_TASK_MISSING");
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
const claimOwnerToken = randomUUID();
|
const claimOwnerToken = randomUUID();
|
||||||
let claimed = await claimPlanningTaskCreation(sessionId, claimOwnerToken, new Date().toISOString());
|
let claimed = await claimPlanningTaskCreation(sessionId, claimOwnerToken, new Date().toISOString(), claimEpoch);
|
||||||
if (!claimed) {
|
if (!claimed) {
|
||||||
session = (await getDurablePlanningSession(sessionId).catch(() => undefined)) ?? session;
|
session = (await getDurablePlanningSession(sessionId).catch(() => undefined)) ?? session;
|
||||||
const recovered = await findCreatedTask();
|
const recovered = await findCreatedTask();
|
||||||
@@ -4013,17 +4073,27 @@ export async function createTaskFromPlanSession(
|
|||||||
await markSessionComplete();
|
await markSessionComplete();
|
||||||
return { task: linked, alreadyCreated: true };
|
return { task: linked, alreadyCreated: true };
|
||||||
}
|
}
|
||||||
|
if (await clearStaleLinkedTask(session.createdTaskId)) {
|
||||||
|
claimed = await claimPlanningTaskCreation(sessionId, claimOwnerToken, new Date().toISOString(), claimEpoch);
|
||||||
|
}
|
||||||
|
if (!claimed) {
|
||||||
|
throw new GenerationInProgressError("Planning task creation is already in progress");
|
||||||
|
}
|
||||||
}
|
}
|
||||||
const startedAt = session.claimStartedAt ? Date.parse(session.claimStartedAt) : Number.NaN;
|
if (!claimed) {
|
||||||
const leaseExpired = session.createClaimStatus === "creating" && Number.isFinite(startedAt) && Date.now() - startedAt >= 30_000;
|
const startedAt = session.claimStartedAt ? Date.parse(session.claimStartedAt) : Number.NaN;
|
||||||
if (!leaseExpired || !session.claimOwnerToken) {
|
const leaseExpired = session.createClaimStatus === "creating" && Number.isFinite(startedAt) && Date.now() - startedAt >= 30_000;
|
||||||
throw new GenerationInProgressError("Planning task creation is already in progress");
|
if (!leaseExpired || !session.claimOwnerToken) {
|
||||||
|
throw new GenerationInProgressError("Planning task creation is already in progress");
|
||||||
|
}
|
||||||
|
await releasePlanningTaskCreation(sessionId, session.claimOwnerToken);
|
||||||
|
claimed = await claimPlanningTaskCreation(sessionId, claimOwnerToken, new Date().toISOString(), claimEpoch);
|
||||||
|
if (!claimed) throw new GenerationInProgressError("Planning task creation is already in progress");
|
||||||
}
|
}
|
||||||
await releasePlanningTaskCreation(sessionId, session.claimOwnerToken);
|
|
||||||
claimed = await claimPlanningTaskCreation(sessionId, claimOwnerToken, new Date().toISOString());
|
|
||||||
if (!claimed) throw new GenerationInProgressError("Planning task creation is already in progress");
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// FNXC:PlanningMultiTask 2026-07-24-03:40: review finding — a post-insert failure (e.g. finalize) lands in the raced-insert catch; without this marker the task WE created was mislabeled alreadyCreated:true.
|
||||||
|
let insertedTask: Task | undefined;
|
||||||
try {
|
try {
|
||||||
const planMd = formatPlanningPlanMd(summary);
|
const planMd = formatPlanningPlanMd(summary);
|
||||||
const originalRequest = session.initialPlan?.trim() || summary.description.trim();
|
const originalRequest = session.initialPlan?.trim() || summary.description.trim();
|
||||||
@@ -4036,22 +4106,35 @@ export async function createTaskFromPlanSession(
|
|||||||
...(options?.baseBranch?.trim() ? { baseBranch: options.baseBranch.trim() } : {}),
|
...(options?.baseBranch?.trim() ? { baseBranch: options.baseBranch.trim() } : {}),
|
||||||
proposalClaimId,
|
proposalClaimId,
|
||||||
});
|
});
|
||||||
|
insertedTask = task;
|
||||||
|
// FNXC:PlanningMultiTask 2026-07-24-03:20: best-effort side effects must be LOUD on failure (review finding) — a task missing its plan document with no signal is undebuggable.
|
||||||
|
const sideEffect = async (label: string, work: () => Promise<unknown> | unknown): Promise<void> => {
|
||||||
|
try {
|
||||||
|
await work();
|
||||||
|
} catch (err) {
|
||||||
|
diagnostics.warn(label, { sessionId, taskId: task.id, message: err instanceof Error ? err.message : String(err), operation: "create-task-session" });
|
||||||
|
}
|
||||||
|
};
|
||||||
if (summary.suggestedSize) {
|
if (summary.suggestedSize) {
|
||||||
await Promise.resolve(store.updateTask?.(task.id, { size: summary.suggestedSize })).catch(() => undefined);
|
await sideEffect("Planning create-task size update failed", () => store.updateTask?.(task.id, { size: summary.suggestedSize }));
|
||||||
}
|
}
|
||||||
await Promise.resolve(store.upsertTaskDocument?.(task.id, { key: "plan", content: planMd, author: "planning", metadata: { planningSessionId: sessionId, source: "planning-mode" } })).catch(() => undefined);
|
await sideEffect("Planning create-task plan document write failed", () => store.upsertTaskDocument?.(task.id, { key: "plan", content: planMd, author: "planning", metadata: { planningSessionId: sessionId, source: "planning-mode" } }));
|
||||||
if (originalRequest) {
|
if (originalRequest) {
|
||||||
await Promise.resolve(store.upsertTaskDocument?.(task.id, { key: "original-description", content: originalRequest, author: "planning", metadata: { planningSessionId: sessionId, source: "planning-mode-initial-plan" } })).catch(() => undefined);
|
await sideEffect("Planning create-task original description document write failed", () => store.upsertTaskDocument?.(task.id, { key: "original-description", content: originalRequest, author: "planning", metadata: { planningSessionId: sessionId, source: "planning-mode-initial-plan" } }));
|
||||||
}
|
}
|
||||||
await Promise.resolve(store.logEntry?.(task.id, "Created via Planning Mode", `Initial plan: ${(session.initialPlan ?? "").slice(0, 200)}`)).catch(() => undefined);
|
await sideEffect("Planning create-task log entry failed", () => store.logEntry?.(task.id, "Created via Planning Mode", `Initial plan: ${(session.initialPlan ?? "").slice(0, 200)}`));
|
||||||
await finalizePlanningTaskCreation(sessionId, claimOwnerToken, task.id);
|
await finalizePlanningTaskCreation(sessionId, claimOwnerToken, task.id, claimEpoch);
|
||||||
await markSessionComplete();
|
await markSessionComplete();
|
||||||
return { task, alreadyCreated: false };
|
return { task, alreadyCreated: false };
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
// A raced insert under the same key is the idempotent success case; anything else releases the claim.
|
// A raced insert under the same key is the idempotent success case; anything else releases the claim.
|
||||||
const raced = await findCreatedTask().catch(() => undefined);
|
const raced = await findCreatedTask().catch(() => undefined);
|
||||||
await releasePlanningTaskCreation(sessionId, claimOwnerToken).catch(() => undefined);
|
await releasePlanningTaskCreation(sessionId, claimOwnerToken).catch(() => undefined);
|
||||||
if (raced) return returnExisting(raced);
|
if (raced) {
|
||||||
|
const recovered = await returnExisting(raced);
|
||||||
|
// A post-insert failure (e.g. finalize) lands here for the task WE just created — it is not "already created".
|
||||||
|
return insertedTask && raced.id === insertedTask.id ? { task: recovered.task, alreadyCreated: false } : recovered;
|
||||||
|
}
|
||||||
throw err;
|
throw err;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -4063,26 +4146,28 @@ export async function getDurablePlanningSession(sessionId: string): Promise<Sess
|
|||||||
return row?.type === "planning" ? restoreClaimSession(row) : undefined;
|
return row?.type === "planning" ? restoreClaimSession(row) : undefined;
|
||||||
}
|
}
|
||||||
|
|
||||||
/** Atomically claim a planning session for its one task creation. */
|
/** Atomically claim a planning session for one creation epoch's task. */
|
||||||
export async function claimPlanningTaskCreation(sessionId: string, ownerToken: string, startedAt: string): Promise<Session | undefined> {
|
export async function claimPlanningTaskCreation(sessionId: string, ownerToken: string, startedAt: string, expectedTaskCreationEpoch?: number): Promise<Session | undefined> {
|
||||||
if (!_aiSessionStore || typeof (_aiSessionStore as unknown as { claimPlanningTaskCreation?: unknown }).claimPlanningTaskCreation !== "function") {
|
if (!_aiSessionStore || typeof (_aiSessionStore as unknown as { claimPlanningTaskCreation?: unknown }).claimPlanningTaskCreation !== "function") {
|
||||||
const session = await getSession(sessionId);
|
const session = await getSession(sessionId);
|
||||||
if (!session || session.createClaimStatus === "creating" || session.createClaimStatus === "created") return undefined;
|
if (!session || session.createClaimStatus === "creating" || session.createClaimStatus === "created") return undefined;
|
||||||
|
if (expectedTaskCreationEpoch !== undefined && (session.taskCreationEpoch ?? 0) !== expectedTaskCreationEpoch) return undefined;
|
||||||
Object.assign(session, { createClaimStatus: "creating", claimOwnerToken: ownerToken, claimStartedAt: startedAt });
|
Object.assign(session, { createClaimStatus: "creating", claimOwnerToken: ownerToken, claimStartedAt: startedAt });
|
||||||
return session;
|
return session;
|
||||||
}
|
}
|
||||||
const row = await _aiSessionStore.claimPlanningTaskCreation(sessionId, ownerToken, startedAt);
|
const row = await _aiSessionStore.claimPlanningTaskCreation(sessionId, ownerToken, startedAt, expectedTaskCreationEpoch);
|
||||||
return row ? restoreClaimSession(row) : undefined;
|
return row ? restoreClaimSession(row) : undefined;
|
||||||
}
|
}
|
||||||
|
|
||||||
export async function finalizePlanningTaskCreation(sessionId: string, ownerToken: string, taskId: string): Promise<Session | undefined> {
|
export async function finalizePlanningTaskCreation(sessionId: string, ownerToken: string, taskId: string, expectedTaskCreationEpoch?: number): Promise<Session | undefined> {
|
||||||
if (!_aiSessionStore || typeof (_aiSessionStore as unknown as { finalizePlanningTaskCreation?: unknown }).finalizePlanningTaskCreation !== "function") {
|
if (!_aiSessionStore || typeof (_aiSessionStore as unknown as { finalizePlanningTaskCreation?: unknown }).finalizePlanningTaskCreation !== "function") {
|
||||||
const session = await getSession(sessionId);
|
const session = await getSession(sessionId);
|
||||||
if (!session || session.claimOwnerToken !== ownerToken) return undefined;
|
if (!session || session.claimOwnerToken !== ownerToken) return undefined;
|
||||||
|
if (expectedTaskCreationEpoch !== undefined && (session.taskCreationEpoch ?? 0) !== expectedTaskCreationEpoch) return undefined;
|
||||||
Object.assign(session, { createClaimStatus: "created", createdTaskId: taskId, claimOwnerToken: undefined, claimStartedAt: undefined });
|
Object.assign(session, { createClaimStatus: "created", createdTaskId: taskId, claimOwnerToken: undefined, claimStartedAt: undefined });
|
||||||
return session;
|
return session;
|
||||||
}
|
}
|
||||||
const row = await _aiSessionStore.finalizePlanningTaskCreation(sessionId, ownerToken, taskId);
|
const row = await _aiSessionStore.finalizePlanningTaskCreation(sessionId, ownerToken, taskId, expectedTaskCreationEpoch);
|
||||||
return row ? restoreClaimSession(row) : undefined;
|
return row ? restoreClaimSession(row) : undefined;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -1328,10 +1328,36 @@ export function registerPlanningSubtaskRoutes(ctx: ApiRoutesContext, deps: Plann
|
|||||||
const returnLinkedTask = async (candidate = session) => {
|
const returnLinkedTask = async (candidate = session) => {
|
||||||
if (!candidate?.createdTaskId) return false;
|
if (!candidate?.createdTaskId) return false;
|
||||||
const linkedTask = await scopedStore.getTask(candidate.createdTaskId).catch(() => null);
|
const linkedTask = await scopedStore.getTask(candidate.createdTaskId).catch(() => null);
|
||||||
if (!linkedTask) throw conflict("PLANNING_CREATED_TASK_MISSING");
|
if (linkedTask) {
|
||||||
await markSessionComplete();
|
await markSessionComplete();
|
||||||
res.status(200).json({ task: linkedTask, alreadyCreated: true });
|
res.status(200).json({ task: linkedTask, alreadyCreated: true });
|
||||||
return true;
|
return true;
|
||||||
|
}
|
||||||
|
/*
|
||||||
|
FNXC:PlanningMultiTask 2026-07-24-03:20:
|
||||||
|
Reported bug: deleting the task created from a plan left the session permanently
|
||||||
|
dead-ended on PLANNING_CREATED_TASK_MISSING — Retry create replayed the same 409
|
||||||
|
forever. Distinguish "task deleted" from "transient read failure" using the
|
||||||
|
include-archived task scan (the same crash-window authority findCreatedTask uses):
|
||||||
|
if the linked id is still LISTED but getTask failed, keep failing closed (never fork
|
||||||
|
on a flaky read); if it is absent from the full list, the linkage is stale — clear it
|
||||||
|
so this request falls through and creates a fresh task under the current epoch key.
|
||||||
|
*/
|
||||||
|
const allTasks = await scopedStore.listTasks({ includeArchived: true }).catch(() => null);
|
||||||
|
const stillListed = allTasks === null || allTasks.some((task) => task.id === candidate.createdTaskId);
|
||||||
|
if (stillListed) throw conflict("PLANNING_CREATED_TASK_MISSING");
|
||||||
|
await runPlanningCreateSideEffect(
|
||||||
|
"Planning create-task stale linkage clear failed",
|
||||||
|
() => updatePlanningCreateClaim(sessionId, { createClaimStatus: "none", createdTaskId: undefined, claimOwnerToken: undefined, claimStartedAt: undefined }),
|
||||||
|
{ sessionId, staleTaskId: candidate.createdTaskId },
|
||||||
|
);
|
||||||
|
if (session) {
|
||||||
|
session.createdTaskId = undefined;
|
||||||
|
session.createClaimStatus = "none";
|
||||||
|
session.claimOwnerToken = undefined;
|
||||||
|
session.claimStartedAt = undefined;
|
||||||
|
}
|
||||||
|
return false;
|
||||||
};
|
};
|
||||||
|
|
||||||
// A task row is the crash-window authority. Reconcile it before trying to claim.
|
// A task row is the crash-window authority. Reconcile it before trying to claim.
|
||||||
@@ -1351,7 +1377,7 @@ export function registerPlanningSubtaskRoutes(ctx: ApiRoutesContext, deps: Plann
|
|||||||
const claimStartedAt = new Date().toISOString();
|
const claimStartedAt = new Date().toISOString();
|
||||||
const hasDurableClaimStore = typeof (aiSessionStore as unknown as { claimPlanningTaskCreation?: unknown } | undefined)?.claimPlanningTaskCreation === "function";
|
const hasDurableClaimStore = typeof (aiSessionStore as unknown as { claimPlanningTaskCreation?: unknown } | undefined)?.claimPlanningTaskCreation === "function";
|
||||||
let claimed = hasDurableClaimStore
|
let claimed = hasDurableClaimStore
|
||||||
? await claimPlanningTaskCreation(sessionId, claimOwnerToken, claimStartedAt)
|
? await claimPlanningTaskCreation(sessionId, claimOwnerToken, claimStartedAt, claimEpoch)
|
||||||
: session
|
: session
|
||||||
? session.createClaimStatus !== "creating" && session.createClaimStatus !== "created"
|
? session.createClaimStatus !== "creating" && session.createClaimStatus !== "created"
|
||||||
? (await updatePlanningCreateClaim(sessionId, { createClaimStatus: "creating", claimOwnerToken, claimStartedAt, createdTaskId: undefined }), session)
|
? (await updatePlanningCreateClaim(sessionId, { createClaimStatus: "creating", claimOwnerToken, claimStartedAt, createdTaskId: undefined }), session)
|
||||||
@@ -1374,7 +1400,7 @@ export function registerPlanningSubtaskRoutes(ctx: ApiRoutesContext, deps: Plann
|
|||||||
const leaseExpired = session?.createClaimStatus === "creating" && Number.isFinite(startedAt) && Date.now() - startedAt >= 30_000;
|
const leaseExpired = session?.createClaimStatus === "creating" && Number.isFinite(startedAt) && Date.now() - startedAt >= 30_000;
|
||||||
if (!leaseExpired || !session?.claimOwnerToken) throw conflict("Planning task creation is already in progress");
|
if (!leaseExpired || !session?.claimOwnerToken) throw conflict("Planning task creation is already in progress");
|
||||||
await releasePlanningTaskCreation(sessionId, session.claimOwnerToken);
|
await releasePlanningTaskCreation(sessionId, session.claimOwnerToken);
|
||||||
claimed = await claimPlanningTaskCreation(sessionId, claimOwnerToken, new Date().toISOString());
|
claimed = await claimPlanningTaskCreation(sessionId, claimOwnerToken, new Date().toISOString(), claimEpoch);
|
||||||
if (!claimed) throw conflict("Planning task creation is already in progress");
|
if (!claimed) throw conflict("Planning task creation is already in progress");
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1448,7 +1474,7 @@ export function registerPlanningSubtaskRoutes(ctx: ApiRoutesContext, deps: Plann
|
|||||||
// Write the linkage before responding. If this write is interrupted, the next retry
|
// Write the linkage before responding. If this write is interrupted, the next retry
|
||||||
// reconciles the unique proposalClaimId task mapping above and never inserts another task.
|
// reconciles the unique proposalClaimId task mapping above and never inserts another task.
|
||||||
if (hasDurableClaimStore) {
|
if (hasDurableClaimStore) {
|
||||||
await finalizePlanningTaskCreation(sessionId, claimOwnerToken, task.id);
|
await finalizePlanningTaskCreation(sessionId, claimOwnerToken, task.id, claimEpoch);
|
||||||
} else if (session) {
|
} else if (session) {
|
||||||
await updatePlanningCreateClaim(sessionId, { createClaimStatus: "created", createdTaskId: task.id, claimOwnerToken: undefined, claimStartedAt: undefined });
|
await updatePlanningCreateClaim(sessionId, { createClaimStatus: "created", createdTaskId: task.id, claimOwnerToken: undefined, claimStartedAt: undefined });
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user