FN-7951: harden runGenerationWithTimeout abort cancellation across planning surfaces
Ensures aborted AI generation (timeout, user-stop, displacement, retries) actually tears down the in-flight agent session instead of only rejecting the Promise.race waiter, since provider SDKs may ignore AbortSignal. - Add a once-only onAbort teardown hook to GenerationGuard, invoked for timeout, user-stop, and displaced abort causes so consumers can dispose their in-flight session exactly once. - Give planning's local generation runner (runGenerationWithTimeout) the same guaranteed once-only abortTeardown for timeout, user-stop, displacement, stuck, and loop aborts, replacing the ad hoc dispose-on-timeout-only logic. - Forward the AbortSignal into planning's history-replay prompt, turn prompts, and JSON-parse-retry prompts, and short-circuit with createAbortError() when the signal is already aborted before/after each prompt call. - Wire subtask-breakdown's onTimeout/onUserStop handlers to the new onAbort hook instead of disposing the agent directly, keeping teardown centralized in the guard. - Add GenerationInProgressError / TargetGenerationInProgressError handling in mission-routes to return 409 Conflict instead of a generic 500 when a generation is already running. - Extend mission-interview and milestone-slice-interview generation paths with matching abort-forwarding and teardown behavior, plus new/expanded tests covering cancellation across timeout, user-stop, displacement, and retry paths. - Add a patch changeset documenting the fix for @runfusion/fusion. Files changed: .changeset/harden-generation-abort.md | 7 ++ .../src/__tests__/ai-session-timeout.test.ts | 41 +++++-- .../__tests__/milestone-slice-interview.test.ts | 72 ++++++++++++- .../src/__tests__/mission-interview.test.ts | 64 ++++++++++- .../planning-generation-cancellation.test.ts | 82 ++++++++++++++ .../src/__tests__/subtask-breakdown.test.ts | 21 +++- packages/dashboard/src/ai-session-timeout.ts | 33 +++++- .../dashboard/src/milestone-slice-interview.ts | 120 +++++++++++++++++++-- packages/dashboard/src/mission-interview.ts | 119 ++++++++++++++++++-- packages/dashboard/src/mission-routes.ts | 12 +++ packages/dashboard/src/planning.ts | 70 +++++++++--- packages/dashboard/src/subtask-breakdown.ts | 10 +- 12 files changed, 589 insertions(+), 62 deletions(-) Fusion-Task-Id: FN-7951 Fusion-Task-Lineage: debcd6a9-f54e-4ef3-87e1-4f06be0b5f64 Co-authored-by: Fusion (runfusion.ai) <noreply@runfusion.ai>
This commit is contained in:
7
.changeset/harden-generation-abort.md
Normal file
7
.changeset/harden-generation-abort.md
Normal file
@@ -0,0 +1,7 @@
|
|||||||
|
---
|
||||||
|
"@runfusion/fusion": patch
|
||||||
|
---
|
||||||
|
|
||||||
|
summary: Stop abandoned AI-session prompts when planning and interview generations are aborted.
|
||||||
|
category: fix
|
||||||
|
dev: Forwards AbortSignal into guarded prompt calls and disposes in-flight agent sessions on abort.
|
||||||
@@ -24,12 +24,17 @@ describe("GenerationGuard", () => {
|
|||||||
expect(guard.has("s1")).toBe(false);
|
expect(guard.has("s1")).toBe(false);
|
||||||
});
|
});
|
||||||
|
|
||||||
it("fires onTimeout, aborts the operation, and rejects with AbortError when the timer fires", async () => {
|
it("fires onTimeout, tears down, aborts the operation, and rejects with AbortError when the timer fires", async () => {
|
||||||
const guard = new GenerationGuard();
|
const guard = new GenerationGuard();
|
||||||
const onTimeout = vi.fn();
|
const onTimeout = vi.fn();
|
||||||
|
const onAbort = vi.fn();
|
||||||
|
|
||||||
let opSettled = false;
|
let opSettled = false;
|
||||||
const promise = guard.run("s1", 1_000, { onTimeout }, async () => {
|
let promptSignalAborted = false;
|
||||||
|
const promise = guard.run("s1", 1_000, { onTimeout, onAbort }, async (abortSignal) => {
|
||||||
|
abortSignal.addEventListener("abort", () => {
|
||||||
|
promptSignalAborted = true;
|
||||||
|
}, { once: true });
|
||||||
// Simulate a hung prompt() that never resolves on its own.
|
// Simulate a hung prompt() that never resolves on its own.
|
||||||
await new Promise<void>(() => { /* intentionally never resolves */ });
|
await new Promise<void>(() => { /* intentionally never resolves */ });
|
||||||
opSettled = true;
|
opSettled = true;
|
||||||
@@ -48,6 +53,8 @@ describe("GenerationGuard", () => {
|
|||||||
expect(outcome.ok).toBe(false);
|
expect(outcome.ok).toBe(false);
|
||||||
expect(outcome.ok ? null : outcome.err).toSatisfy((err: unknown) => isAbortError(err));
|
expect(outcome.ok ? null : outcome.err).toSatisfy((err: unknown) => isAbortError(err));
|
||||||
expect(onTimeout).toHaveBeenCalledTimes(1);
|
expect(onTimeout).toHaveBeenCalledTimes(1);
|
||||||
|
expect(onAbort).toHaveBeenCalledTimes(1);
|
||||||
|
expect(promptSignalAborted).toBe(true);
|
||||||
expect(opSettled).toBe(false);
|
expect(opSettled).toBe(false);
|
||||||
expect(guard.has("s1")).toBe(false);
|
expect(guard.has("s1")).toBe(false);
|
||||||
});
|
});
|
||||||
@@ -56,8 +63,13 @@ describe("GenerationGuard", () => {
|
|||||||
const guard = new GenerationGuard();
|
const guard = new GenerationGuard();
|
||||||
const onTimeout = vi.fn();
|
const onTimeout = vi.fn();
|
||||||
const onUserStop = vi.fn();
|
const onUserStop = vi.fn();
|
||||||
|
const onAbort = vi.fn();
|
||||||
|
let promptSignalAborted = false;
|
||||||
|
|
||||||
const promise = guard.run("s1", 10_000, { onTimeout, onUserStop }, async () => {
|
const promise = guard.run("s1", 10_000, { onTimeout, onUserStop, onAbort }, async (abortSignal) => {
|
||||||
|
abortSignal.addEventListener("abort", () => {
|
||||||
|
promptSignalAborted = true;
|
||||||
|
}, { once: true });
|
||||||
await new Promise<void>(() => { /* hang */ });
|
await new Promise<void>(() => { /* hang */ });
|
||||||
});
|
});
|
||||||
const settled = promise.then(
|
const settled = promise.then(
|
||||||
@@ -72,6 +84,8 @@ describe("GenerationGuard", () => {
|
|||||||
expect(outcome.ok).toBe(false);
|
expect(outcome.ok).toBe(false);
|
||||||
expect(onTimeout).not.toHaveBeenCalled();
|
expect(onTimeout).not.toHaveBeenCalled();
|
||||||
expect(onUserStop).toHaveBeenCalledTimes(1);
|
expect(onUserStop).toHaveBeenCalledTimes(1);
|
||||||
|
expect(onAbort).toHaveBeenCalledTimes(1);
|
||||||
|
expect(promptSignalAborted).toBe(true);
|
||||||
|
|
||||||
// Subsequent stop is a no-op.
|
// Subsequent stop is a no-op.
|
||||||
expect(guard.stop("s1")).toBe(false);
|
expect(guard.stop("s1")).toBe(false);
|
||||||
@@ -80,12 +94,19 @@ describe("GenerationGuard", () => {
|
|||||||
it("re-entrant run for the same id aborts the prior generation", async () => {
|
it("re-entrant run for the same id aborts the prior generation", async () => {
|
||||||
const guard = new GenerationGuard();
|
const guard = new GenerationGuard();
|
||||||
const firstUserStop = vi.fn();
|
const firstUserStop = vi.fn();
|
||||||
|
const firstAbort = vi.fn();
|
||||||
|
let firstPromptSignalAborted = false;
|
||||||
|
|
||||||
const first = guard.run(
|
const first = guard.run(
|
||||||
"s1",
|
"s1",
|
||||||
10_000,
|
10_000,
|
||||||
{ onTimeout: vi.fn(), onUserStop: firstUserStop },
|
{ onTimeout: vi.fn(), onUserStop: firstUserStop, onAbort: firstAbort },
|
||||||
async () => { await new Promise<void>(() => { /* hang */ }); },
|
async (abortSignal) => {
|
||||||
|
abortSignal.addEventListener("abort", () => {
|
||||||
|
firstPromptSignalAborted = true;
|
||||||
|
}, { once: true });
|
||||||
|
await new Promise<void>(() => { /* hang */ });
|
||||||
|
},
|
||||||
);
|
);
|
||||||
const firstSettled = first.then(
|
const firstSettled = first.then(
|
||||||
() => ({ ok: true as const }),
|
() => ({ ok: true as const }),
|
||||||
@@ -101,16 +122,20 @@ describe("GenerationGuard", () => {
|
|||||||
// user-facing stop, so onUserStop is intentionally not fired for the
|
// user-facing stop, so onUserStop is intentionally not fired for the
|
||||||
// displaced generation. The displaced caller still observes AbortError.
|
// displaced generation. The displaced caller still observes AbortError.
|
||||||
expect(firstUserStop).not.toHaveBeenCalled();
|
expect(firstUserStop).not.toHaveBeenCalled();
|
||||||
|
expect(firstAbort).toHaveBeenCalledTimes(1);
|
||||||
|
expect(firstPromptSignalAborted).toBe(true);
|
||||||
expect(guard.has("s1")).toBe(false);
|
expect(guard.has("s1")).toBe(false);
|
||||||
});
|
});
|
||||||
|
|
||||||
it("reset() aborts every active generation", async () => {
|
it("reset() aborts every active generation", async () => {
|
||||||
const guard = new GenerationGuard();
|
const guard = new GenerationGuard();
|
||||||
|
|
||||||
const a = guard.run("a", 10_000, { onTimeout: vi.fn() }, async () => {
|
const abortA = vi.fn();
|
||||||
|
const abortB = vi.fn();
|
||||||
|
const a = guard.run("a", 10_000, { onTimeout: vi.fn(), onAbort: abortA }, async () => {
|
||||||
await new Promise<void>(() => { /* hang */ });
|
await new Promise<void>(() => { /* hang */ });
|
||||||
});
|
});
|
||||||
const b = guard.run("b", 10_000, { onTimeout: vi.fn() }, async () => {
|
const b = guard.run("b", 10_000, { onTimeout: vi.fn(), onAbort: abortB }, async () => {
|
||||||
await new Promise<void>(() => { /* hang */ });
|
await new Promise<void>(() => { /* hang */ });
|
||||||
});
|
});
|
||||||
const aSettled = a.catch((err) => err);
|
const aSettled = a.catch((err) => err);
|
||||||
@@ -120,6 +145,8 @@ describe("GenerationGuard", () => {
|
|||||||
|
|
||||||
expect(isAbortError(await aSettled)).toBe(true);
|
expect(isAbortError(await aSettled)).toBe(true);
|
||||||
expect(isAbortError(await bSettled)).toBe(true);
|
expect(isAbortError(await bSettled)).toBe(true);
|
||||||
|
expect(abortA).toHaveBeenCalledTimes(1);
|
||||||
|
expect(abortB).toHaveBeenCalledTimes(1);
|
||||||
expect(guard.has("a")).toBe(false);
|
expect(guard.has("a")).toBe(false);
|
||||||
expect(guard.has("b")).toBe(false);
|
expect(guard.has("b")).toBe(false);
|
||||||
});
|
});
|
||||||
|
|||||||
@@ -38,6 +38,7 @@ import {
|
|||||||
getTargetInterviewSummary,
|
getTargetInterviewSummary,
|
||||||
getRateLimitResetTime,
|
getRateLimitResetTime,
|
||||||
InvalidSessionStateError,
|
InvalidSessionStateError,
|
||||||
|
TargetGenerationInProgressError,
|
||||||
TargetInvalidSessionStateError,
|
TargetInvalidSessionStateError,
|
||||||
milestoneSliceInterviewStreamManager,
|
milestoneSliceInterviewStreamManager,
|
||||||
parseTargetInterviewResponse,
|
parseTargetInterviewResponse,
|
||||||
@@ -913,13 +914,16 @@ describe("milestone-slice-interview module", () => {
|
|||||||
vi.useFakeTimers();
|
vi.useFakeTimers();
|
||||||
|
|
||||||
let resolveHungPrompt: (() => void) | undefined;
|
let resolveHungPrompt: (() => void) | undefined;
|
||||||
|
let promptSignal: AbortSignal | undefined;
|
||||||
|
const dispose = vi.fn();
|
||||||
mockCreateFnAgent.mockImplementationOnce(async () => ({
|
mockCreateFnAgent.mockImplementationOnce(async () => ({
|
||||||
session: {
|
session: {
|
||||||
state: { messages: [] },
|
state: { messages: [] },
|
||||||
prompt: vi.fn(async () => {
|
prompt: vi.fn(async (_message: string, options?: { signal?: AbortSignal }) => {
|
||||||
|
promptSignal = options?.signal;
|
||||||
await new Promise<void>((resolve) => { resolveHungPrompt = resolve; });
|
await new Promise<void>((resolve) => { resolveHungPrompt = resolve; });
|
||||||
}),
|
}),
|
||||||
dispose: vi.fn(),
|
dispose,
|
||||||
},
|
},
|
||||||
}));
|
}));
|
||||||
|
|
||||||
@@ -942,6 +946,8 @@ describe("milestone-slice-interview module", () => {
|
|||||||
|
|
||||||
const session = getTargetInterviewSession(sessionId);
|
const session = getTargetInterviewSession(sessionId);
|
||||||
expect(session?.error).toMatch(/timed out/i);
|
expect(session?.error).toMatch(/timed out/i);
|
||||||
|
expect(promptSignal?.aborted).toBe(true);
|
||||||
|
expect(dispose).toHaveBeenCalledTimes(1);
|
||||||
|
|
||||||
resolveHungPrompt?.();
|
resolveHungPrompt?.();
|
||||||
await vi.advanceTimersByTimeAsync(0);
|
await vi.advanceTimersByTimeAsync(0);
|
||||||
@@ -951,13 +957,16 @@ describe("milestone-slice-interview module", () => {
|
|||||||
|
|
||||||
it("stopMilestoneSliceInterviewGeneration aborts an in-flight session and marks it stopped", async () => {
|
it("stopMilestoneSliceInterviewGeneration aborts an in-flight session and marks it stopped", async () => {
|
||||||
let resolveHungPrompt: (() => void) | undefined;
|
let resolveHungPrompt: (() => void) | undefined;
|
||||||
|
let promptSignal: AbortSignal | undefined;
|
||||||
|
const dispose = vi.fn();
|
||||||
mockCreateFnAgent.mockImplementationOnce(async () => ({
|
mockCreateFnAgent.mockImplementationOnce(async () => ({
|
||||||
session: {
|
session: {
|
||||||
state: { messages: [] },
|
state: { messages: [] },
|
||||||
prompt: vi.fn(async () => {
|
prompt: vi.fn(async (_message: string, options?: { signal?: AbortSignal }) => {
|
||||||
|
promptSignal = options?.signal;
|
||||||
await new Promise<void>((resolve) => { resolveHungPrompt = resolve; });
|
await new Promise<void>((resolve) => { resolveHungPrompt = resolve; });
|
||||||
}),
|
}),
|
||||||
dispose: vi.fn(),
|
dispose,
|
||||||
},
|
},
|
||||||
}));
|
}));
|
||||||
|
|
||||||
@@ -982,10 +991,65 @@ describe("milestone-slice-interview module", () => {
|
|||||||
|
|
||||||
const session = getTargetInterviewSession(sessionId);
|
const session = getTargetInterviewSession(sessionId);
|
||||||
expect(session?.error).toMatch(/stopped by user/i);
|
expect(session?.error).toMatch(/stopped by user/i);
|
||||||
|
expect(promptSignal?.aborted).toBe(true);
|
||||||
|
expect(dispose).toHaveBeenCalledTimes(1);
|
||||||
expect(stopMilestoneSliceInterviewGeneration(sessionId)).toBe(false);
|
expect(stopMilestoneSliceInterviewGeneration(sessionId)).toBe(false);
|
||||||
|
|
||||||
resolveHungPrompt?.();
|
resolveHungPrompt?.();
|
||||||
await new Promise((resolve) => setTimeout(resolve, 0));
|
await new Promise((resolve) => setTimeout(resolve, 0));
|
||||||
});
|
});
|
||||||
|
|
||||||
|
it("rejects an overlapping submit instead of crashing when it races the prior generation's displaced-abort teardown", async () => {
|
||||||
|
/*
|
||||||
|
FNXC:AiSessionCancellation 2026-07-13-00:10:
|
||||||
|
Regression test for the onAbort teardown disposing the shared session.agent
|
||||||
|
on EVERY abort cause, including "displaced" (a re-entrant generationGuard.run()
|
||||||
|
call for the same session id). Before the TargetGenerationInProgressError
|
||||||
|
guard, a second overlapping submitTargetInterviewResponse call would observe
|
||||||
|
session.agent === undefined (cleared by the first call's displaced-abort
|
||||||
|
teardown) and crash with a TypeError instead of a clean, recoverable error.
|
||||||
|
*/
|
||||||
|
let promptCallIndex = 0;
|
||||||
|
let resolveHungPrompt: (() => void) | undefined;
|
||||||
|
const agent = {
|
||||||
|
session: {
|
||||||
|
state: { messages: [] as Array<{ role: string; content: string }> },
|
||||||
|
prompt: vi.fn(async () => {
|
||||||
|
promptCallIndex += 1;
|
||||||
|
if (promptCallIndex === 1) {
|
||||||
|
agent.session.state.messages.push({ role: "assistant", content: createQuestionJson("q-init") });
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
// Second call (the first submit's turn): hang so it is still
|
||||||
|
// registered in the guard when the second submit races in.
|
||||||
|
await new Promise<void>((resolve) => { resolveHungPrompt = resolve; });
|
||||||
|
}),
|
||||||
|
dispose: vi.fn(),
|
||||||
|
},
|
||||||
|
};
|
||||||
|
mockCreateFnAgent.mockImplementationOnce(async () => agent);
|
||||||
|
|
||||||
|
const sessionId = await createTargetInterviewSession(
|
||||||
|
"10.0.1.12",
|
||||||
|
"milestone",
|
||||||
|
"milestone-racing",
|
||||||
|
"Racing milestone interview",
|
||||||
|
undefined,
|
||||||
|
"/tmp/project",
|
||||||
|
MOCK_TASK_STORE,
|
||||||
|
);
|
||||||
|
await waitForCurrentQuestion(sessionId);
|
||||||
|
|
||||||
|
const first = submitTargetInterviewResponse(sessionId, { "q-init": "answer-1" });
|
||||||
|
// Give the first submit's generationGuard.run() a tick to register and
|
||||||
|
// reach the hung prompt() call before the second submit races in.
|
||||||
|
await new Promise((resolve) => setTimeout(resolve, 10));
|
||||||
|
|
||||||
|
await expect(submitTargetInterviewResponse(sessionId, { "q-init": "answer-2" }))
|
||||||
|
.rejects.toThrow(TargetGenerationInProgressError);
|
||||||
|
|
||||||
|
resolveHungPrompt?.();
|
||||||
|
await expect(first).resolves.toBeDefined();
|
||||||
|
});
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|||||||
@@ -37,6 +37,7 @@ import {
|
|||||||
getMissionInterviewSummary,
|
getMissionInterviewSummary,
|
||||||
getRateLimitResetTime,
|
getRateLimitResetTime,
|
||||||
listMissionInterviewDrafts,
|
listMissionInterviewDrafts,
|
||||||
|
GenerationInProgressError,
|
||||||
InvalidSessionStateError,
|
InvalidSessionStateError,
|
||||||
missionInterviewStreamManager,
|
missionInterviewStreamManager,
|
||||||
parseMissionAgentResponse,
|
parseMissionAgentResponse,
|
||||||
@@ -1175,13 +1176,16 @@ describe("mission-interview module", () => {
|
|||||||
vi.useFakeTimers();
|
vi.useFakeTimers();
|
||||||
|
|
||||||
let resolveHungPrompt: (() => void) | undefined;
|
let resolveHungPrompt: (() => void) | undefined;
|
||||||
|
let promptSignal: AbortSignal | undefined;
|
||||||
|
const dispose = vi.fn();
|
||||||
mockCreateFnAgent.mockImplementationOnce(async () => ({
|
mockCreateFnAgent.mockImplementationOnce(async () => ({
|
||||||
session: {
|
session: {
|
||||||
state: { messages: [] },
|
state: { messages: [] },
|
||||||
prompt: vi.fn(async () => {
|
prompt: vi.fn(async (_message: string, options?: { signal?: AbortSignal }) => {
|
||||||
|
promptSignal = options?.signal;
|
||||||
await new Promise<void>((resolve) => { resolveHungPrompt = resolve; });
|
await new Promise<void>((resolve) => { resolveHungPrompt = resolve; });
|
||||||
}),
|
}),
|
||||||
dispose: vi.fn(),
|
dispose,
|
||||||
},
|
},
|
||||||
}));
|
}));
|
||||||
|
|
||||||
@@ -1200,6 +1204,8 @@ describe("mission-interview module", () => {
|
|||||||
|
|
||||||
const session = getMissionInterviewSession(sessionId);
|
const session = getMissionInterviewSession(sessionId);
|
||||||
expect(session?.error).toMatch(/timed out/i);
|
expect(session?.error).toMatch(/timed out/i);
|
||||||
|
expect(promptSignal?.aborted).toBe(true);
|
||||||
|
expect(dispose).toHaveBeenCalledTimes(1);
|
||||||
|
|
||||||
resolveHungPrompt?.();
|
resolveHungPrompt?.();
|
||||||
await vi.advanceTimersByTimeAsync(0);
|
await vi.advanceTimersByTimeAsync(0);
|
||||||
@@ -1211,13 +1217,16 @@ describe("mission-interview module", () => {
|
|||||||
// Real timers — we want the guard.run() registration to actually happen
|
// Real timers — we want the guard.run() registration to actually happen
|
||||||
// through normal microtask scheduling without us racing it.
|
// through normal microtask scheduling without us racing it.
|
||||||
let resolveHungPrompt: (() => void) | undefined;
|
let resolveHungPrompt: (() => void) | undefined;
|
||||||
|
let promptSignal: AbortSignal | undefined;
|
||||||
|
const dispose = vi.fn();
|
||||||
mockCreateFnAgent.mockImplementationOnce(async () => ({
|
mockCreateFnAgent.mockImplementationOnce(async () => ({
|
||||||
session: {
|
session: {
|
||||||
state: { messages: [] },
|
state: { messages: [] },
|
||||||
prompt: vi.fn(async () => {
|
prompt: vi.fn(async (_message: string, options?: { signal?: AbortSignal }) => {
|
||||||
|
promptSignal = options?.signal;
|
||||||
await new Promise<void>((resolve) => { resolveHungPrompt = resolve; });
|
await new Promise<void>((resolve) => { resolveHungPrompt = resolve; });
|
||||||
}),
|
}),
|
||||||
dispose: vi.fn(),
|
dispose,
|
||||||
},
|
},
|
||||||
}));
|
}));
|
||||||
|
|
||||||
@@ -1242,6 +1251,8 @@ describe("mission-interview module", () => {
|
|||||||
|
|
||||||
const session = getMissionInterviewSession(sessionId);
|
const session = getMissionInterviewSession(sessionId);
|
||||||
expect(session?.error).toMatch(/stopped by user/i);
|
expect(session?.error).toMatch(/stopped by user/i);
|
||||||
|
expect(promptSignal?.aborted).toBe(true);
|
||||||
|
expect(dispose).toHaveBeenCalledTimes(1);
|
||||||
|
|
||||||
// Stop is idempotent — no in-flight generation after first call.
|
// Stop is idempotent — no in-flight generation after first call.
|
||||||
expect(stopMissionInterviewGeneration(sessionId)).toBe(false);
|
expect(stopMissionInterviewGeneration(sessionId)).toBe(false);
|
||||||
@@ -1249,5 +1260,50 @@ describe("mission-interview module", () => {
|
|||||||
resolveHungPrompt?.();
|
resolveHungPrompt?.();
|
||||||
await new Promise((resolve) => setTimeout(resolve, 0));
|
await new Promise((resolve) => setTimeout(resolve, 0));
|
||||||
});
|
});
|
||||||
|
|
||||||
|
it("rejects an overlapping submit instead of crashing when it races the prior generation's displaced-abort teardown", async () => {
|
||||||
|
/*
|
||||||
|
FNXC:AiSessionCancellation 2026-07-13-00:10:
|
||||||
|
Regression test for the onAbort teardown disposing the shared session.agent
|
||||||
|
on EVERY abort cause, including "displaced" (a re-entrant generationGuard.run()
|
||||||
|
call for the same session id). Before the GenerationInProgressError guard, a
|
||||||
|
second overlapping submitMissionInterviewResponse call would observe
|
||||||
|
session.agent === undefined (cleared by the first call's displaced-abort
|
||||||
|
teardown) and crash with a TypeError instead of a clean, recoverable error.
|
||||||
|
*/
|
||||||
|
let promptCallIndex = 0;
|
||||||
|
let resolveHungPrompt: (() => void) | undefined;
|
||||||
|
const agent = {
|
||||||
|
session: {
|
||||||
|
state: { messages: [] as Array<{ role: string; content: string }> },
|
||||||
|
prompt: vi.fn(async () => {
|
||||||
|
promptCallIndex += 1;
|
||||||
|
if (promptCallIndex === 1) {
|
||||||
|
agent.session.state.messages.push({ role: "assistant", content: createQuestionJson("q-init") });
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
// Second call (the first submit's turn): hang so it is still
|
||||||
|
// registered in the guard when the second submit races in.
|
||||||
|
await new Promise<void>((resolve) => { resolveHungPrompt = resolve; });
|
||||||
|
}),
|
||||||
|
dispose: vi.fn(),
|
||||||
|
},
|
||||||
|
};
|
||||||
|
mockCreateFnAgent.mockImplementationOnce(async () => agent);
|
||||||
|
|
||||||
|
const sessionId = await createMissionInterviewSession("10.0.0.12", "Racing mission interview", "/tmp/project", MOCK_TASK_STORE);
|
||||||
|
await waitForCurrentQuestion(sessionId);
|
||||||
|
|
||||||
|
const first = submitMissionInterviewResponse(sessionId, { "q-init": "answer-1" });
|
||||||
|
// Give the first submit's generationGuard.run() a tick to register and
|
||||||
|
// reach the hung prompt() call before the second submit races in.
|
||||||
|
await new Promise((resolve) => setTimeout(resolve, 10));
|
||||||
|
|
||||||
|
await expect(submitMissionInterviewResponse(sessionId, { "q-init": "answer-2" }))
|
||||||
|
.rejects.toThrow(GenerationInProgressError);
|
||||||
|
|
||||||
|
resolveHungPrompt?.();
|
||||||
|
await expect(first).resolves.toBeDefined();
|
||||||
|
});
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|||||||
@@ -0,0 +1,82 @@
|
|||||||
|
// @vitest-environment node
|
||||||
|
|
||||||
|
import { beforeEach, describe, expect, it, vi } from "vitest";
|
||||||
|
import type { TaskStore } from "@fusion/core";
|
||||||
|
|
||||||
|
vi.mock("@fusion/engine", () => ({
|
||||||
|
listCliAdapterDescriptors: () => [],
|
||||||
|
resolveMcpServersForStore: async () => ({ servers: [] }),
|
||||||
|
buildSessionSkillContextSync: () => ({
|
||||||
|
skillSelectionContext: undefined,
|
||||||
|
resolvedSkillNames: ["fusion"],
|
||||||
|
skillSource: "role-fallback" as const,
|
||||||
|
}),
|
||||||
|
createFnAgent: vi.fn(),
|
||||||
|
createWorkflowAuthoringTools: () => [],
|
||||||
|
createChatTaskDocumentTools: () => [],
|
||||||
|
}));
|
||||||
|
|
||||||
|
import {
|
||||||
|
__resetPlanningState,
|
||||||
|
__setCreateFnAgent,
|
||||||
|
createSessionWithAgent,
|
||||||
|
getSession,
|
||||||
|
planningStreamManager,
|
||||||
|
stopGeneration,
|
||||||
|
} from "../planning.js";
|
||||||
|
|
||||||
|
const MOCK_TASK_STORE = {
|
||||||
|
listTasks: vi.fn(async () => []),
|
||||||
|
getTask: vi.fn(async () => {
|
||||||
|
throw new Error("not found");
|
||||||
|
}),
|
||||||
|
} as unknown as TaskStore;
|
||||||
|
|
||||||
|
describe("planning generation cancellation", () => {
|
||||||
|
beforeEach(() => {
|
||||||
|
__resetPlanningState();
|
||||||
|
});
|
||||||
|
|
||||||
|
it("forwards AbortSignal and disposes the in-flight planning prompt on user stop", async () => {
|
||||||
|
let resolveHungPrompt: (() => void) | undefined;
|
||||||
|
let promptSignal: AbortSignal | undefined;
|
||||||
|
let promptResolvedAfterAbort = false;
|
||||||
|
const dispose = vi.fn();
|
||||||
|
|
||||||
|
__setCreateFnAgent(vi.fn(async () => ({
|
||||||
|
session: {
|
||||||
|
state: { messages: [] },
|
||||||
|
prompt: vi.fn(async (_message: string, options?: { signal?: AbortSignal }) => {
|
||||||
|
promptSignal = options?.signal;
|
||||||
|
await new Promise<void>((resolve) => { resolveHungPrompt = resolve; });
|
||||||
|
promptResolvedAfterAbort = Boolean(promptSignal?.aborted);
|
||||||
|
}),
|
||||||
|
dispose,
|
||||||
|
},
|
||||||
|
})) as any);
|
||||||
|
|
||||||
|
const sessionId = await createSessionWithAgent(
|
||||||
|
"10.0.2.10",
|
||||||
|
"Plan a cancellable session",
|
||||||
|
"/tmp/project",
|
||||||
|
MOCK_TASK_STORE,
|
||||||
|
);
|
||||||
|
|
||||||
|
planningStreamManager.consumeInitialTurn(sessionId)?.();
|
||||||
|
for (let i = 0; i < 10 && !promptSignal; i++) {
|
||||||
|
await new Promise((resolve) => setTimeout(resolve, 0));
|
||||||
|
}
|
||||||
|
expect(promptSignal).toBeDefined();
|
||||||
|
|
||||||
|
expect(stopGeneration(sessionId)).toBe(true);
|
||||||
|
await new Promise((resolve) => setTimeout(resolve, 0));
|
||||||
|
|
||||||
|
expect(promptSignal?.aborted).toBe(true);
|
||||||
|
expect(dispose).toHaveBeenCalledTimes(1);
|
||||||
|
expect(getSession(sessionId)?.error).toMatch(/stopped by user/i);
|
||||||
|
|
||||||
|
resolveHungPrompt?.();
|
||||||
|
await new Promise((resolve) => setTimeout(resolve, 0));
|
||||||
|
expect(promptResolvedAfterAbort).toBe(true);
|
||||||
|
});
|
||||||
|
});
|
||||||
@@ -1094,8 +1094,10 @@ describe("subtask generation timeout / abort", () => {
|
|||||||
vi.useFakeTimers();
|
vi.useFakeTimers();
|
||||||
|
|
||||||
let resolveHungPrompt: (() => void) | undefined;
|
let resolveHungPrompt: (() => void) | undefined;
|
||||||
|
let promptSignal: AbortSignal | undefined;
|
||||||
const dispose = vi.fn();
|
const dispose = vi.fn();
|
||||||
const hungPromptCallable = vi.fn(async () => {
|
const hungPromptCallable = vi.fn(async (_message: string, options?: { signal?: AbortSignal }) => {
|
||||||
|
promptSignal = options?.signal;
|
||||||
// Simulate a stalled provider stream that never terminates on its own.
|
// Simulate a stalled provider stream that never terminates on its own.
|
||||||
await new Promise<void>((resolve) => { resolveHungPrompt = resolve; });
|
await new Promise<void>((resolve) => { resolveHungPrompt = resolve; });
|
||||||
return undefined;
|
return undefined;
|
||||||
@@ -1131,6 +1133,8 @@ describe("subtask generation timeout / abort", () => {
|
|||||||
expect(after?.status).toBe("error");
|
expect(after?.status).toBe("error");
|
||||||
expect(after?.error).toMatch(/timed out/i);
|
expect(after?.error).toMatch(/timed out/i);
|
||||||
expect(events).toContainEqual(expect.objectContaining({ type: "error", data: expect.stringMatching(/timed out/i) }));
|
expect(events).toContainEqual(expect.objectContaining({ type: "error", data: expect.stringMatching(/timed out/i) }));
|
||||||
|
expect(promptSignal?.aborted).toBe(true);
|
||||||
|
expect(dispose).toHaveBeenCalledTimes(1);
|
||||||
|
|
||||||
// The hung prompt is still pending; release it so its microtask completes.
|
// The hung prompt is still pending; release it so its microtask completes.
|
||||||
resolveHungPrompt?.();
|
resolveHungPrompt?.();
|
||||||
@@ -1270,13 +1274,16 @@ describe("subtask generation timeout / abort", () => {
|
|||||||
vi.useFakeTimers();
|
vi.useFakeTimers();
|
||||||
|
|
||||||
let resolveHungPrompt: (() => void) | undefined;
|
let resolveHungPrompt: (() => void) | undefined;
|
||||||
|
let promptSignal: AbortSignal | undefined;
|
||||||
|
const dispose = vi.fn();
|
||||||
mockCreateFnAgent.mockImplementation(async () => ({
|
mockCreateFnAgent.mockImplementation(async () => ({
|
||||||
session: {
|
session: {
|
||||||
state: { messages: [] },
|
state: { messages: [] },
|
||||||
prompt: vi.fn(async () => {
|
prompt: vi.fn(async (_message: string, options?: { signal?: AbortSignal }) => {
|
||||||
|
promptSignal = options?.signal;
|
||||||
await new Promise<void>((resolve) => { resolveHungPrompt = resolve; });
|
await new Promise<void>((resolve) => { resolveHungPrompt = resolve; });
|
||||||
}),
|
}),
|
||||||
dispose: vi.fn(),
|
dispose,
|
||||||
},
|
},
|
||||||
}));
|
}));
|
||||||
|
|
||||||
@@ -1285,7 +1292,11 @@ describe("subtask generation timeout / abort", () => {
|
|||||||
undefined,
|
undefined,
|
||||||
"/tmp/project",
|
"/tmp/project",
|
||||||
);
|
);
|
||||||
await Promise.resolve();
|
for (let i = 0; i < 10 && !promptSignal; i++) {
|
||||||
|
await vi.advanceTimersByTimeAsync(0);
|
||||||
|
await Promise.resolve();
|
||||||
|
}
|
||||||
|
expect(promptSignal).toBeDefined();
|
||||||
|
|
||||||
expect(stopSubtaskGeneration(created.sessionId)).toBe(true);
|
expect(stopSubtaskGeneration(created.sessionId)).toBe(true);
|
||||||
await vi.advanceTimersByTimeAsync(0);
|
await vi.advanceTimersByTimeAsync(0);
|
||||||
@@ -1293,6 +1304,8 @@ describe("subtask generation timeout / abort", () => {
|
|||||||
const after = getSubtaskSession(created.sessionId);
|
const after = getSubtaskSession(created.sessionId);
|
||||||
expect(after?.status).toBe("error");
|
expect(after?.status).toBe("error");
|
||||||
expect(after?.error).toMatch(/stopped by user/i);
|
expect(after?.error).toMatch(/stopped by user/i);
|
||||||
|
expect(promptSignal?.aborted).toBe(true);
|
||||||
|
expect(dispose).toHaveBeenCalledTimes(1);
|
||||||
|
|
||||||
// Stop is idempotent — no in-flight generation after first call.
|
// Stop is idempotent — no in-flight generation after first call.
|
||||||
expect(stopSubtaskGeneration(created.sessionId)).toBe(false);
|
expect(stopSubtaskGeneration(created.sessionId)).toBe(false);
|
||||||
|
|||||||
@@ -14,11 +14,18 @@ export interface TimeoutHandlers {
|
|||||||
onTimeout: () => void;
|
onTimeout: () => void;
|
||||||
/** Fired when abort happens for a non-timeout reason (e.g. manual stop). */
|
/** Fired when abort happens for a non-timeout reason (e.g. manual stop). */
|
||||||
onUserStop?: () => void;
|
onUserStop?: () => void;
|
||||||
|
/**
|
||||||
|
* FNXC:AiSessionCancellation 2026-07-13-00:00:
|
||||||
|
* FN-7951 requires aborting generation to stop the underlying prompt work, not just reject the Promise.race waiter. Producers use this once-only hook to dispose the in-flight agent session for timeout, user-stop, displacement, and reset because provider SDKs may ignore AbortSignal.
|
||||||
|
*/
|
||||||
|
onAbort?: () => void;
|
||||||
}
|
}
|
||||||
|
|
||||||
interface ActiveEntry {
|
interface ActiveEntry {
|
||||||
abort: AbortController;
|
abort: AbortController;
|
||||||
timer: ReturnType<typeof setTimeout>;
|
timer: ReturnType<typeof setTimeout>;
|
||||||
|
onAbort?: () => void;
|
||||||
|
onAbortFired: boolean;
|
||||||
}
|
}
|
||||||
|
|
||||||
type AbortCause = "timeout" | "user-stop" | "displaced";
|
type AbortCause = "timeout" | "user-stop" | "displaced";
|
||||||
@@ -47,17 +54,22 @@ export class GenerationGuard {
|
|||||||
this.cancelInternal(sessionId, "displaced");
|
this.cancelInternal(sessionId, "displaced");
|
||||||
|
|
||||||
const abort = new AbortController();
|
const abort = new AbortController();
|
||||||
const timer = setTimeout(() => {
|
const entry: ActiveEntry = {
|
||||||
|
abort,
|
||||||
|
timer: undefined as unknown as ReturnType<typeof setTimeout>,
|
||||||
|
onAbort: handlers.onAbort,
|
||||||
|
onAbortFired: false,
|
||||||
|
};
|
||||||
|
entry.timer = setTimeout(() => {
|
||||||
this.abortCause.set(abort, "timeout");
|
this.abortCause.set(abort, "timeout");
|
||||||
try {
|
try {
|
||||||
handlers.onTimeout();
|
handlers.onTimeout();
|
||||||
} catch {
|
} catch {
|
||||||
// swallow — handler errors must not prevent abort
|
// swallow — handler errors must not prevent abort
|
||||||
}
|
}
|
||||||
|
this.fireAbortTeardown(entry);
|
||||||
abort.abort();
|
abort.abort();
|
||||||
}, timeoutMs);
|
}, timeoutMs);
|
||||||
|
|
||||||
const entry: ActiveEntry = { abort, timer };
|
|
||||||
this.active.set(sessionId, entry);
|
this.active.set(sessionId, entry);
|
||||||
|
|
||||||
const abortPromise = new Promise<never>((_, reject) => {
|
const abortPromise = new Promise<never>((_, reject) => {
|
||||||
@@ -83,7 +95,7 @@ export class GenerationGuard {
|
|||||||
}
|
}
|
||||||
throw err;
|
throw err;
|
||||||
} finally {
|
} finally {
|
||||||
clearTimeout(timer);
|
clearTimeout(entry.timer);
|
||||||
this.abortCause.delete(abort);
|
this.abortCause.delete(abort);
|
||||||
if (this.active.get(sessionId) === entry) {
|
if (this.active.get(sessionId) === entry) {
|
||||||
this.active.delete(sessionId);
|
this.active.delete(sessionId);
|
||||||
@@ -121,10 +133,23 @@ export class GenerationGuard {
|
|||||||
if (!entry) return false;
|
if (!entry) return false;
|
||||||
clearTimeout(entry.timer);
|
clearTimeout(entry.timer);
|
||||||
this.abortCause.set(entry.abort, cause);
|
this.abortCause.set(entry.abort, cause);
|
||||||
|
this.fireAbortTeardown(entry);
|
||||||
entry.abort.abort();
|
entry.abort.abort();
|
||||||
this.active.delete(sessionId);
|
this.active.delete(sessionId);
|
||||||
return true;
|
return true;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private fireAbortTeardown(entry: ActiveEntry): void {
|
||||||
|
if (entry.onAbortFired) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
entry.onAbortFired = true;
|
||||||
|
try {
|
||||||
|
entry.onAbort?.();
|
||||||
|
} catch {
|
||||||
|
// swallow — teardown errors must not prevent abort propagation
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
export function createAbortError(): Error {
|
export function createAbortError(): Error {
|
||||||
|
|||||||
@@ -30,7 +30,7 @@ import {
|
|||||||
resetDiagnosticsSink,
|
resetDiagnosticsSink,
|
||||||
nonfatal,
|
nonfatal,
|
||||||
} from "./ai-session-diagnostics.js";
|
} from "./ai-session-diagnostics.js";
|
||||||
import { GenerationGuard, isAbortError } from "./ai-session-timeout.js";
|
import { createAbortError, GenerationGuard, isAbortError } from "./ai-session-timeout.js";
|
||||||
|
|
||||||
// Re-export JSON parsing utilities from mission-interview for external consumers
|
// Re-export JSON parsing utilities from mission-interview for external consumers
|
||||||
export {
|
export {
|
||||||
@@ -737,6 +737,36 @@ function disposeAgentForRetry(session: TargetInterviewSession): void {
|
|||||||
session.agent = undefined;
|
session.agent = undefined;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:AiSessionCancellation 2026-07-13-00:10:
|
||||||
|
guard.run()'s onAbort teardown fires for EVERY abort cause, including "displaced" (a re-entrant
|
||||||
|
generationGuard.run() call for the same session id triggers cancelInternal("displaced") on the
|
||||||
|
prior entry before the new op runs). Retry flows call disposeAgentForRetry(session) themselves and
|
||||||
|
then assign a brand-new session.agent BEFORE the retry's own generationGuard.run() call displaces
|
||||||
|
the stale (already-forgotten) entry from session creation/history-replay. If the stale entry's
|
||||||
|
onAbort teardown reads session.agent dynamically at teardown time (as disposeAgentForRetry does),
|
||||||
|
it disposes the FRESH agent the retry just installed — not the stale one — and the retry's own
|
||||||
|
operation then crashes on `session.agent!` being undefined. Capture the exact agent instance a
|
||||||
|
generation started with and only tear down / clear that specific instance, so a later displacement
|
||||||
|
can never dispose an agent installed by a newer call.
|
||||||
|
*/
|
||||||
|
function disposeAgentGeneration(session: TargetInterviewSession, agent: AgentResult | undefined): void {
|
||||||
|
if (!agent) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
nonfatal(
|
||||||
|
() => agent.session.dispose?.(),
|
||||||
|
diagnostics,
|
||||||
|
"Error disposing agent for retry",
|
||||||
|
{ sessionId: session.id, operation: "dispose-retry" }
|
||||||
|
);
|
||||||
|
|
||||||
|
if (session.agent === agent) {
|
||||||
|
session.agent = undefined;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// ── AI Agent Integration ───────────────────────────────────────────────────
|
// ── AI Agent Integration ───────────────────────────────────────────────────
|
||||||
|
|
||||||
function getSystemPrompt(targetType: TargetType): string {
|
function getSystemPrompt(targetType: TargetType): string {
|
||||||
@@ -886,6 +916,7 @@ async function ensureInterviewAgent(
|
|||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
const replayAgent = session.agent;
|
||||||
await generationGuard.run(
|
await generationGuard.run(
|
||||||
session.id,
|
session.id,
|
||||||
GENERATION_TIMEOUT_MS,
|
GENERATION_TIMEOUT_MS,
|
||||||
@@ -898,14 +929,28 @@ async function ensureInterviewAgent(
|
|||||||
session,
|
session,
|
||||||
"Generation stopped by user. You can retry or start a new session.",
|
"Generation stopped by user. You can retry or start a new session.",
|
||||||
),
|
),
|
||||||
|
onAbort: () => disposeAgentGeneration(session, replayAgent),
|
||||||
|
},
|
||||||
|
async (abortSignal) => {
|
||||||
|
/*
|
||||||
|
FNXC:AiSessionCancellation 2026-07-13-00:00:
|
||||||
|
FN-7951 requires every milestone/slice interview prompt, including history replay, to receive the generation AbortSignal. Promise.race only stops the caller from awaiting; signal forwarding plus guard-level session teardown is the cancellation contract.
|
||||||
|
*/
|
||||||
|
if (abortSignal.aborted) {
|
||||||
|
throw createAbortError();
|
||||||
|
}
|
||||||
|
await session.agent!.session.prompt(
|
||||||
|
[
|
||||||
|
"Previous conversation summary:",
|
||||||
|
historySummary,
|
||||||
|
"Use this context when handling the next user response.",
|
||||||
|
].join("\n\n"),
|
||||||
|
{ signal: abortSignal },
|
||||||
|
);
|
||||||
|
if (abortSignal.aborted) {
|
||||||
|
throw createAbortError();
|
||||||
|
}
|
||||||
},
|
},
|
||||||
() => session.agent!.session.prompt(
|
|
||||||
[
|
|
||||||
"Previous conversation summary:",
|
|
||||||
historySummary,
|
|
||||||
"Use this context when handling the next user response.",
|
|
||||||
].join("\n\n"),
|
|
||||||
),
|
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -951,6 +996,7 @@ async function continueAgentConversation(session: TargetInterviewSession, messag
|
|||||||
throw new TargetInvalidSessionStateError("AI agent not initialized");
|
throw new TargetInvalidSessionStateError("AI agent not initialized");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
const generationAgent = session.agent;
|
||||||
try {
|
try {
|
||||||
await generationGuard.run(
|
await generationGuard.run(
|
||||||
session.id,
|
session.id,
|
||||||
@@ -964,12 +1010,23 @@ async function continueAgentConversation(session: TargetInterviewSession, messag
|
|||||||
session,
|
session,
|
||||||
"Generation stopped by user. You can retry or start a new session.",
|
"Generation stopped by user. You can retry or start a new session.",
|
||||||
),
|
),
|
||||||
|
onAbort: () => disposeAgentGeneration(session, generationAgent),
|
||||||
},
|
},
|
||||||
async () => {
|
async (abortSignal) => {
|
||||||
const agent = session.agent!;
|
const agent = session.agent!;
|
||||||
session.thinkingOutput = "";
|
session.thinkingOutput = "";
|
||||||
|
|
||||||
await agent.session.prompt(message);
|
/*
|
||||||
|
FNXC:AiSessionCancellation 2026-07-13-00:00:
|
||||||
|
Milestone/slice interview turns and parse-retry prompts must pass the active AbortSignal to prompt() and short-circuit after abort. The GenerationGuard also tears down the agent session because provider SDKs may ignore the signal.
|
||||||
|
*/
|
||||||
|
if (abortSignal.aborted) {
|
||||||
|
throw createAbortError();
|
||||||
|
}
|
||||||
|
await agent.session.prompt(message, { signal: abortSignal });
|
||||||
|
if (abortSignal.aborted) {
|
||||||
|
throw createAbortError();
|
||||||
|
}
|
||||||
|
|
||||||
// Get the response text from the agent's state
|
// Get the response text from the agent's state
|
||||||
interface AgentMessage {
|
interface AgentMessage {
|
||||||
@@ -1010,12 +1067,19 @@ async function continueAgentConversation(session: TargetInterviewSession, messag
|
|||||||
);
|
);
|
||||||
try {
|
try {
|
||||||
session.thinkingOutput = "";
|
session.thinkingOutput = "";
|
||||||
|
if (abortSignal.aborted) {
|
||||||
|
throw createAbortError();
|
||||||
|
}
|
||||||
await agent.session.prompt(
|
await agent.session.prompt(
|
||||||
"Your previous response could not be parsed as JSON. " +
|
"Your previous response could not be parsed as JSON. " +
|
||||||
'Please respond with ONLY a valid JSON object: either {"type":"question","data":{...}} ' +
|
'Please respond with ONLY a valid JSON object: either {"type":"question","data":{...}} ' +
|
||||||
'or {"type":"complete","data":{"title":"...","description":"...","planningNotes":"...","verification":"..."}}' +
|
'or {"type":"complete","data":{"title":"...","description":"...","planningNotes":"...","verification":"..."}}' +
|
||||||
". No markdown, no explanation, just the JSON."
|
". No markdown, no explanation, just the JSON.",
|
||||||
|
{ signal: abortSignal },
|
||||||
);
|
);
|
||||||
|
if (abortSignal.aborted) {
|
||||||
|
throw createAbortError();
|
||||||
|
}
|
||||||
|
|
||||||
const retryMessage = (agent.session.state.messages as AgentMessage[])
|
const retryMessage = (agent.session.state.messages as AgentMessage[])
|
||||||
.filter((m: AgentMessage) => m.role === "assistant")
|
.filter((m: AgentMessage) => m.role === "assistant")
|
||||||
@@ -1034,6 +1098,9 @@ async function continueAgentConversation(session: TargetInterviewSession, messag
|
|||||||
}
|
}
|
||||||
responseText = retryText;
|
responseText = retryText;
|
||||||
} catch (retryErr) {
|
} catch (retryErr) {
|
||||||
|
if (isAbortError(retryErr)) {
|
||||||
|
throw retryErr;
|
||||||
|
}
|
||||||
diagnostics.errorFromException("Retry prompt failed for session", retryErr, { sessionId: session.id, operation: "retry-prompt" });
|
diagnostics.errorFromException("Retry prompt failed for session", retryErr, { sessionId: session.id, operation: "retry-prompt" });
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
@@ -1158,6 +1225,14 @@ export async function submitTargetInterviewResponse(
|
|||||||
throw new TargetSessionNotFoundError(`Interview session ${sessionId} not found or expired`);
|
throw new TargetSessionNotFoundError(`Interview session ${sessionId} not found or expired`);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:AiSessionCancellation 2026-07-13-00:10:
|
||||||
|
Reject an overlapping submit instead of letting generationGuard.run()'s displaced-abort teardown dispose the shared session.agent out from under this call (see TargetGenerationInProgressError doc).
|
||||||
|
*/
|
||||||
|
if (generationGuard.has(sessionId)) {
|
||||||
|
throw new TargetGenerationInProgressError("Generation already in progress for this response");
|
||||||
|
}
|
||||||
|
|
||||||
if (!session.currentQuestion) {
|
if (!session.currentQuestion) {
|
||||||
throw new TargetInvalidSessionStateError("No active question in session");
|
throw new TargetInvalidSessionStateError("No active question in session");
|
||||||
}
|
}
|
||||||
@@ -1224,6 +1299,18 @@ export async function retryTargetInterviewSession(
|
|||||||
throw new TargetInvalidSessionStateError(`Interview session ${sessionId} is not in an error state`);
|
throw new TargetInvalidSessionStateError(`Interview session ${sessionId} is not in an error state`);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:AiSessionCancellation 2026-07-13-00:10:
|
||||||
|
A session can be observed in "error" (persisted status) while its original fire-and-forget
|
||||||
|
initializeAgent() first turn is still actually in flight (createTargetInterviewSession never
|
||||||
|
awaits it). Retrying while that generation is still registered would race two concurrent
|
||||||
|
continueAgentConversation calls over the single shared session.agent slot. Reject cleanly instead,
|
||||||
|
matching the mission-interview.ts retryMissionInterviewSession guard for the identical race.
|
||||||
|
*/
|
||||||
|
if (generationGuard.has(sessionId)) {
|
||||||
|
throw new TargetGenerationInProgressError("Generation already in progress for this session");
|
||||||
|
}
|
||||||
|
|
||||||
disposeAgentForRetry(session);
|
disposeAgentForRetry(session);
|
||||||
|
|
||||||
session.error = undefined;
|
session.error = undefined;
|
||||||
@@ -1435,6 +1522,17 @@ export class TargetInvalidSessionStateError extends Error {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:AiSessionCancellation 2026-07-13-00:10:
|
||||||
|
FN-7951's onAbort teardown disposes the shared session.agent on every abort cause, including "displaced" (a re-entrant generationGuard.run() call for the same session id). The continued-conversation operation reads session.agent synchronously at the start of its op closure, so a second overlapping call for the same session would observe session.agent === undefined (cleared by the first call's displaced-abort teardown) and crash with a TypeError instead of a clean, recoverable error. Reject overlapping generations up front so the shared agent handle is never raced.
|
||||||
|
*/
|
||||||
|
export class TargetGenerationInProgressError extends Error {
|
||||||
|
constructor(message: string) {
|
||||||
|
super(message);
|
||||||
|
this.name = "TargetGenerationInProgressError";
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Reset all milestone/slice interview state. Used for testing only.
|
* Reset all milestone/slice interview state. Used for testing only.
|
||||||
*/
|
*/
|
||||||
|
|||||||
@@ -27,7 +27,7 @@ import {
|
|||||||
resetDiagnosticsSink,
|
resetDiagnosticsSink,
|
||||||
nonfatal,
|
nonfatal,
|
||||||
} from "./ai-session-diagnostics.js";
|
} from "./ai-session-diagnostics.js";
|
||||||
import { GenerationGuard, isAbortError } from "./ai-session-timeout.js";
|
import { createAbortError, GenerationGuard, isAbortError } from "./ai-session-timeout.js";
|
||||||
|
|
||||||
import { buildSessionSkillContextSync, createFnAgent as engineCreateFnAgent, resolveMcpServersForStore } from "@fusion/engine";
|
import { buildSessionSkillContextSync, createFnAgent as engineCreateFnAgent, resolveMcpServersForStore } from "@fusion/engine";
|
||||||
import { createPlanningBoardTools } from "./planning-board-tools.js";
|
import { createPlanningBoardTools } from "./planning-board-tools.js";
|
||||||
@@ -832,6 +832,36 @@ function disposeMissionAgentForRetry(session: MissionInterviewSession): void {
|
|||||||
session.agent = undefined;
|
session.agent = undefined;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:AiSessionCancellation 2026-07-13-00:10:
|
||||||
|
guard.run()'s onAbort teardown fires for EVERY abort cause, including "displaced" (a re-entrant
|
||||||
|
generationGuard.run() call for the same session id triggers cancelInternal("displaced") on the
|
||||||
|
prior entry before the new op runs). Retry/rewind flows call disposeMissionAgentForRetry(session)
|
||||||
|
themselves and then assign a brand-new session.agent BEFORE the retry's own generationGuard.run()
|
||||||
|
call displaces the stale (already-forgotten) entry from session creation/history-replay. If the
|
||||||
|
stale entry's onAbort teardown reads session.agent dynamically at teardown time (as
|
||||||
|
disposeMissionAgentForRetry does), it disposes the FRESH agent the retry just installed — not the
|
||||||
|
stale one — and the retry's own operation then crashes on `session.agent!` being undefined.
|
||||||
|
Capture the exact agent instance a generation started with and only tear down / clear that
|
||||||
|
specific instance, so a later displacement can never dispose an agent installed by a newer call.
|
||||||
|
*/
|
||||||
|
function disposeMissionAgentGeneration(session: MissionInterviewSession, agent: AgentResult | undefined): void {
|
||||||
|
if (!agent) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
nonfatal(
|
||||||
|
() => agent.session.dispose?.(),
|
||||||
|
diagnostics,
|
||||||
|
"Error disposing agent for retry",
|
||||||
|
{ sessionId: session.id, operation: "dispose-retry" }
|
||||||
|
);
|
||||||
|
|
||||||
|
if (session.agent === agent) {
|
||||||
|
session.agent = undefined;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// ── AI Agent Integration ───────────────────────────────────────────────────
|
// ── AI Agent Integration ───────────────────────────────────────────────────
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -1024,6 +1054,7 @@ async function ensureMissionInterviewAgent(
|
|||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
const replayAgent = session.agent;
|
||||||
await generationGuard.run(
|
await generationGuard.run(
|
||||||
session.id,
|
session.id,
|
||||||
GENERATION_TIMEOUT_MS,
|
GENERATION_TIMEOUT_MS,
|
||||||
@@ -1036,14 +1067,28 @@ async function ensureMissionInterviewAgent(
|
|||||||
session,
|
session,
|
||||||
"Generation stopped by user. You can retry or start a new session.",
|
"Generation stopped by user. You can retry or start a new session.",
|
||||||
),
|
),
|
||||||
|
onAbort: () => disposeMissionAgentGeneration(session, replayAgent),
|
||||||
|
},
|
||||||
|
async (abortSignal) => {
|
||||||
|
/*
|
||||||
|
FNXC:AiSessionCancellation 2026-07-13-00:00:
|
||||||
|
FN-7951 requires every mission-interview prompt, including history replay, to receive the generation AbortSignal. Promise.race only stops the caller from awaiting; signal forwarding plus guard-level session teardown is the cancellation contract.
|
||||||
|
*/
|
||||||
|
if (abortSignal.aborted) {
|
||||||
|
throw createAbortError();
|
||||||
|
}
|
||||||
|
await session.agent!.session.prompt(
|
||||||
|
[
|
||||||
|
"Previous conversation summary:",
|
||||||
|
historySummary,
|
||||||
|
"Use this context when handling the next user response.",
|
||||||
|
].join("\n\n"),
|
||||||
|
{ signal: abortSignal },
|
||||||
|
);
|
||||||
|
if (abortSignal.aborted) {
|
||||||
|
throw createAbortError();
|
||||||
|
}
|
||||||
},
|
},
|
||||||
() => session.agent!.session.prompt(
|
|
||||||
[
|
|
||||||
"Previous conversation summary:",
|
|
||||||
historySummary,
|
|
||||||
"Use this context when handling the next user response.",
|
|
||||||
].join("\n\n"),
|
|
||||||
),
|
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1078,6 +1123,7 @@ async function continueAgentConversation(session: MissionInterviewSession, messa
|
|||||||
throw new InvalidSessionStateError("AI agent not initialized");
|
throw new InvalidSessionStateError("AI agent not initialized");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
const generationAgent = session.agent;
|
||||||
try {
|
try {
|
||||||
await generationGuard.run(
|
await generationGuard.run(
|
||||||
session.id,
|
session.id,
|
||||||
@@ -1091,12 +1137,23 @@ async function continueAgentConversation(session: MissionInterviewSession, messa
|
|||||||
session,
|
session,
|
||||||
"Generation stopped by user. You can retry or start a new session.",
|
"Generation stopped by user. You can retry or start a new session.",
|
||||||
),
|
),
|
||||||
|
onAbort: () => disposeMissionAgentGeneration(session, generationAgent),
|
||||||
},
|
},
|
||||||
async () => {
|
async (abortSignal) => {
|
||||||
const agent = session.agent!;
|
const agent = session.agent!;
|
||||||
session.thinkingOutput = "";
|
session.thinkingOutput = "";
|
||||||
|
|
||||||
await agent.session.prompt(message);
|
/*
|
||||||
|
FNXC:AiSessionCancellation 2026-07-13-00:00:
|
||||||
|
Mission interview turns and parse-retry prompts must pass the active AbortSignal to prompt() and short-circuit after abort. The GenerationGuard also tears down the agent session because provider SDKs may ignore the signal.
|
||||||
|
*/
|
||||||
|
if (abortSignal.aborted) {
|
||||||
|
throw createAbortError();
|
||||||
|
}
|
||||||
|
await agent.session.prompt(message, { signal: abortSignal });
|
||||||
|
if (abortSignal.aborted) {
|
||||||
|
throw createAbortError();
|
||||||
|
}
|
||||||
|
|
||||||
// Get the response text from the agent's state
|
// Get the response text from the agent's state
|
||||||
interface AgentMessage {
|
interface AgentMessage {
|
||||||
@@ -1137,12 +1194,19 @@ async function continueAgentConversation(session: MissionInterviewSession, messa
|
|||||||
);
|
);
|
||||||
try {
|
try {
|
||||||
session.thinkingOutput = "";
|
session.thinkingOutput = "";
|
||||||
|
if (abortSignal.aborted) {
|
||||||
|
throw createAbortError();
|
||||||
|
}
|
||||||
await agent.session.prompt(
|
await agent.session.prompt(
|
||||||
"Your previous response could not be parsed as JSON. " +
|
"Your previous response could not be parsed as JSON. " +
|
||||||
'Please respond with ONLY a valid JSON object: either {"type":"question","data":{...}} ' +
|
'Please respond with ONLY a valid JSON object: either {"type":"question","data":{...}} ' +
|
||||||
'or {"type":"complete","data":{"missionTitle":"...","missionDescription":"...","milestones":[...]}}. ' +
|
'or {"type":"complete","data":{"missionTitle":"...","missionDescription":"...","milestones":[...]}}. ' +
|
||||||
"No markdown, no explanation, just the JSON."
|
"No markdown, no explanation, just the JSON.",
|
||||||
|
{ signal: abortSignal },
|
||||||
);
|
);
|
||||||
|
if (abortSignal.aborted) {
|
||||||
|
throw createAbortError();
|
||||||
|
}
|
||||||
|
|
||||||
const retryMessage = (agent.session.state.messages as AgentMessage[])
|
const retryMessage = (agent.session.state.messages as AgentMessage[])
|
||||||
.filter((m: AgentMessage) => m.role === "assistant")
|
.filter((m: AgentMessage) => m.role === "assistant")
|
||||||
@@ -1161,6 +1225,9 @@ async function continueAgentConversation(session: MissionInterviewSession, messa
|
|||||||
}
|
}
|
||||||
responseText = retryText;
|
responseText = retryText;
|
||||||
} catch (retryErr) {
|
} catch (retryErr) {
|
||||||
|
if (isAbortError(retryErr)) {
|
||||||
|
throw retryErr;
|
||||||
|
}
|
||||||
diagnostics.errorFromException("Retry prompt failed for session", retryErr, { sessionId: session.id, operation: "retry-prompt" });
|
diagnostics.errorFromException("Retry prompt failed for session", retryErr, { sessionId: session.id, operation: "retry-prompt" });
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
@@ -1303,6 +1370,14 @@ export async function submitMissionInterviewResponse(
|
|||||||
if (store && !session.store) session.store = store;
|
if (store && !session.store) session.store = store;
|
||||||
if (rootDir && !session.rootDir) session.rootDir = rootDir;
|
if (rootDir && !session.rootDir) session.rootDir = rootDir;
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:AiSessionCancellation 2026-07-13-00:10:
|
||||||
|
Reject an overlapping submit instead of letting generationGuard.run()'s displaced-abort teardown dispose the shared session.agent out from under this call (see GenerationInProgressError doc).
|
||||||
|
*/
|
||||||
|
if (generationGuard.has(sessionId)) {
|
||||||
|
throw new GenerationInProgressError("Generation already in progress for this response");
|
||||||
|
}
|
||||||
|
|
||||||
if (!session.currentQuestion) {
|
if (!session.currentQuestion) {
|
||||||
throw new InvalidSessionStateError("No active question in session");
|
throw new InvalidSessionStateError("No active question in session");
|
||||||
}
|
}
|
||||||
@@ -1368,6 +1443,17 @@ export async function retryMissionInterviewSession(
|
|||||||
throw new InvalidSessionStateError(`Mission interview session ${sessionId} is not in an error state`);
|
throw new InvalidSessionStateError(`Mission interview session ${sessionId} is not in an error state`);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:AiSessionCancellation 2026-07-13-00:10:
|
||||||
|
A session can be observed in "error" (persisted status) while its original fire-and-forget
|
||||||
|
initializeAgent() first turn is still actually in flight (createMissionInterviewSession never
|
||||||
|
awaits it). Retrying while that generation is still registered would race two concurrent
|
||||||
|
continueAgentConversation calls over the single shared session.agent slot. Reject cleanly instead.
|
||||||
|
*/
|
||||||
|
if (generationGuard.has(sessionId)) {
|
||||||
|
throw new GenerationInProgressError("Generation already in progress for this session");
|
||||||
|
}
|
||||||
|
|
||||||
disposeMissionAgentForRetry(session);
|
disposeMissionAgentForRetry(session);
|
||||||
|
|
||||||
session.error = undefined;
|
session.error = undefined;
|
||||||
@@ -1570,3 +1656,14 @@ export class InvalidSessionStateError extends Error {
|
|||||||
this.name = "InvalidSessionStateError";
|
this.name = "InvalidSessionStateError";
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:AiSessionCancellation 2026-07-13-00:10:
|
||||||
|
FN-7951's onAbort teardown disposes the single shared session.agent on every abort cause, including "displaced" (a re-entrant generationGuard.run() call for the same session id). continueAgentConversation's operation reads session.agent synchronously at the start of its op closure, so a second overlapping call for the same session would observe session.agent === undefined (cleared by the first call's displaced-abort teardown) and crash with a TypeError instead of a clean, recoverable error. Reject overlapping generations up front so the shared agent handle is never raced.
|
||||||
|
*/
|
||||||
|
export class GenerationInProgressError extends Error {
|
||||||
|
constructor(message: string) {
|
||||||
|
super(message);
|
||||||
|
this.name = "GenerationInProgressError";
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -638,6 +638,8 @@ export function createMissionRouter(
|
|||||||
throw notFound(errMsg);
|
throw notFound(errMsg);
|
||||||
} else if (errName === "InvalidSessionStateError") {
|
} else if (errName === "InvalidSessionStateError") {
|
||||||
throw badRequest(errMsg);
|
throw badRequest(errMsg);
|
||||||
|
} else if (errName === "GenerationInProgressError") {
|
||||||
|
throw conflict(errMsg);
|
||||||
} else {
|
} else {
|
||||||
throw internalError(errMsg || "Failed to process response");
|
throw internalError(errMsg || "Failed to process response");
|
||||||
}
|
}
|
||||||
@@ -689,6 +691,8 @@ export function createMissionRouter(
|
|||||||
throw notFound(errMsg);
|
throw notFound(errMsg);
|
||||||
} else if (errName === "InvalidSessionStateError") {
|
} else if (errName === "InvalidSessionStateError") {
|
||||||
throw badRequest(errMsg);
|
throw badRequest(errMsg);
|
||||||
|
} else if (errName === "GenerationInProgressError") {
|
||||||
|
throw conflict(errMsg);
|
||||||
} else {
|
} else {
|
||||||
throw internalError(errMsg || "Failed to retry interview session");
|
throw internalError(errMsg || "Failed to retry interview session");
|
||||||
}
|
}
|
||||||
@@ -3383,6 +3387,8 @@ export function createMissionRouter(
|
|||||||
throw notFound(errMsg);
|
throw notFound(errMsg);
|
||||||
} else if (errName === "TargetInvalidSessionStateError") {
|
} else if (errName === "TargetInvalidSessionStateError") {
|
||||||
throw badRequest(errMsg);
|
throw badRequest(errMsg);
|
||||||
|
} else if (errName === "TargetGenerationInProgressError") {
|
||||||
|
throw conflict(errMsg);
|
||||||
} else {
|
} else {
|
||||||
throw internalError(errMsg || "Failed to process response");
|
throw internalError(errMsg || "Failed to process response");
|
||||||
}
|
}
|
||||||
@@ -3547,6 +3553,8 @@ export function createMissionRouter(
|
|||||||
throw notFound(errMsg);
|
throw notFound(errMsg);
|
||||||
} else if (errName === "TargetInvalidSessionStateError") {
|
} else if (errName === "TargetInvalidSessionStateError") {
|
||||||
throw badRequest(errMsg);
|
throw badRequest(errMsg);
|
||||||
|
} else if (errName === "TargetGenerationInProgressError") {
|
||||||
|
throw conflict(errMsg);
|
||||||
} else {
|
} else {
|
||||||
throw internalError(errMsg || "Failed to retry interview session");
|
throw internalError(errMsg || "Failed to retry interview session");
|
||||||
}
|
}
|
||||||
@@ -3732,6 +3740,8 @@ export function createMissionRouter(
|
|||||||
throw notFound(errMsg);
|
throw notFound(errMsg);
|
||||||
} else if (errName === "TargetInvalidSessionStateError") {
|
} else if (errName === "TargetInvalidSessionStateError") {
|
||||||
throw badRequest(errMsg);
|
throw badRequest(errMsg);
|
||||||
|
} else if (errName === "TargetGenerationInProgressError") {
|
||||||
|
throw conflict(errMsg);
|
||||||
} else {
|
} else {
|
||||||
throw internalError(errMsg || "Failed to process response");
|
throw internalError(errMsg || "Failed to process response");
|
||||||
}
|
}
|
||||||
@@ -3896,6 +3906,8 @@ export function createMissionRouter(
|
|||||||
throw notFound(errMsg);
|
throw notFound(errMsg);
|
||||||
} else if (errName === "TargetInvalidSessionStateError") {
|
} else if (errName === "TargetInvalidSessionStateError") {
|
||||||
throw badRequest(errMsg);
|
throw badRequest(errMsg);
|
||||||
|
} else if (errName === "TargetGenerationInProgressError") {
|
||||||
|
throw conflict(errMsg);
|
||||||
} else {
|
} else {
|
||||||
throw internalError(errMsg || "Failed to retry interview session");
|
throw internalError(errMsg || "Failed to retry interview session");
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -433,6 +433,8 @@ interface ActivePlanningGeneration {
|
|||||||
abortController: AbortController;
|
abortController: AbortController;
|
||||||
timer: NodeJS.Timeout;
|
timer: NodeJS.Timeout;
|
||||||
abortReason?: PlanningGenerationAbortReason;
|
abortReason?: PlanningGenerationAbortReason;
|
||||||
|
abortTeardownFired: boolean;
|
||||||
|
abortTeardown: () => void;
|
||||||
markProgress: (output: string) => void;
|
markProgress: (output: string) => void;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -736,6 +738,9 @@ function cleanupInMemorySession(sessionId: string): boolean {
|
|||||||
const activeGeneration = activeGenerations.get(sessionId);
|
const activeGeneration = activeGenerations.get(sessionId);
|
||||||
if (activeGeneration) {
|
if (activeGeneration) {
|
||||||
clearTimeout(activeGeneration.timer);
|
clearTimeout(activeGeneration.timer);
|
||||||
|
activeGeneration.abortReason = "user-stop";
|
||||||
|
activeGeneration.abortTeardown();
|
||||||
|
activeGeneration.abortController.abort();
|
||||||
activeGenerations.delete(sessionId);
|
activeGenerations.delete(sessionId);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1896,7 +1901,21 @@ async function ensureSessionAgent(
|
|||||||
}
|
}
|
||||||
|
|
||||||
const contextMessage = buildHistoryReplayPrompt(historyForReplay);
|
const contextMessage = buildHistoryReplayPrompt(historyForReplay);
|
||||||
await session.agent.session.prompt(contextMessage);
|
await runGenerationWithTimeout(session, async (abortSignal) => {
|
||||||
|
/*
|
||||||
|
FNXC:AiSessionCancellation 2026-07-13-00:00:
|
||||||
|
Planning history replay is an agent prompt surface too. Forward the generation AbortSignal and rely on runGenerationWithTimeout to tear down the in-flight session because Promise.race alone cannot cancel prompt() work.
|
||||||
|
*/
|
||||||
|
if (abortSignal.aborted) {
|
||||||
|
throw createAbortError();
|
||||||
|
}
|
||||||
|
await (session.agent!.session.prompt as (input: string, options?: { signal?: AbortSignal }) => Promise<void>)(contextMessage, {
|
||||||
|
signal: abortSignal,
|
||||||
|
});
|
||||||
|
if (abortSignal.aborted) {
|
||||||
|
throw createAbortError();
|
||||||
|
}
|
||||||
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
async function maybeNotifyPlanningAwaitingInput(session: Session, question: PlanningQuestion): Promise<void> {
|
async function maybeNotifyPlanningAwaitingInput(session: Session, question: PlanningQuestion): Promise<void> {
|
||||||
@@ -1997,6 +2016,7 @@ async function runGenerationWithTimeout<T>(session: Session, operation: (abortSi
|
|||||||
if (existing) {
|
if (existing) {
|
||||||
clearTimeout(existing.timer);
|
clearTimeout(existing.timer);
|
||||||
existing.abortReason = "displaced";
|
existing.abortReason = "displaced";
|
||||||
|
existing.abortTeardown();
|
||||||
existing.abortController.abort();
|
existing.abortController.abort();
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -2017,8 +2037,8 @@ async function runGenerationWithTimeout<T>(session: Session, operation: (abortSi
|
|||||||
operation: "planning-generation-watchdog",
|
operation: "planning-generation-watchdog",
|
||||||
});
|
});
|
||||||
setSessionError(session, message);
|
setSessionError(session, message);
|
||||||
disposeSessionAgentForRetry(session);
|
|
||||||
}
|
}
|
||||||
|
generationRecord.abortTeardown();
|
||||||
abortController.abort();
|
abortController.abort();
|
||||||
};
|
};
|
||||||
|
|
||||||
@@ -2029,6 +2049,18 @@ async function runGenerationWithTimeout<T>(session: Session, operation: (abortSi
|
|||||||
const generationRecord: ActivePlanningGeneration = {
|
const generationRecord: ActivePlanningGeneration = {
|
||||||
abortController,
|
abortController,
|
||||||
timer: scheduleInactivityTimer(),
|
timer: scheduleInactivityTimer(),
|
||||||
|
abortTeardownFired: false,
|
||||||
|
abortTeardown: () => {
|
||||||
|
if (generationRecord.abortTeardownFired) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
generationRecord.abortTeardownFired = true;
|
||||||
|
/*
|
||||||
|
FNXC:AiSessionCancellation 2026-07-13-00:00:
|
||||||
|
Planning has a local generation runner instead of GenerationGuard. Abort teardown must run once for timeout, user-stop, displacement, stuck, and loop aborts so an abandoned prompt cannot continue after the Promise.race waiter rejects.
|
||||||
|
*/
|
||||||
|
disposeSessionAgentForRetry(session);
|
||||||
|
},
|
||||||
markProgress: (output: string) => {
|
markProgress: (output: string) => {
|
||||||
const signature = normalizeGenerationProgress(output);
|
const signature = normalizeGenerationProgress(output);
|
||||||
if (!signature) {
|
if (!signature) {
|
||||||
@@ -2139,11 +2171,19 @@ async function continueAgentConversation(session: Session, message: string): Pro
|
|||||||
// Clear thinking output for this turn
|
// Clear thinking output for this turn
|
||||||
session.thinkingOutput = "";
|
session.thinkingOutput = "";
|
||||||
|
|
||||||
// Send message to agent using .prompt() - it will stream thinking via onThinking callback.
|
/*
|
||||||
// Pass abort signal so timeout/user-stop can cancel the underlying prompt when supported.
|
FNXC:AiSessionCancellation 2026-07-13-00:00:
|
||||||
|
Planning turns and parse-retry prompts must pass the active AbortSignal to prompt() and short-circuit after abort. The local generation runner also tears down the agent session because provider SDKs may ignore the signal.
|
||||||
|
*/
|
||||||
|
if (abortSignal.aborted) {
|
||||||
|
throw createAbortError();
|
||||||
|
}
|
||||||
await (session.agent.session.prompt as (input: string, options?: { signal?: AbortSignal }) => Promise<void>)(message, {
|
await (session.agent.session.prompt as (input: string, options?: { signal?: AbortSignal }) => Promise<void>)(message, {
|
||||||
signal: abortSignal,
|
signal: abortSignal,
|
||||||
});
|
});
|
||||||
|
if (abortSignal.aborted) {
|
||||||
|
throw createAbortError();
|
||||||
|
}
|
||||||
|
|
||||||
// Get the response text from the agent's state
|
// Get the response text from the agent's state
|
||||||
interface AgentMessage {
|
interface AgentMessage {
|
||||||
@@ -2209,12 +2249,18 @@ async function continueAgentConversation(session: Session, message: string): Pro
|
|||||||
);
|
);
|
||||||
try {
|
try {
|
||||||
session.thinkingOutput = "";
|
session.thinkingOutput = "";
|
||||||
|
if (abortSignal.aborted) {
|
||||||
|
throw createAbortError();
|
||||||
|
}
|
||||||
await (session.agent.session.prompt as (input: string, options?: { signal?: AbortSignal }) => Promise<void>)(
|
await (session.agent.session.prompt as (input: string, options?: { signal?: AbortSignal }) => Promise<void>)(
|
||||||
"Your previous response could not be parsed as JSON. " +
|
"Your previous response could not be parsed as JSON. " +
|
||||||
'Please respond with ONLY a valid JSON object: either {"type":"question","data":{...}} ' +
|
'Please respond with ONLY a valid JSON object: either {"type":"question","data":{...}} ' +
|
||||||
'or {"type":"complete","data":{...}}. No markdown, no explanation, just the JSON.',
|
'or {"type":"complete","data":{...}}. No markdown, no explanation, just the JSON.',
|
||||||
{ signal: abortSignal },
|
{ signal: abortSignal },
|
||||||
);
|
);
|
||||||
|
if (abortSignal.aborted) {
|
||||||
|
throw createAbortError();
|
||||||
|
}
|
||||||
|
|
||||||
// Get the new response text
|
// Get the new response text
|
||||||
const retryMessage = (session.agent.session.state.messages as AgentMessage[])
|
const retryMessage = (session.agent.session.state.messages as AgentMessage[])
|
||||||
@@ -2238,6 +2284,9 @@ async function continueAgentConversation(session: Session, message: string): Pro
|
|||||||
responseText = retryText;
|
responseText = retryText;
|
||||||
markPlanningGenerationProgress(session.id, responseText);
|
markPlanningGenerationProgress(session.id, responseText);
|
||||||
} catch (retryErr) {
|
} catch (retryErr) {
|
||||||
|
if (retryErr instanceof Error && retryErr.name === "AbortError") {
|
||||||
|
throw retryErr;
|
||||||
|
}
|
||||||
// Retry prompt itself failed — give up
|
// Retry prompt itself failed — give up
|
||||||
diagnostics.errorFromException(
|
diagnostics.errorFromException(
|
||||||
"Retry prompt failed for session",
|
"Retry prompt failed for session",
|
||||||
@@ -2738,20 +2787,11 @@ export function stopGeneration(sessionId: string): boolean {
|
|||||||
}
|
}
|
||||||
|
|
||||||
activeGeneration.abortReason = "user-stop";
|
activeGeneration.abortReason = "user-stop";
|
||||||
activeGeneration.abortController.abort();
|
|
||||||
clearTimeout(activeGeneration.timer);
|
clearTimeout(activeGeneration.timer);
|
||||||
|
activeGeneration.abortTeardown();
|
||||||
|
activeGeneration.abortController.abort();
|
||||||
activeGenerations.delete(sessionId);
|
activeGenerations.delete(sessionId);
|
||||||
|
|
||||||
if (session.agent) {
|
|
||||||
nonfatal(
|
|
||||||
() => session.agent?.session.dispose?.(),
|
|
||||||
diagnostics,
|
|
||||||
"Error disposing agent for stop-generation",
|
|
||||||
{ sessionId, operation: "stop-generation-dispose" },
|
|
||||||
);
|
|
||||||
session.agent = undefined;
|
|
||||||
}
|
|
||||||
|
|
||||||
setSessionError(session, PLANNING_USER_STOP_ERROR_MESSAGE);
|
setSessionError(session, PLANNING_USER_STOP_ERROR_MESSAGE);
|
||||||
return true;
|
return true;
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -517,19 +517,18 @@ async function generateSubtasks(
|
|||||||
GENERATION_TIMEOUT_MS,
|
GENERATION_TIMEOUT_MS,
|
||||||
{
|
{
|
||||||
onTimeout: () => {
|
onTimeout: () => {
|
||||||
disposeSubtaskAgentForRetry(session);
|
|
||||||
setSubtaskError(
|
setSubtaskError(
|
||||||
sessionId,
|
sessionId,
|
||||||
"AI generation timed out. You can retry or start a new session.",
|
"AI generation timed out. You can retry or start a new session.",
|
||||||
);
|
);
|
||||||
},
|
},
|
||||||
onUserStop: () => {
|
onUserStop: () => {
|
||||||
disposeSubtaskAgentForRetry(session);
|
|
||||||
setSubtaskError(
|
setSubtaskError(
|
||||||
sessionId,
|
sessionId,
|
||||||
"Generation stopped by user. You can retry or start a new session.",
|
"Generation stopped by user. You can retry or start a new session.",
|
||||||
);
|
);
|
||||||
},
|
},
|
||||||
|
onAbort: () => disposeSubtaskAgentForRetry(session),
|
||||||
},
|
},
|
||||||
async (abortSignal) => {
|
async (abortSignal) => {
|
||||||
/*
|
/*
|
||||||
@@ -582,6 +581,13 @@ async function generateSubtasks(
|
|||||||
}
|
}
|
||||||
session.agent = agent;
|
session.agent = agent;
|
||||||
|
|
||||||
|
/*
|
||||||
|
FNXC:AiSessionCancellation 2026-07-13-00:00:
|
||||||
|
Subtask generation already forwarded the AbortSignal; keep the explicit pre/post abort checks and pair them with guard-level session teardown because Promise.race alone cannot cancel agent.session.prompt().
|
||||||
|
*/
|
||||||
|
if (abortSignal.aborted) {
|
||||||
|
throw createAbortError();
|
||||||
|
}
|
||||||
await agent.session.prompt(session.initialDescription, { signal: abortSignal });
|
await agent.session.prompt(session.initialDescription, { signal: abortSignal });
|
||||||
|
|
||||||
if (abortSignal.aborted) {
|
if (abortSignal.aborted) {
|
||||||
|
|||||||
Reference in New Issue
Block a user