fix(agent): classify next queued steering source

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
Fatih Al-Aziz
2026-08-03 13:18:45 +07:00
parent 1e448304ac
commit 2416134cf0
2 changed files with 93 additions and 1 deletions
+3 -1
View File
@@ -1373,8 +1373,10 @@ export class Agent {
if (this.#steeringQueue.length === 0) {
return { queued: false };
}
const messageCount = this.#steeringMode === "one-at-a-time" ? 1 : this.#steeringQueue.length;
let hasAgentSteering = false;
for (const message of this.#steeringQueue) {
for (let i = 0; i < messageCount; i++) {
const message = this.#steeringQueue[i];
const role = "role" in message ? message.role : undefined;
const attribution = "attribution" in message ? message.attribution : undefined;
if (role !== "user") continue;
+90
View File
@@ -80,6 +80,96 @@ describe("Agent", () => {
expect(skippedContent.text).not.toContain("queued user message");
});
it("classifies one-at-a-time steering from the next queued mixed source", async () => {
const cases = [
{
order: ["system", "agent"] as const,
expected: "pending system advisory",
unexpected: "pending parent steering message",
},
{
order: ["agent", "system"] as const,
expected: "pending parent steering message",
unexpected: "pending system advisory",
},
];
for (const scenario of cases) {
const toolSchema = z.object({ value: z.string() });
const executed: string[] = [];
let agent: Agent;
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") {
for (const source of scenario.order) {
if (source === "agent") {
agent.steer({
role: "user",
content: "parent steering",
attribution: "agent",
timestamp: Date.now(),
});
} else {
agent.steer({
role: "custom",
customType: "advisor",
content: "advisor steering",
display: true,
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"] },
{ 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<AgentEvent, { type: "tool_execution_end" }> =>
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 ${scenario.expected}`);
expect(skippedContent.text).not.toContain(scenario.unexpected);
}
});
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 });