fix(agent): prevented consuming legacy steering queue during mid-batch interrupts
- Stopped calling the consuming `getSteeringMessages` getter during mid-batch interrupt polls to prevent stranding or dropping messages before they reach the injection boundary. - Skip subsequent steering checks in the poll loop once an interrupt has already triggered. - Added a regression test to ensure legacy steering remains queued until the injection boundary when no non-consuming peek exists.
This commit is contained in:
@@ -8,6 +8,8 @@
|
||||
|
||||
### Fixed
|
||||
|
||||
- Fixed an issue where legacy steering messages were prematurely consumed and dropped during in-flight tool execution polls
|
||||
|
||||
- Fixed an issue where skipped tool results in queued messages were incorrectly treated as completed, preventing necessary retries.
|
||||
- Improved branch summaries to preserve informative tool results from abandoned branches while filtering out redundant output.
|
||||
- Fixed interruptible tool waits to properly abort on host-provided IRC interrupts in addition to user steering.
|
||||
|
||||
@@ -1671,7 +1671,6 @@ async function executeToolCalls(
|
||||
const {
|
||||
hasSteeringMessages,
|
||||
hasIrcInterrupts,
|
||||
getSteeringMessages,
|
||||
interruptMode = "immediate",
|
||||
getToolContext,
|
||||
transformToolCallArguments,
|
||||
@@ -1724,19 +1723,18 @@ async function executeToolCalls(
|
||||
const checkSteering = async (): Promise<void> => {
|
||||
// `signal` (external/user abort) is checked separately from the internal
|
||||
// abort controllers: once the run is externally aborted it is unwinding
|
||||
// and the interrupt would be redundant.
|
||||
if (!shouldInterruptImmediately || signal?.aborted) {
|
||||
// and the interrupt would be redundant. Once an IRC/steering interrupt has
|
||||
// already fired, do not poll again — especially not through a consuming
|
||||
// legacy steering getter.
|
||||
if (!shouldInterruptImmediately || interruptState.triggered || signal?.aborted) {
|
||||
return;
|
||||
}
|
||||
// Prefer non-consuming peeks so queues still own their messages until
|
||||
// the injection boundary. Fall back to consuming steering only for older
|
||||
// integrations that never supplied a peek.
|
||||
// Mid-batch steering detection must be non-consuming. If a direct
|
||||
// integration only provides getSteeringMessages(), the queue drains at the
|
||||
// injection boundary below; polling it here would strand or drop messages.
|
||||
let steeringQueued = false;
|
||||
if (hasSteeringMessages) {
|
||||
steeringQueued = await hasSteeringMessages();
|
||||
} else if (getSteeringMessages) {
|
||||
const msgs = await getSteeringMessages();
|
||||
steeringQueued = (msgs?.length ?? 0) > 0;
|
||||
}
|
||||
if (steeringQueued) {
|
||||
// User steering upgrades an in-flight IRC interrupt: it aborts the
|
||||
@@ -1747,7 +1745,6 @@ async function executeToolCalls(
|
||||
}
|
||||
return;
|
||||
}
|
||||
if (interruptState.triggered) return;
|
||||
if (hasIrcInterrupts && (await hasIrcInterrupts())) {
|
||||
// Peer IRC only aborts interruptible waits: a foreground bash / write
|
||||
// mid-execution keeps running so we never leave partial side effects.
|
||||
@@ -2078,13 +2075,13 @@ async function executeToolCalls(
|
||||
|
||||
// While an interruptible tool is in flight (e.g. a `job`/`irc` wait
|
||||
// blocking on external work), queued steering or interrupting IRC would
|
||||
// otherwise wait out the tool's own window. Poll the non-consuming queues
|
||||
// otherwise wait out the tool's own window. Poll only non-consuming queues
|
||||
// and abort the shared tool signal so the boundary dequeue below injects
|
||||
// the message promptly. Gated on immediate-interrupt mode + an
|
||||
// interruptible tool; checkSteering is idempotent (no-op once triggered).
|
||||
const watchSteeringWhileRunning =
|
||||
shouldInterruptImmediately &&
|
||||
(hasSteeringMessages !== undefined || getSteeringMessages !== undefined || hasIrcInterrupts !== undefined) &&
|
||||
(hasSteeringMessages !== undefined || hasIrcInterrupts !== undefined) &&
|
||||
records.some(r => r.tool?.interruptible === true);
|
||||
const steeringWatchTimer = watchSteeringWhileRunning
|
||||
? setInterval(() => void checkSteering(), STEERING_INTERRUPT_POLL_MS)
|
||||
|
||||
@@ -1258,6 +1258,75 @@ describe("agentLoop with AgentMessage", () => {
|
||||
).toBe(true);
|
||||
});
|
||||
|
||||
it("keeps legacy steering queued until the injection boundary when no non-consuming peek exists", async () => {
|
||||
const toolSchema = type({ value: "string" });
|
||||
const executed: string[] = [];
|
||||
let steerReady = false;
|
||||
let steeringDrained = false;
|
||||
const steeringMessage = createUserMessage("queued steering");
|
||||
|
||||
const tool: AgentTool<typeof toolSchema, { value: string }> = {
|
||||
name: "echo",
|
||||
label: "Echo",
|
||||
description: "Echo tool",
|
||||
parameters: toolSchema,
|
||||
concurrency: "exclusive",
|
||||
async execute(_toolCallId, params) {
|
||||
executed.push(params.value);
|
||||
if (params.value === "first") steerReady = true;
|
||||
return {
|
||||
content: [{ type: "text", text: `ok:${params.value}` }],
|
||||
details: { value: params.value },
|
||||
};
|
||||
},
|
||||
};
|
||||
|
||||
const context: AgentContext = { systemPrompt: [""], messages: [], tools: [tool] };
|
||||
const mock = createMockModel({
|
||||
responses: [
|
||||
{
|
||||
content: [
|
||||
{ type: "toolCall", id: "tool-1", name: "echo", arguments: { value: "first" } },
|
||||
{ type: "toolCall", id: "tool-2", name: "echo", arguments: { value: "second" } },
|
||||
],
|
||||
},
|
||||
{ content: ["done"] },
|
||||
],
|
||||
});
|
||||
const config: AgentLoopConfig = {
|
||||
model: mock.model,
|
||||
convertToLlm: identityConverter,
|
||||
interruptMode: "immediate",
|
||||
// Legacy integrations may only have the consuming dequeue. The
|
||||
// mid-batch poll must not call it; the boundary below owns the drain.
|
||||
getSteeringMessages: async () => {
|
||||
if (steerReady && !steeringDrained) {
|
||||
steeringDrained = true;
|
||||
return [steeringMessage];
|
||||
}
|
||||
return [];
|
||||
},
|
||||
};
|
||||
|
||||
const events: AgentEvent[] = [];
|
||||
for await (const event of agentLoop([createUserMessage("start")], context, config, undefined, mock.stream)) {
|
||||
events.push(event);
|
||||
}
|
||||
|
||||
expect(executed).toEqual(["first", "second"]);
|
||||
expect(steeringDrained).toBe(true);
|
||||
expect(
|
||||
events.some(
|
||||
e => e.type === "message_start" && e.message.role === "user" && e.message.content === "queued steering",
|
||||
),
|
||||
).toBe(true);
|
||||
expect(
|
||||
mock.calls[1]?.context.messages.some(
|
||||
m => m.role === "user" && typeof m.content === "string" && m.content === "queued steering",
|
||||
),
|
||||
).toBe(true);
|
||||
});
|
||||
|
||||
it("does not abort a non-interruptible foreground tool when only IRC is queued", async () => {
|
||||
const toolSchema = type({});
|
||||
let ircReady = false;
|
||||
|
||||
@@ -137,7 +137,6 @@ describe("executeBash", () => {
|
||||
expect(result.output.trim()).toBe(tempDir);
|
||||
});
|
||||
|
||||
|
||||
it("honors symlinked cwd requests in persistent shells", async () => {
|
||||
if (process.platform === "win32") {
|
||||
return;
|
||||
|
||||
Reference in New Issue
Block a user