diff --git a/packages/agent/src/agent-loop.ts b/packages/agent/src/agent-loop.ts index 36ba3587e..4e697ed4a 100644 --- a/packages/agent/src/agent-loop.ts +++ b/packages/agent/src/agent-loop.ts @@ -49,6 +49,7 @@ import type { AgentMessage, AgentTool, AgentToolResult, + AsideMessage, StreamFn, } from "./types"; import { yieldIfDue } from "./utils/yield"; @@ -465,6 +466,23 @@ function cloneAssistantMessageForToolCallCap(message: AssistantMessage): Assista }; } +/** + * Resolve aside entries at the moment the loop is about to inject them. Each entry + * is either a ready {@link AgentMessage} or a sync thunk evaluated here so the + * producer can make the final inject-or-drop decision (return null) against + * up-to-the-injection state — e.g. dropping late diagnostics a newer edit + * superseded. Kept sync so it can never stall the loop. + */ +function resolveAsides(entries: AsideMessage[] | undefined): AgentMessage[] { + if (!entries || entries.length === 0) return []; + const out: AgentMessage[] = []; + for (const entry of entries) { + const message = typeof entry === "function" ? entry() : entry; + if (message) out.push(message); + } + return out; +} + async function runLoopBody( currentContext: AgentContext, newMessages: AgentMessage[], @@ -648,13 +666,13 @@ async function runLoopBody( stream.push({ type: "turn_end", message, toolResults }); const steering = steeringMessagesFromExecution ?? ((await config.getSteeringMessages?.()) || []); - const asides = (await config.getAsideMessages?.()) || []; + const asides = resolveAsides(await config.getAsideMessages?.()); pendingMessages = asides.length > 0 ? [...steering, ...asides] : steering; } // Agent would stop here. Drain non-interrupting asides + follow-up messages. await config.onBeforeYield?.(); - const asideMessages = (await config.getAsideMessages?.()) || []; + const asideMessages = resolveAsides(await config.getAsideMessages?.()); const followUpMessages = (await config.getFollowUpMessages?.()) || []; if (asideMessages.length > 0 || followUpMessages.length > 0) { // Set as pending so the inner loop processes them before stopping. diff --git a/packages/agent/src/agent.ts b/packages/agent/src/agent.ts index 2dbe08de0..3e6a5dfb1 100644 --- a/packages/agent/src/agent.ts +++ b/packages/agent/src/agent.ts @@ -33,6 +33,7 @@ import type { AgentState, AgentTool, AgentToolContext, + AsideMessage, StreamFn, ToolCallContext, } from "./types"; @@ -319,7 +320,7 @@ export class Agent { #onAssistantMessageEvent?: (message: AssistantMessage, event: AssistantMessageEvent) => void; #onHarmonyLeak?: (event: HarmonyAuditEvent) => void | Promise; #onBeforeYield?: () => Promise | void; - #asideMessageProvider?: () => AgentMessage[] | Promise; + #asideMessageProvider?: () => AsideMessage[] | Promise; #telemetry?: AgentLoopConfig["telemetry"]; #appendOnlyContext?: AppendOnlyContextManager; @@ -635,7 +636,7 @@ export class Agent { * completions, late LSP diagnostics) drained at each step boundary. Never * aborts in-flight tools. See `AgentLoopConfig.getAsideMessages`. */ - setAsideMessageProvider(fn: (() => AgentMessage[] | Promise) | undefined): void { + setAsideMessageProvider(fn: (() => AsideMessage[] | Promise) | undefined): void { this.#asideMessageProvider = fn; } diff --git a/packages/agent/src/types.ts b/packages/agent/src/types.ts index de32f298d..52db396e0 100644 --- a/packages/agent/src/types.ts +++ b/packages/agent/src/types.ts @@ -26,6 +26,14 @@ export type StreamFn = ( ...args: Parameters ) => AssistantMessageEventStream | Promise; +/** + * An aside entry: a ready {@link AgentMessage}, or a sync thunk evaluated at + * injection time that returns the message to inject or `null` to skip it. Thunks + * let the producer make the final inject-or-drop decision against current state + * (e.g. dropping late diagnostics a newer edit superseded). + */ +export type AsideMessage = AgentMessage | (() => AgentMessage | null); + /** * Configuration for the agent loop. */ @@ -142,7 +150,7 @@ export interface AgentLoopConfig extends SimpleStreamOptions { * fully stop. Returned messages are appended to the context with normal * message events and keep the loop running so the model can react. */ - getAsideMessages?: () => Promise; + getAsideMessages?: () => Promise; /** * Hook fired right before the loop would exit. * diff --git a/packages/agent/test/agent-loop.test.ts b/packages/agent/test/agent-loop.test.ts index 743b1f362..908a568ef 100644 --- a/packages/agent/test/agent-loop.test.ts +++ b/packages/agent/test/agent-loop.test.ts @@ -841,6 +841,34 @@ describe("agentLoop with AgentMessage", () => { ); expect(sawAsideInContext).toBe(true); }); + + it("evaluates aside thunks at injection and skips ones that return null", async () => { + const context: AgentContext = { systemPrompt: [""], messages: [], tools: [] }; + const mock = createMockModel({ responses: [{ content: ["done"] }] }); + let polls = 0; + const config: AgentLoopConfig = { + model: mock.model, + convertToLlm: identityConverter, + // A lazy aside that decides, at injection time, NOT to inject (e.g. superseded). + getAsideMessages: async () => { + polls++; + return [() => null]; + }, + }; + + const events: AgentEvent[] = []; + const stream = agentLoop([createUserMessage("hi")], context, config, undefined, mock.stream); + for await (const event of stream) { + events.push(event); + } + + // The thunk was consulted... + expect(polls).toBeGreaterThan(0); + // ...but a null result injects nothing and triggers no wasted continuation turn. + const userStarts = events.filter(e => e.type === "message_start" && e.message.role === "user"); + expect(userStarts).toHaveLength(1); // only the original prompt + expect(mock.calls).toHaveLength(1); + }); }); it("refreshes tools and system prompt between same-turn model calls", async () => { diff --git a/packages/coding-agent/src/session/agent-session.ts b/packages/coding-agent/src/session/agent-session.ts index 51f083502..2a3204943 100644 --- a/packages/coding-agent/src/session/agent-session.ts +++ b/packages/coding-agent/src/session/agent-session.ts @@ -1191,7 +1191,7 @@ export class AgentSession { // Background-job completions / late diagnostics are pulled into the run at // each step boundary as non-interrupting asides (see Agent.getAsideMessages), // so they reach the model between requests without waiting for a yield. - this.agent.setAsideMessageProvider(() => this.yieldQueue.drainMessages()); + this.agent.setAsideMessageProvider(() => this.yieldQueue.drainLazy()); this.#convertToLlm = config.convertToLlm ?? convertToLlm; this.#rebuildSystemPrompt = config.rebuildSystemPrompt; this.#getMcpServerInstructions = config.getMcpServerInstructions; diff --git a/packages/coding-agent/src/session/yield-queue.ts b/packages/coding-agent/src/session/yield-queue.ts index bb92a8e5e..9473c941a 100644 --- a/packages/coding-agent/src/session/yield-queue.ts +++ b/packages/coding-agent/src/session/yield-queue.ts @@ -103,21 +103,21 @@ export class YieldQueue { } /** - * Build and remove all queued messages, applying each dispatcher's staleness - * filter. No injection side effects — used for pull-based delivery at agent - * step boundaries (see `Agent.setAsideMessageProvider`), so background-job - * completions and late diagnostics reach the model between requests without - * the agent having to stop. + * Snapshot and remove all queued entries, returning one lazy thunk per kind. + * Each thunk applies the dispatcher's staleness filter and builds the batched + * message only when called — so the consumer (the agent loop) decides, at the + * moment it injects, whether the message is still worth delivering (a thunk may + * return null to skip). Background-job completions and late diagnostics reach + * the model between requests without the agent having to stop. */ - drainMessages(): AgentMessage[] { - const messages: AgentMessage[] = []; + drainLazy(): Array<() => AgentMessage | null> { + const thunks: Array<() => AgentMessage | null> = []; for (const [kind, dispatcher] of this.#dispatchers) { const entries = this.#drain(kind); if (entries.length === 0) continue; - const message = this.#build(kind, dispatcher, entries); - if (message) messages.push(message); + thunks.push(() => this.#build(kind, dispatcher, entries)); } - return messages; + return thunks; } clear(): void { diff --git a/packages/coding-agent/src/tools/index.ts b/packages/coding-agent/src/tools/index.ts index 85f832cbe..b70d56684 100644 --- a/packages/coding-agent/src/tools/index.ts +++ b/packages/coding-agent/src/tools/index.ts @@ -133,8 +133,9 @@ export interface DeferredDiagnosticsEntry { /** True when any message is error severity. */ errored: boolean; /** - * Evaluated at flush time: drop the entry when a newer edit to the same file - * has superseded it, so the model never sees diagnostics for stale content. + * Evaluated at injection time (in the dispatcher's stale check): drop the entry + * when a newer mutation to the same file has superseded it, so the model never + * sees diagnostics for stale content. */ isStale(): boolean; } diff --git a/packages/coding-agent/test/session/yield-queue.test.ts b/packages/coding-agent/test/session/yield-queue.test.ts index 8264d4ad4..fa749448a 100644 --- a/packages/coding-agent/test/session/yield-queue.test.ts +++ b/packages/coding-agent/test/session/yield-queue.test.ts @@ -152,22 +152,45 @@ describe("YieldQueue", () => { expect(harness.streamingMessages.map(messageText)).toEqual(["second", "first"]); }); - test("drainMessages builds non-stale entries, clears the queue, returns nothing on re-drain", async () => { + test("drainLazy snapshots+clears immediately but defers build+staleness to the thunk", () => { const harness = createHarness(true); + const staleIds = new Set(); harness.queue.register("items", { - isStale: entry => entry.stale === true, + isStale: entry => staleIds.has(entry.id), build: entries => userMessage(entries.map(entry => entry.id).join(",")), }); - harness.queue.enqueue("items", { id: "keep" }); - harness.queue.enqueue("items", { id: "drop", stale: true }); + harness.queue.enqueue("items", { id: "a" }); + harness.queue.enqueue("items", { id: "b" }); - const drained = harness.queue.drainMessages(); - expect(drained.map(messageText)).toEqual(["keep"]); - // Pull-based drain has no injection side effects and empties the queue. + // Snapshot + clear happens at drain; the queue is emptied immediately. + const thunks = harness.queue.drainLazy(); + expect(thunks).toHaveLength(1); + expect(harness.queue.has()).toBe(false); + + // A mutation AFTER drainLazy but BEFORE the thunk runs supersedes "b". + staleIds.add("b"); + + // The thunk evaluates staleness at call time (injection), dropping "b". + const message = thunks[0]!(); + expect(message && messageText(message)).toBe("a"); + // No injection side effects from the pull path. expect(harness.streamingMessages).toHaveLength(0); expect(harness.idleBatches).toHaveLength(0); - expect(harness.queue.has()).toBe(false); - expect(harness.queue.drainMessages()).toEqual([]); + }); + + test("drainLazy thunk returns null when everything is stale by injection time", () => { + const harness = createHarness(true); + let stale = false; + harness.queue.register("items", { + isStale: () => stale, + build: entries => userMessage(entries.map(entry => entry.id).join(",")), + }); + harness.queue.enqueue("items", { id: "x" }); + + const thunks = harness.queue.drainLazy(); + stale = true; // superseded between drain and injection + + expect(thunks[0]!()).toBeNull(); }); });