fix(coding-agent): coalesced queued steer flushes to avoid AgentBusyError interrupt races

- Coalesced repeated interrupt-and-flush calls into one queued-steer resume flow.
- Retried queued continues on AgentBusyError after waitForIdle up to a 30s timeout.
- Handled queue-flush failures in input controllers with warning logs and TUI error display.
This commit is contained in:
can1357
2026-06-13 15:52:25 +02:00
parent 69051241f1
commit 39efdfe91b
5 changed files with 211 additions and 6 deletions
+1
View File
@@ -21,6 +21,7 @@
### Fixed
- Fixed empty-Enter steering injection so repeated submits coalesce into one interrupt/resume cycle instead of racing `agent.continue()` and surfacing `AgentBusyError`; failed flushes are now reported in the TUI instead of becoming unhandled rejections.
- Fixed model auth gateway probing to avoid skipping candidates with unknown `maxTokens` limits (`null`)
- Fixed model listings so providers registered via extensions are now included from `-e` and configured `extensions` sources
- Fixed `/mcp reauth`, `/mcp test`, and `/mcp unauth` to find and operate on MCP servers reported by `/mcp list` even when they are only runtime-discovered and not stored in writable config, including namespaced plugin servers like `cloudflare:cloudflare-api`
@@ -412,7 +412,14 @@ export class InputController {
const queuedMessages = this.ctx.session.getQueuedMessages();
if (queuedMessages.steering.length > 0) {
this.ctx.notifyInterrupting();
await this.ctx.session.interruptAndFlushQueuedMessages({ reason: USER_INTERRUPT_LABEL });
try {
await this.ctx.session.interruptAndFlushQueuedMessages({ reason: USER_INTERRUPT_LABEL });
} catch (error) {
logger.warn("queued steer interrupt failed", {
error: error instanceof Error ? error.message : String(error),
});
this.ctx.showError(error instanceof Error ? error.message : String(error));
}
this.ctx.updatePendingMessagesDisplay();
this.ctx.ui.requestRender();
return;
@@ -696,7 +703,14 @@ export class InputController {
if (!text) {
// Mirror the empty-submit steer flush against the focused session.
if (target.isStreaming && target.getQueuedMessages().steering.length > 0) {
await target.interruptAndFlushQueuedMessages({ reason: USER_INTERRUPT_LABEL });
try {
await target.interruptAndFlushQueuedMessages({ reason: USER_INTERRUPT_LABEL });
} catch (error) {
logger.warn("focused queued steer interrupt failed", {
error: error instanceof Error ? error.message : String(error),
});
this.ctx.showError(error instanceof Error ? error.message : String(error));
}
this.ctx.updatePendingMessagesDisplay();
this.ctx.ui.requestRender();
}
@@ -1084,6 +1084,7 @@ export class AgentSession {
// both #endInFlight (normal) and #resetInFlight (abort).
#pendingAgentEndEmit: AgentSessionEvent | undefined;
#resumingQueuedMessages = false;
#queuedMessageFlushStartPromise: Promise<{ continuation: Promise<void> | undefined }> | undefined;
#obfuscator: SecretObfuscator | undefined;
#checkpointState: CheckpointState | undefined = undefined;
#pendingRewindReport: string | undefined = undefined;
@@ -5510,20 +5511,61 @@ export class AgentSession {
* messages drain instead of waiting for another natural turn boundary.
*/
async interruptAndFlushQueuedMessages(options?: { reason?: string }): Promise<void> {
const existingStart = this.#queuedMessageFlushStartPromise;
if (existingStart) {
const { continuation } = await existingStart;
await continuation;
return;
}
if (!this.agent.hasQueuedMessages()) return;
const start = this.#beginQueuedMessageFlush(options);
this.#queuedMessageFlushStartPromise = start;
let continuation: Promise<void> | undefined;
try {
({ continuation } = await start);
} finally {
if (this.#queuedMessageFlushStartPromise === start) {
this.#queuedMessageFlushStartPromise = undefined;
}
}
await continuation;
}
async #beginQueuedMessageFlush(options?: { reason?: string }): Promise<{ continuation: Promise<void> | undefined }> {
if (!this.agent.hasQueuedMessages()) return { continuation: undefined };
this.#resumingQueuedMessages = true;
try {
await this.abort({ reason: options?.reason });
if (!this.agent.hasQueuedMessages()) return;
if (this.isCompacting || this.isGeneratingHandoff) return;
if (!this.agent.hasQueuedMessages()) return { continuation: undefined };
await this.#maybeRestoreRetryFallbackPrimary();
this.#resumingQueuedMessages = false;
await this.agent.continue();
if (!this.agent.hasQueuedMessages()) return { continuation: undefined };
const continuation = this.#continueQueuedMessagesWithIdleRetry();
return { continuation };
} finally {
this.#resumingQueuedMessages = false;
}
}
async #continueQueuedMessagesWithIdleRetry(): Promise<void> {
const deadline = Date.now() + 30_000;
for (;;) {
if (!this.agent.hasQueuedMessages()) return;
try {
await this.agent.continue();
return;
} catch (err) {
if (!(err instanceof AgentBusyError)) {
throw err;
}
if (Date.now() >= deadline) {
throw new Error("Timed out waiting for prior agent run to finish before continuing queued messages.");
}
await this.agent.waitForIdle();
}
}
}
/**
* Start a new session, optionally with initial messages and parent tracking.
* Clears all messages and starts a new session.
@@ -226,6 +226,152 @@ describe("AgentSession concurrent prompt guard", () => {
).toBe(true);
expect(session.getQueuedMessages()).toEqual({ steering: [], followUp: [] });
});
it("coalesces repeated interrupt-and-flush requests for one queued steer", async () => {
const model = getBundledModel("anthropic", "claude-sonnet-4-5")!;
const callMessages: Message[][] = [];
const agent = new Agent({
getApiKey: () => "test-key",
initialState: {
model,
systemPrompt: ["Test"],
tools: [],
},
convertToLlm,
streamFn: (_model, context, options) => {
const callIndex = callMessages.length;
callMessages.push([...context.messages]);
const stream = new AssistantMessageEventStream();
queueMicrotask(() => {
stream.push({ type: "start", partial: createAssistantMessage("") });
if (callIndex > 0) {
stream.push({ type: "done", reason: "stop", message: createAssistantMessage("Handled steer") });
}
});
options?.signal?.addEventListener(
"abort",
() => {
stream.push({
type: "error",
reason: "aborted",
error: createAssistantMessage("Interrupted"),
});
},
{ once: true },
);
return stream;
},
});
const sessionManager = SessionManager.inMemory();
const settings = Settings.isolated();
const authStorage = await AuthStorage.create(path.join(tempDir, "testauth-interrupt-flush-repeat.db"));
authStorages.push(authStorage);
const modelRegistry = new ModelRegistry(authStorage, path.join(tempDir, "models-interrupt-flush-repeat.yml"));
authStorage.setRuntimeApiKey("anthropic", "test-key");
session = new AgentSession({
agent,
sessionManager,
settings,
modelRegistry,
});
const firstPrompt = session.prompt("First message").catch(() => {});
await waitFor(() => session.isStreaming && callMessages.length === 1);
await session.steer("Send this once");
expect(session.getQueuedMessages().steering).toEqual(["Send this once"]);
const flushes = Array.from({ length: 6 }, () =>
session.interruptAndFlushQueuedMessages({ reason: "Interrupted by user" }),
);
await expect(Promise.all(flushes)).resolves.toEqual([
undefined,
undefined,
undefined,
undefined,
undefined,
undefined,
]);
await firstPrompt;
expect(callMessages).toHaveLength(2);
const resumedCall = callMessages[1];
expect(
resumedCall?.some(message => {
if (typeof message.content === "string") {
return message.content.includes("Send this once");
}
return message.content.some(content => content.type === "text" && content.text.includes("Send this once"));
}),
).toBe(true);
expect(session.getQueuedMessages()).toEqual({ steering: [], followUp: [] });
});
it("does not resume when the queued steer is cleared during interrupt-and-flush", async () => {
const model = getBundledModel("anthropic", "claude-sonnet-4-5")!;
const callMessages: Message[][] = [];
const agent = new Agent({
getApiKey: () => "test-key",
initialState: {
model,
systemPrompt: ["Test"],
tools: [],
},
convertToLlm,
streamFn: (_model, context, options) => {
callMessages.push([...context.messages]);
const stream = new AssistantMessageEventStream();
queueMicrotask(() => {
stream.push({ type: "start", partial: createAssistantMessage("") });
});
options?.signal?.addEventListener(
"abort",
() => {
stream.push({
type: "error",
reason: "aborted",
error: createAssistantMessage("Interrupted"),
});
},
{ once: true },
);
return stream;
},
});
const sessionManager = SessionManager.inMemory();
const settings = Settings.isolated();
const authStorage = await AuthStorage.create(path.join(tempDir, "testauth-interrupt-flush-clear.db"));
authStorages.push(authStorage);
const modelRegistry = new ModelRegistry(authStorage, path.join(tempDir, "models-interrupt-flush-clear.yml"));
authStorage.setRuntimeApiKey("anthropic", "test-key");
session = new AgentSession({
agent,
sessionManager,
settings,
modelRegistry,
});
const firstPrompt = session.prompt("First message").catch(() => {});
await waitFor(() => session.isStreaming && callMessages.length === 1);
await session.steer("Restore this to the editor instead");
const abort = session.abort.bind(session);
vi.spyOn(session, "abort").mockImplementation(async options => {
await abort(options);
session.clearQueue();
});
await expect(session.interruptAndFlushQueuedMessages({ reason: "Interrupted by user" })).resolves.toBeUndefined();
await firstPrompt;
expect(callMessages).toHaveLength(1);
expect(session.getQueuedMessages()).toEqual({ steering: [], followUp: [] });
});
it("delivers queued steering after interrupting mid-tool execution (queue survives external abort)", async () => {
// Regression: pressing Enter with a queued steer while a tool was running
@@ -149,6 +149,8 @@ async function createContext() {
showModelSelector,
updateEditorBorderColor: vi.fn(),
hasActiveBtw: vi.fn(() => false),
notifyInterrupting: vi.fn(),
showError: vi.fn(),
} as unknown as InteractiveModeContext;
return {