fix(agent): resolved deferred asides at injection to drop stale messages
- Introduced `AsideMessage` as a message-or-thunk union so aside providers can defer injection decisions. - Updated agent loop handling to resolve aside thunks at injection time and skip entries that returned `null`, then switched the session yield queue to `drainLazy` for deferred message building. - Added tests validating lazy aside evaluation and staleness-aware dropping when everything becomes stale after dequeueing.
This commit is contained in:
@@ -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.
|
||||
|
||||
@@ -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<void>;
|
||||
#onBeforeYield?: () => Promise<void> | void;
|
||||
#asideMessageProvider?: () => AgentMessage[] | Promise<AgentMessage[]>;
|
||||
#asideMessageProvider?: () => AsideMessage[] | Promise<AsideMessage[]>;
|
||||
#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<AgentMessage[]>) | undefined): void {
|
||||
setAsideMessageProvider(fn: (() => AsideMessage[] | Promise<AsideMessage[]>) | undefined): void {
|
||||
this.#asideMessageProvider = fn;
|
||||
}
|
||||
|
||||
|
||||
@@ -26,6 +26,14 @@ export type StreamFn = (
|
||||
...args: Parameters<typeof streamSimple>
|
||||
) => AssistantMessageEventStream | Promise<AssistantMessageEventStream>;
|
||||
|
||||
/**
|
||||
* 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<AgentMessage[]>;
|
||||
getAsideMessages?: () => Promise<AsideMessage[]>;
|
||||
/**
|
||||
* Hook fired right before the loop would exit.
|
||||
*
|
||||
|
||||
@@ -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 () => {
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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<string>();
|
||||
harness.queue.register<Entry>("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<Entry>("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();
|
||||
});
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user