feat(agent): added interruptible tool polling for queued steering
- Added an `interruptible` field to AgentTool and documented when it is honored. - Updated immediate-mode tool execution to poll steering during in-flight interruptible calls and abort them when steering is queued. - Marked the coding `job` tool as interruptible and added tests covering mid-wait aborts versus boundary-only steering drain.
This commit is contained in:
@@ -3,6 +3,7 @@
|
||||
## [Unreleased]
|
||||
### Added
|
||||
|
||||
- Added the `interruptible` tool field: when set, the agent loop may abort the tool mid-execution to deliver a queued steering message (honored only in `immediate` interrupt mode).
|
||||
- Added support for `gemini` and `gemma` as valid owned tool syntax values in environment configuration
|
||||
|
||||
## [15.13.2] - 2026-06-15
|
||||
|
||||
@@ -74,6 +74,14 @@ const ABORTED: unique symbol = Symbol("agent-loop-aborted");
|
||||
*/
|
||||
const MAX_PAUSED_TURN_CONTINUATIONS = 8;
|
||||
|
||||
/**
|
||||
* Cadence (ms) for polling queued steering while an `interruptible` tool is in
|
||||
* flight, so a steer cuts the wait short instead of sitting idle until the
|
||||
* tool's own window elapses. A cheap synchronous queue check; latency-bounded
|
||||
* at one tick.
|
||||
*/
|
||||
const STEERING_INTERRUPT_POLL_MS = 250;
|
||||
|
||||
class HarmonyLeakInterruption extends Error {
|
||||
constructor(
|
||||
readonly detection: HarmonyDetection,
|
||||
@@ -1797,7 +1805,24 @@ async function executeToolCalls(
|
||||
}
|
||||
}
|
||||
|
||||
await Promise.allSettled(tasks);
|
||||
// While an interruptible tool is in flight (e.g. a `job` poll blocking on
|
||||
// background work), a queued steer would otherwise wait out the tool's own
|
||||
// window. Poll the steering queue and let checkSteering() abort the shared
|
||||
// tool signal so the wait returns early; the boundary dequeue below then
|
||||
// injects it. Gated on immediate-interrupt mode + an interruptible tool;
|
||||
// checkSteering is idempotent (no-op once triggered).
|
||||
const watchSteeringWhileRunning =
|
||||
shouldInterruptImmediately &&
|
||||
(hasSteeringMessages !== undefined || getSteeringMessages !== undefined) &&
|
||||
records.some(r => r.tool?.interruptible === true);
|
||||
const steeringWatchTimer = watchSteeringWhileRunning
|
||||
? setInterval(() => void checkSteering(), STEERING_INTERRUPT_POLL_MS)
|
||||
: undefined;
|
||||
try {
|
||||
await Promise.allSettled(tasks);
|
||||
} finally {
|
||||
if (steeringWatchTimer !== undefined) clearInterval(steeringWatchTimer);
|
||||
}
|
||||
// Yield after batch tool execution to let GC and I/O catch up,
|
||||
// especially when tool results are large (e.g. bash output).
|
||||
await yieldIfDue();
|
||||
|
||||
@@ -503,6 +503,15 @@ export interface AgentTool<TParameters extends TSchema = TSchema, TDetails = any
|
||||
concurrency?: "shared" | "exclusive" | ((args: Partial<Static<TParameters>>) => "shared" | "exclusive");
|
||||
/** If true, argument validation errors are non-fatal: raw args are passed to execute() instead of returning an error to the LLM. */
|
||||
lenientArgValidation?: boolean;
|
||||
/**
|
||||
* If true, the agent loop may abort this tool mid-execution to deliver a
|
||||
* queued steering message (instead of waiting for the tool to finish on its
|
||||
* own). Set only on tools that purely *wait* and observe their abort signal
|
||||
* cleanly (e.g. the `job` poll), so the abort surfaces the tool's current
|
||||
* snapshot rather than corrupting a side effect. Honored only when
|
||||
* `interruptMode` is "immediate".
|
||||
*/
|
||||
interruptible?: boolean;
|
||||
/**
|
||||
* Controls how the INTENT_FIELD (`_i`) is handled for this tool.
|
||||
* - `"require"` (default): `_i` is injected and required in the parameter schema.
|
||||
|
||||
@@ -842,6 +842,149 @@ describe("agentLoop with AgentMessage", () => {
|
||||
expect(sawInterruptInContext).toBe(true);
|
||||
});
|
||||
|
||||
it("drains queued steering by aborting an interruptible tool mid-wait", async () => {
|
||||
const toolSchema = z.object({});
|
||||
let steerReady = false;
|
||||
let drained = false;
|
||||
let observedAbort = false;
|
||||
let resolvedByTimeout = false;
|
||||
|
||||
const tool: AgentTool<typeof toolSchema, Record<string, never>> = {
|
||||
name: "wait",
|
||||
label: "Wait",
|
||||
description: "Blocks until aborted (mimics a job poll)",
|
||||
parameters: toolSchema,
|
||||
interruptible: true,
|
||||
async execute(_toolCallId, _params, signal) {
|
||||
steerReady = true;
|
||||
const { promise, resolve } = Promise.withResolvers<void>();
|
||||
if (signal?.aborted) {
|
||||
resolve();
|
||||
} else {
|
||||
const timer = setTimeout(() => {
|
||||
resolvedByTimeout = true;
|
||||
resolve();
|
||||
}, 2000);
|
||||
signal?.addEventListener(
|
||||
"abort",
|
||||
() => {
|
||||
clearTimeout(timer);
|
||||
resolve();
|
||||
},
|
||||
{ once: true },
|
||||
);
|
||||
}
|
||||
await promise;
|
||||
observedAbort = signal?.aborted === true;
|
||||
return { content: [{ type: "text", text: "waited" }], details: {} };
|
||||
},
|
||||
};
|
||||
|
||||
const context: AgentContext = { systemPrompt: [""], messages: [], tools: [tool] };
|
||||
const mock = createMockModel({
|
||||
responses: [
|
||||
{ content: [{ type: "toolCall", id: "tool-1", name: "wait", arguments: {} }] },
|
||||
{ content: ["done"] },
|
||||
],
|
||||
});
|
||||
const config: AgentLoopConfig = {
|
||||
model: mock.model,
|
||||
convertToLlm: identityConverter,
|
||||
interruptMode: "immediate",
|
||||
hasSteeringMessages: () => steerReady && !drained,
|
||||
getSteeringMessages: async () => {
|
||||
if (steerReady && !drained) {
|
||||
drained = true;
|
||||
return [createUserMessage("interrupt")];
|
||||
}
|
||||
return [];
|
||||
},
|
||||
};
|
||||
|
||||
const events: AgentEvent[] = [];
|
||||
for await (const event of agentLoop([createUserMessage("start")], context, config, undefined, mock.stream)) {
|
||||
events.push(event);
|
||||
}
|
||||
|
||||
expect(observedAbort).toBe(true);
|
||||
expect(resolvedByTimeout).toBe(false);
|
||||
expect(drained).toBe(true);
|
||||
expect(
|
||||
events.some(e => e.type === "message_start" && e.message.role === "user" && e.message.content === "interrupt"),
|
||||
).toBe(true);
|
||||
});
|
||||
|
||||
it("does not abort a non-interruptible tool mid-wait; steering still drains at the boundary", async () => {
|
||||
const toolSchema = z.object({});
|
||||
let steerReady = false;
|
||||
let drained = false;
|
||||
let observedAbort = false;
|
||||
let resolvedByTimeout = false;
|
||||
|
||||
const tool: AgentTool<typeof toolSchema, Record<string, never>> = {
|
||||
name: "wait",
|
||||
label: "Wait",
|
||||
description: "Blocks on its own window (no interruptible flag)",
|
||||
parameters: toolSchema,
|
||||
async execute(_toolCallId, _params, signal) {
|
||||
steerReady = true;
|
||||
const { promise, resolve } = Promise.withResolvers<void>();
|
||||
if (signal?.aborted) {
|
||||
resolve();
|
||||
} else {
|
||||
const timer = setTimeout(() => {
|
||||
resolvedByTimeout = true;
|
||||
resolve();
|
||||
}, 300);
|
||||
signal?.addEventListener(
|
||||
"abort",
|
||||
() => {
|
||||
clearTimeout(timer);
|
||||
resolve();
|
||||
},
|
||||
{ once: true },
|
||||
);
|
||||
}
|
||||
await promise;
|
||||
observedAbort = signal?.aborted === true;
|
||||
return { content: [{ type: "text", text: "waited" }], details: {} };
|
||||
},
|
||||
};
|
||||
|
||||
const context: AgentContext = { systemPrompt: [""], messages: [], tools: [tool] };
|
||||
const mock = createMockModel({
|
||||
responses: [
|
||||
{ content: [{ type: "toolCall", id: "tool-1", name: "wait", arguments: {} }] },
|
||||
{ content: ["done"] },
|
||||
],
|
||||
});
|
||||
const config: AgentLoopConfig = {
|
||||
model: mock.model,
|
||||
convertToLlm: identityConverter,
|
||||
interruptMode: "immediate",
|
||||
hasSteeringMessages: () => steerReady && !drained,
|
||||
getSteeringMessages: async () => {
|
||||
if (steerReady && !drained) {
|
||||
drained = true;
|
||||
return [createUserMessage("interrupt")];
|
||||
}
|
||||
return [];
|
||||
},
|
||||
};
|
||||
|
||||
const events: AgentEvent[] = [];
|
||||
for await (const event of agentLoop([createUserMessage("start")], context, config, undefined, mock.stream)) {
|
||||
events.push(event);
|
||||
}
|
||||
|
||||
expect(observedAbort).toBe(false);
|
||||
expect(resolvedByTimeout).toBe(true);
|
||||
expect(drained).toBe(true);
|
||||
expect(
|
||||
events.some(e => e.type === "message_start" && e.message.role === "user" && e.message.content === "interrupt"),
|
||||
).toBe(true);
|
||||
});
|
||||
|
||||
it("leaves steering queued when the run is aborted while interrupted tools settle", async () => {
|
||||
// Regression: the mid-batch steering poll used to DEQUEUE the message into
|
||||
// a loop-local variable. An external abort while the in-flight tools were
|
||||
|
||||
@@ -8,6 +8,7 @@
|
||||
|
||||
### Changed
|
||||
|
||||
- Changed the `job` poll to return early when a steering message is queued, draining the steer immediately instead of waiting out the poll window.
|
||||
- Capped unexpected-stop auto-continuation to three retry attempts before giving up on repeated stops
|
||||
- Updated the `edit` tool's hashline prompt, grammar, and docs to recommend the `.=` inclusive range separator (`SWAP 1.=3:`); the legacy `..` form still parses.
|
||||
|
||||
|
||||
@@ -87,6 +87,7 @@ export class JobTool implements AgentTool<typeof jobSchema, JobToolDetails> {
|
||||
readonly description: string;
|
||||
readonly parameters = jobSchema;
|
||||
readonly strict = true;
|
||||
readonly interruptible = true;
|
||||
readonly loadMode = "discoverable";
|
||||
constructor(private readonly session: ToolSession) {
|
||||
this.description = prompt.render(jobDescription);
|
||||
|
||||
Reference in New Issue
Block a user