FN-170: wire native chat interrupts into cancellation
Chat cancellation now uses the runtime-native interrupt seam before teardown while preserving fail-soft queue behavior. - Interrupt active runtimes before disposal for Stop and Force send. - Bound hanging or rejecting interrupts and retain dispose-only compatibility. - Add regression coverage and document the cancellation contract. Files changed: .changeset/fn-170-chat-native-interrupt.md | 7 + docs/dashboard-guide.md | 2 +- .../dashboard/src/__tests__/chat-manager.test.ts | 171 +++++++++++++++++++++ packages/dashboard/src/chat.ts | 55 ++++++- 4 files changed, 230 insertions(+), 5 deletions(-) Fusion-Task-Id: FN-170 Fusion-Task-Lineage: e81c131d-d6e5-498f-888d-f53a9a2c1201 Co-authored-by: Fusion <noreply@runfusion.ai>
This commit is contained in:
7
.changeset/fn-170-chat-native-interrupt.md
Normal file
7
.changeset/fn-170-chat-native-interrupt.md
Normal file
@@ -0,0 +1,7 @@
|
||||
---
|
||||
"@runfusion/fusion": patch
|
||||
---
|
||||
|
||||
summary: Make chat Stop and Force send interrupt active model turns before teardown.
|
||||
category: fix
|
||||
dev: `ChatManager.cancelGeneration` now makes a duck-typed, bounded native-interrupt request before disposal, mirroring the engine abort-then-dispose seam; `beginGeneration` remains controller-only.
|
||||
@@ -282,7 +282,7 @@ Direct Chat and task-detail Chat share one browser-local, text-only pending queu
|
||||
|
||||
Both model-loop surfaces expose the same queue controls: edit an entry, move it earlier or later, delete it, or **Force send** a selected entry. Duplicate text is selected by its position in the list, not by its content. A blank edit is rejected without deleting the queued entry, and queue controls remain named and touch-reachable on narrow screens.
|
||||
|
||||
Normal completion and **Stop** release only the FIFO front after cancellation and authoritative history reconciliation. **Force send** first fences and cancels the active stream, waits for cancellation plus reconciliation, then dispatches only the selected entry; failures keep the entry in its original queue position. Attachments are never queued. Activity task chat, Chat Rooms, and CLI-backed chat intentionally keep their separate interaction and transport contracts and do not inherit these model-loop queue controls.
|
||||
Normal completion and **Stop** release only the FIFO front after cancellation and authoritative history reconciliation. **Stop** and **Force send** first ask the running model runtime to interrupt through its own interrupt before the response is torn down; runtimes without an interrupt continue through the existing teardown. A runtime that stalls or fails to answer that request cannot delay or break cancellation, history reconciliation, or the queue's failure behavior. **Force send** then dispatches only the selected entry; failures keep the entry in its original queue position. Attachments are never queued. Activity task chat, Chat Rooms, and CLI-backed chat intentionally keep their separate interaction and transport contracts and do not inherit these model-loop queue controls.
|
||||
|
||||
## Automations
|
||||
|
||||
|
||||
@@ -3465,6 +3465,177 @@ describe("ChatManager.sendMessage", () => {
|
||||
}
|
||||
});
|
||||
|
||||
describe("native runtime interruption", () => {
|
||||
it("interrupts a streaming runtime before disposal and durably persists its visible prefix", async () => {
|
||||
let rejectPrompt: ((reason?: unknown) => void) | undefined;
|
||||
const abort = vi.fn().mockImplementation(async () => {
|
||||
rejectPrompt?.(new Error("Runtime interrupted"));
|
||||
});
|
||||
const dispose = vi.fn();
|
||||
__setCreateFnAgent(async (options: any) => ({
|
||||
session: {
|
||||
prompt: vi.fn().mockImplementation(() => {
|
||||
options.onText("Distinct interrupted prefix");
|
||||
return new Promise<void>((_resolve, reject) => {
|
||||
rejectPrompt = reject;
|
||||
});
|
||||
}),
|
||||
abort,
|
||||
dispose,
|
||||
state: { messages: [] },
|
||||
},
|
||||
}));
|
||||
mockChatStore.addMessage.mockImplementation(async (_sessionId: string, input: any) => ({
|
||||
id: input.role === "assistant" ? "assistant-interrupted-1" : "user-1",
|
||||
sessionId: "chat-001",
|
||||
role: input.role,
|
||||
content: input.content,
|
||||
metadata: input.metadata ?? null,
|
||||
createdAt: "2026-08-23T00:00:00.000Z",
|
||||
}));
|
||||
|
||||
const events: Array<{ type: string; data: any }> = [];
|
||||
const unsubscribe = chatStreamManager.subscribe("chat-001", (event) => events.push(event));
|
||||
const chatManager = createChatManager();
|
||||
const sendPromise = chatManager.sendMessage("chat-001", "Hello");
|
||||
await new Promise((resolve) => setTimeout(resolve, 0));
|
||||
|
||||
let sentinel: ReturnType<typeof setTimeout> | undefined;
|
||||
try {
|
||||
const cancellation = await Promise.race([
|
||||
chatManager.cancelGeneration("chat-001"),
|
||||
new Promise<never>((_resolve, reject) => {
|
||||
sentinel = setTimeout(() => reject(new Error("Cancellation did not settle")), 250);
|
||||
}),
|
||||
]);
|
||||
await sendPromise;
|
||||
|
||||
expect(cancellation).toEqual({
|
||||
success: true,
|
||||
interrupted: true,
|
||||
message: expect.objectContaining({ content: "Distinct interrupted prefix" }),
|
||||
});
|
||||
expect(abort).toHaveBeenCalledTimes(1);
|
||||
expect(abort.mock.invocationCallOrder[0]).toBeLessThan(dispose.mock.invocationCallOrder[0]!);
|
||||
const assistantCalls = mockChatStore.addMessage.mock.calls.filter((call) => call[1].role === "assistant");
|
||||
expect(assistantCalls).toHaveLength(1);
|
||||
expect(assistantCalls[0][1]).toEqual(expect.objectContaining({
|
||||
content: "Distinct interrupted prefix",
|
||||
metadata: expect.objectContaining({ interrupted: true }),
|
||||
}));
|
||||
const doneEvents = events.filter((event) => event.type === "done");
|
||||
expect(doneEvents).toHaveLength(1);
|
||||
expect(doneEvents[0].data).toEqual(expect.objectContaining({ interrupted: true }));
|
||||
} finally {
|
||||
if (sentinel !== undefined) clearTimeout(sentinel);
|
||||
unsubscribe();
|
||||
}
|
||||
});
|
||||
|
||||
it("requests the native interrupt before disposing a pre-seeded generation", async () => {
|
||||
const chatManager = createChatManager();
|
||||
const abortController = new AbortController();
|
||||
const abort = vi.fn().mockResolvedValue(undefined);
|
||||
const dispose = vi.fn();
|
||||
(chatManager as any).activeGenerations.set("chat-001", {
|
||||
abortController,
|
||||
agentResult: { session: { abort, dispose } },
|
||||
generationId: 1,
|
||||
cancellationRequested: false,
|
||||
});
|
||||
|
||||
await expect(chatManager.cancelGeneration("chat-001")).resolves.toEqual({ success: true, interrupted: false });
|
||||
expect(abortController.signal.aborted).toBe(true);
|
||||
expect(abort).toHaveBeenCalledTimes(1);
|
||||
expect(dispose).toHaveBeenCalledTimes(1);
|
||||
expect(abort.mock.invocationCallOrder[0]).toBeLessThan(dispose.mock.invocationCallOrder[0]!);
|
||||
});
|
||||
|
||||
it("keeps dispose-only cancellation for ACP/Grok-style sessions without abort", async () => {
|
||||
// FNXC:ChatCancellation 2026-08-23-02:53: ACP/Grok adapters expose dispose but no abort, so Force send must retain this fail-soft path.
|
||||
const chatManager = createChatManager();
|
||||
const abortController = new AbortController();
|
||||
const dispose = vi.fn().mockImplementation(() => {
|
||||
throw new Error("dispose-only runtime");
|
||||
});
|
||||
(chatManager as any).activeGenerations.set("chat-001", {
|
||||
abortController,
|
||||
agentResult: { session: { dispose } },
|
||||
generationId: 1,
|
||||
cancellationRequested: false,
|
||||
});
|
||||
|
||||
await expect(chatManager.cancelGeneration("chat-001")).resolves.toEqual({ success: true, interrupted: false });
|
||||
expect(dispose).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
it("bounds a hanging native interrupt before disposal", async () => {
|
||||
vi.useFakeTimers();
|
||||
const chatManager = createChatManager();
|
||||
const dispose = vi.fn();
|
||||
(chatManager as any).activeGenerations.set("chat-001", {
|
||||
abortController: new AbortController(),
|
||||
agentResult: { session: { abort: vi.fn(() => new Promise<void>(() => {})), dispose } },
|
||||
generationId: 1,
|
||||
cancellationRequested: false,
|
||||
});
|
||||
|
||||
const cancellation = chatManager.cancelGeneration("chat-001");
|
||||
await vi.advanceTimersByTimeAsync(5_000);
|
||||
await expect(cancellation).resolves.toEqual({ success: true, interrupted: false });
|
||||
expect(dispose).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
it("logs a rejecting native interrupt and still disposes", async () => {
|
||||
const error = vi.fn();
|
||||
__setChatDiagnostics({ log: vi.fn(), warn: vi.fn(), error });
|
||||
const chatManager = createChatManager();
|
||||
const dispose = vi.fn();
|
||||
(chatManager as any).activeGenerations.set("chat-001", {
|
||||
abortController: new AbortController(),
|
||||
agentResult: { session: { abort: vi.fn().mockRejectedValue(new Error("abort failed")), dispose } },
|
||||
generationId: 1,
|
||||
cancellationRequested: false,
|
||||
});
|
||||
|
||||
await expect(chatManager.cancelGeneration("chat-001")).resolves.toEqual({ success: true, interrupted: false });
|
||||
expect(dispose).toHaveBeenCalledTimes(1);
|
||||
expect(error).toHaveBeenCalledWith(
|
||||
"Failed to request runtime session interrupt during chat cancellation:",
|
||||
expect.any(Error),
|
||||
);
|
||||
});
|
||||
|
||||
it("requests the native interrupt at most once across duplicate cancellation", async () => {
|
||||
const chatManager = createChatManager();
|
||||
const abort = vi.fn().mockResolvedValue(undefined);
|
||||
(chatManager as any).activeGenerations.set("chat-001", {
|
||||
abortController: new AbortController(),
|
||||
agentResult: { session: { abort, dispose: vi.fn() } },
|
||||
generationId: 1,
|
||||
cancellationRequested: false,
|
||||
});
|
||||
|
||||
await chatManager.cancelGeneration("chat-001");
|
||||
await chatManager.cancelGeneration("chat-001");
|
||||
expect(abort).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
it("does not natively interrupt or dispose a pre-empted generation", () => {
|
||||
const chatManager = createChatManager();
|
||||
const previous = chatManager.beginGeneration("chat-001");
|
||||
const abort = vi.fn();
|
||||
const dispose = vi.fn();
|
||||
(chatManager as any).activeGenerations.get("chat-001").agentResult = { session: { abort, dispose } };
|
||||
|
||||
chatManager.beginGeneration("chat-001");
|
||||
|
||||
expect(previous.abortController.signal.aborted).toBe(true);
|
||||
expect(abort).not.toHaveBeenCalled();
|
||||
expect(dispose).not.toHaveBeenCalled();
|
||||
});
|
||||
});
|
||||
|
||||
it("treats an idle cancellation as a successful no-op without durable side effects", async () => {
|
||||
const chatManager = createChatManager();
|
||||
const events: Array<{ type: string; data: unknown }> = [];
|
||||
|
||||
@@ -172,6 +172,44 @@ const diagnostics: DiagnosticsLogger = {
|
||||
};
|
||||
|
||||
const SKILL_COMMAND_PATTERN = /(^|\s)\/skill:([^\s]+)/gi;
|
||||
const CHAT_RUNTIME_INTERRUPT_TIMEOUT_MS = 5_000;
|
||||
|
||||
/**
|
||||
* FNXC:ChatCancellation 2026-08-23-02:53:
|
||||
* Force send and Stop must ask a running runtime to interrupt through its own API,
|
||||
* not only a local AbortController that the runtime never receives. Plugin sessions
|
||||
* such as ACP and Grok expose no interrupt, so this guard is duck-typed; a stalled or
|
||||
* rejecting runtime is bounded so existing disposal and durable reconciliation proceed.
|
||||
*/
|
||||
async function requestRuntimeSessionInterrupt(session: unknown): Promise<void> {
|
||||
const sessionWithAbort = session as { abort?: () => Promise<void> | void } | undefined;
|
||||
if (typeof sessionWithAbort?.abort !== "function") {
|
||||
return;
|
||||
}
|
||||
|
||||
let timeout: ReturnType<typeof setTimeout> | undefined;
|
||||
try {
|
||||
const outcome = await Promise.race([
|
||||
Promise.resolve()
|
||||
.then(() => sessionWithAbort.abort!())
|
||||
.then(() => "completed" as const, (err) => ({ error: err })),
|
||||
new Promise<"timed-out">((resolve) => {
|
||||
timeout = setTimeout(() => resolve("timed-out"), CHAT_RUNTIME_INTERRUPT_TIMEOUT_MS);
|
||||
}),
|
||||
]);
|
||||
if (outcome === "timed-out") {
|
||||
diagnostics.error("Timed out requesting runtime session interrupt during chat cancellation");
|
||||
} else if (typeof outcome === "object" && "error" in outcome) {
|
||||
diagnostics.error("Failed to request runtime session interrupt during chat cancellation:", outcome.error);
|
||||
}
|
||||
} catch (err) {
|
||||
diagnostics.error("Failed to request runtime session interrupt during chat cancellation:", err);
|
||||
} finally {
|
||||
if (timeout !== undefined) {
|
||||
clearTimeout(timeout);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
function bareChatSkillCommandName(name: string): string {
|
||||
return name
|
||||
@@ -1815,10 +1853,10 @@ export class ChatManager {
|
||||
// controller so it stops issuing further prompts/tool calls that would
|
||||
// race against the new generation for the same CLI session file.
|
||||
//
|
||||
// We deliberately do NOT dispose its agent here — the previous generation
|
||||
// owns its own dispose in its `finally`. Calling dispose pre-emptively can
|
||||
// yank the underlying CLI process out from under the new generation's
|
||||
// freshly-opened SessionManager pointing at the same session file.
|
||||
// We deliberately do NOT request the runtime-native interrupt or dispose its agent here —
|
||||
// the previous generation owns both in its own teardown. Calling dispose pre-emptively can
|
||||
// yank the underlying CLI process out from under the new generation's freshly-opened
|
||||
// SessionManager pointing at the same session file.
|
||||
const existing = this.activeGenerations.get(sessionId);
|
||||
if (existing) {
|
||||
existing.abortController.abort();
|
||||
@@ -3421,6 +3459,15 @@ export class ChatManager {
|
||||
entry.abortController.abort();
|
||||
|
||||
if (entry.agentResult) {
|
||||
/*
|
||||
* FNXC:ChatCancellation 2026-08-23-02:53:
|
||||
* Force send and Stop must request the runtime-native interrupt before disposal because
|
||||
* the local controller is not passed to prompt(). ACP/Grok sessions lack abort(), and
|
||||
* this bounded request must finish before the existing durable settled barrier can wait.
|
||||
*/
|
||||
if (typeof (entry.agentResult.session as { abort?: unknown }).abort === "function") {
|
||||
await requestRuntimeSessionInterrupt(entry.agentResult.session);
|
||||
}
|
||||
try {
|
||||
entry.agentResult.session.dispose?.();
|
||||
} catch (err) {
|
||||
|
||||
Reference in New Issue
Block a user