diff --git a/packages/agent/src/agent-loop.ts b/packages/agent/src/agent-loop.ts index 42a4b67e1..49549dd41 100644 --- a/packages/agent/src/agent-loop.ts +++ b/packages/agent/src/agent-loop.ts @@ -2872,9 +2872,12 @@ function createSkippedToolResult( if (source === "user") { reason = "queued user message"; blocker = "queued message"; - } else if (source === "system") { - reason = "pending internal steering message"; + } else if (source === "agent") { + reason = "pending parent steering message"; blocker = "steering message"; + } else if (source === "system") { + reason = "pending system advisory"; + blocker = "advisory"; } else if (source === "irc") { reason = "pending peer interrupt"; blocker = "interrupt"; diff --git a/packages/agent/src/agent.ts b/packages/agent/src/agent.ts index 7e6216c3a..8b5f36885 100644 --- a/packages/agent/src/agent.ts +++ b/packages/agent/src/agent.ts @@ -1376,14 +1376,17 @@ export class Agent { if (this.#steeringQueue.length === 0) { return { queued: false }; } + let hasAgentSteering = false; for (const message of this.#steeringQueue) { const role = "role" in message ? message.role : undefined; const attribution = "attribution" in message ? message.attribution : undefined; - if (role === "user" && attribution !== "agent") { + if (role !== "user") continue; + if (attribution !== "agent") { return { queued: true, source: "user" }; } + hasAgentSteering = true; } - return { queued: true, source: "system" }; + return { queued: true, source: hasAgentSteering ? "agent" : "system" }; }, waitForSteeringMessages: signal => this.#waitForSteeringMessages(signal), hasIrcInterrupts: this.hasIrcInterrupts, diff --git a/packages/agent/src/types.ts b/packages/agent/src/types.ts index 51c3bbe1d..fbe7e6c43 100644 --- a/packages/agent/src/types.ts +++ b/packages/agent/src/types.ts @@ -117,8 +117,11 @@ export function isSoftToolRequirement(directive: ToolChoiceDirective | undefined return typeof directive === "object" && directive !== null && (directive as SoftToolRequirement).soft === true; } -/** Source category for a queued steering interrupt observed without consuming the queue. */ -export type SteeringInterruptSource = "user" | "system" | "unknown"; +/** + * Source category for a queued steering interrupt observed without consuming the queue. + * Distinguishes real-user, agent-authored, system/advisor, and unknown steering. + */ +export type SteeringInterruptSource = "user" | "agent" | "system" | "unknown"; /** Non-consuming summary of whether queued steering should interrupt a tool batch. */ export interface SteeringQueueState { diff --git a/packages/agent/test/agent-loop.test.ts b/packages/agent/test/agent-loop.test.ts index 496e7be47..c6e574638 100644 --- a/packages/agent/test/agent-loop.test.ts +++ b/packages/agent/test/agent-loop.test.ts @@ -1532,7 +1532,7 @@ describe("agentLoop with AgentMessage", () => { expect(sawInterruptInContext).toBe(true); }); - it("should skip remaining tool calls with internal steering wording when non-user steering is queued", async () => { + it("should skip remaining tool calls with system advisory wording when advisor steering is queued", async () => { const toolSchema = type({ value: "string" }); const executed: string[] = []; const tool: AgentTool = { @@ -1610,9 +1610,7 @@ describe("agentLoop with AgentMessage", () => { const skippedContent = toolEnds[1].result.content[0]; expect(skippedContent?.type).toBe("text"); if (skippedContent?.type !== "text") throw new Error("skipped tool result must be text"); - expect(skippedContent.text).toContain("Skipped due to pending internal steering message"); - expect(skippedContent.text).toContain("After the steering message is handled on the next step"); - expect(skippedContent.text).not.toContain("advisory"); + expect(skippedContent.text).toContain("Skipped due to pending system advisory"); expect(skippedContent.text).not.toContain("queued user message"); expect(skippedContent.text).toContain("Do not count this skipped result as completed work"); expect(skippedContent.text).toContain("retry the skipped tool if it is still needed"); diff --git a/packages/agent/test/agent.test.ts b/packages/agent/test/agent.test.ts index 0eb82f5cd..a4fbaed10 100644 --- a/packages/agent/test/agent.test.ts +++ b/packages/agent/test/agent.test.ts @@ -17,6 +17,69 @@ describe("Agent", () => { expect(agent.state.messages).not.toContainEqual(message); }); + it("classifies agent-authored steering as a parent steering message", async () => { + const toolSchema = z.object({ value: z.string() }); + const executed: string[] = []; + let agent: Agent; + const tool: AgentTool = { + name: "echo", + label: "Echo", + description: "Echo tool", + parameters: toolSchema, + concurrency: "exclusive", + async execute(_toolCallId, params) { + executed.push(params.value); + if (params.value === "first") { + agent.steer({ + role: "user", + content: "parent steering", + attribution: "agent", + timestamp: Date.now(), + }); + } + return { + content: [{ type: "text", text: `ok:${params.value}` }], + details: { value: params.value }, + }; + }, + }; + 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"] }, + ], + }); + agent = new Agent({ + initialState: { model: mock.model, systemPrompt: ["Test"], tools: [tool], messages: [] }, + streamFn: mock.stream, + interruptMode: "immediate", + }); + const events: AgentEvent[] = []; + const unsubscribe = agent.subscribe(event => events.push(event)); + + await agent.prompt("start"); + unsubscribe(); + + expect(executed).toEqual(["first"]); + const skipped = events.find( + (event): event is Extract => + event.type === "tool_execution_end" && event.toolCallId === "tool-2", + ); + expect(skipped).toBeDefined(); + const skippedContent = skipped?.result.content[0]; + expect(skippedContent?.type).toBe("text"); + if (skippedContent?.type !== "text") throw new Error("skipped tool result must be text"); + expect(skippedContent.text).toContain("Skipped due to pending parent steering message"); + expect(skippedContent.text).toContain("After the steering message is handled on the next step"); + expect(skippedContent.text).not.toContain("pending system advisory"); + expect(skippedContent.text).not.toContain("queued user message"); + }); + it("continue() should process queued follow-up messages after an assistant turn", async () => { const mock = createMockModel({ responses: [{ content: ["Processed"] }] }); const agent = new Agent({ streamFn: mock.stream });