fix(agents): scope trailing progress to the provider turn

Native-async tool fragments from one response share a provider turn; reset
the progress-only fact on turn_start instead of each tool-call fragment.
This commit is contained in:
Ayaan Zaidi 2026-10-01 18:12:06 +05:30
parent e969c5c3f7
commit e27810bce1
No known key found for this signature in database
6 changed files with 75 additions and 59 deletions

View file

@ -13,6 +13,7 @@ import { getReplyPayloadMetadata } from "../../../auto-reply/reply-payload.js";
import { HEARTBEAT_TOKEN } from "../../../auto-reply/tokens.js";
import { runAgentLoop } from "../../../plugin-sdk/agent-core.js";
import { createHookRunner } from "../../../plugins/hooks.js";
import type { Deferred } from "../../../shared/deferred.js";
import { createSubscribedSessionHarness } from "../../embedded-agent-subscribe.e2e-harness.js";
import type { AgentMessage } from "../../runtime/index.js";
import { makeAssistantMessageFixture } from "../../test-helpers/assistant-message-fixtures.js";
@ -781,7 +782,8 @@ describe("trailing source progress at runtime settlement", () => {
contextWindow: 100_000,
maxTokens: 8_000,
};
type ModelCall = { progress?: "sync" | "async"; read?: boolean };
// Native-async calls run while the provider streams; each settles before the next fragment.
type ModelCall = { tools?: Array<"progress" | "read">; async?: boolean };
async function settleRealLoop(calls: ModelCall[]) {
const sessionManager = {};
@ -802,7 +804,12 @@ describe("trailing source progress at runtime settlement", () => {
return { content: [{ type: "text" as const, text: "ok" }], details: {} };
},
});
const progressSettled = createDeferred();
const toolSettled = new Map<string, Deferred>();
const settledFor = (toolCallId: string) => {
const settled = toolSettled.get(toolCallId) ?? createDeferred();
toolSettled.set(toolCallId, settled);
return settled;
};
await runAgentLoop(
[{ role: "user", content: "Run the report.", timestamp: 0 }],
{
@ -842,9 +849,9 @@ describe("trailing source progress at runtime settlement", () => {
if (event.type === "message_end") {
messages.push(event.message);
}
if (event.type === "tool_execution_end" && event.toolName === "message") {
if (event.type === "tool_execution_end") {
await subscription.waitForPendingEvents();
progressSettled.resolve();
settledFor(event.toolCallId).resolve();
}
if (event.type === "agent_end") {
await subscription.waitForPendingEvents();
@ -853,51 +860,47 @@ describe("trailing source progress at runtime settlement", () => {
undefined,
() => {
const call = calls[callIndex++] ?? {};
const toolCall: ToolCall | undefined = call.progress
? {
type: "toolCall",
id: `progress-${callIndex}`,
name: "message",
arguments: {
action: "send",
final: false,
target: "channel:source",
message: "Started the run, I will report back.",
},
...(call.progress === "async" ? { async: true as const } : {}),
}
: call.read
? { type: "toolCall", id: `read-${callIndex}`, name: "read", arguments: {} }
: undefined;
const toolCalls: ToolCall[] = (call.tools ?? []).map((name, index) => ({
type: "toolCall",
id: `${name}-${callIndex}-${index}`,
...(name === "progress"
? {
name: "message",
arguments: {
action: "send",
final: false,
target: "channel:source",
message: "Started the run, I will report back.",
},
}
: { name: "read", arguments: {} }),
...(call.async ? { async: true as const } : {}),
}));
const message = makeAssistantMessageFixture({
api: testModel.api,
provider: testModel.provider,
model: testModel.id,
content: [...(toolCall ? [toolCall] : []), { type: "text", text: "" }],
stopReason: toolCall && call.progress !== "async" ? "toolUse" : "stop",
content: [...toolCalls, { type: "text", text: "" }],
stopReason: toolCalls.length > 0 && !call.async ? "toolUse" : "stop",
errorMessage: undefined,
});
const stream = new AssistantMessageEventStream();
stream.push({ type: "start", partial: { ...message, content: [] } });
if (toolCall) {
stream.push({
type: "toolcall_end",
contentIndex: 0,
toolCall,
partial: { ...message, content: [toolCall] },
});
}
// A native-async send can settle while the provider is still streaming;
// its empty terminal fragment then arrives after the send's receipt.
const finish = () => {
void (async () => {
stream.push({ type: "start", partial: { ...message, content: [] } });
for (const [index, toolCall] of toolCalls.entries()) {
stream.push({
type: "toolcall_end",
contentIndex: index,
toolCall,
partial: { ...message, content: toolCalls.slice(0, index + 1) },
});
if (call.async) {
await settledFor(toolCall.id).promise;
}
}
stream.push({ type: "done", reason: "stop", message });
stream.end();
};
if (call.progress === "async") {
void progressSettled.promise.then(finish);
} else {
finish();
}
})();
return stream;
},
);
@ -929,7 +932,7 @@ describe("trailing source progress at runtime settlement", () => {
}
it("treats progress sent as the last tool batch as the reply", async () => {
const { result, finalizerInstruction } = await settleRealLoop([{ progress: "sync" }, {}]);
const { result, finalizerInstruction } = await settleRealLoop([{ tools: ["progress"] }, {}]);
expect(result.sourceReplyDeliveryState).toBe("delivered");
expect(result.messagingToolSentTargets?.map((send) => send.sourceReplyFinal)).toEqual([true]);
@ -938,8 +941,19 @@ describe("trailing source progress at runtime settlement", () => {
it("still finalizes when the turn continues after an async progress send", async () => {
const { result, finalizerInstruction } = await settleRealLoop([
{ progress: "async" },
{ read: true },
{ tools: ["progress"], async: true },
{ tools: ["read"] },
{},
]);
expect(result.sourceReplyDeliveryState).toBe("missing");
expect(result.messagingToolSentTargets?.map((send) => send.sourceReplyFinal)).toEqual([false]);
expect(finalizerInstruction).toContain("did not produce a user-visible answer");
});
it("still finalizes when async work and progress share one provider response", async () => {
const { result, finalizerInstruction } = await settleRealLoop([
{ tools: ["read", "progress"], async: true },
{},
]);

View file

@ -73,7 +73,7 @@ export function handleMessageEnd(
ctx.state.deterministicApprovalPromptPending = false;
ctx.state.deterministicApprovalPromptSent = false;
ctx.state.currentSourceMessagingToolSentTextsNormalized.length = 0;
ctx.state.toolBatchSourceProgress = undefined;
ctx.state.lastToolTurnOnlySourceProgress = undefined;
ctx.state.lastAssistant = undefined;
return;
}
@ -94,12 +94,6 @@ export function handleMessageEnd(
emitReasoningEnd(ctx);
}
ctx.noteLastAssistant(assistantMessage);
if (
Array.isArray(assistantMessage.content) &&
assistantMessage.content.some((block) => block.type === "toolCall")
) {
ctx.state.toolBatchSourceProgress = "open";
}
if (suppressVisibleAssistantOutput) {
appendRawStream(
() => ({

View file

@ -391,14 +391,13 @@ export async function handleToolExecutionEnd(
ctx.trimMessagingToolSent();
}
}
ctx.state.toolBatchSourceProgress =
ctx.state.turnToolsOnlySourceProgress =
(ctx.state.turnToolsOnlySourceProgress ?? true) &&
sourceReplyFinal === false &&
isMessagingSend &&
!isToolError &&
!messageDelivery?.partialDelivery &&
ctx.state.toolBatchSourceProgress
? "progress"
: undefined;
!messageDelivery?.partialDelivery;
ctx.state.lastToolTurnOnlySourceProgress = ctx.state.turnToolsOnlySourceProgress;
// Track committed reminders only when cron.add completed successfully.
if (
!isToolError &&

View file

@ -75,6 +75,12 @@ export function createEmbeddedAgentSessionEventHandler(ctx: EmbeddedAgentSubscri
case "message_end":
void scheduleEvent(evt, () => handleMessageEnd(ctx, evt));
return;
case "turn_start":
// Async tool fragments share one provider turn; only a new model call starts a batch.
void scheduleEvent(evt, () => {
ctx.state.turnToolsOnlySourceProgress = undefined;
});
return;
case "turn_end":
void scheduleEvent(evt, () => ctx.noteLastAssistant(evt.message));
return;

View file

@ -223,8 +223,10 @@ export type EmbeddedAgentSubscribeState = {
messageToolOnlySourceReplyDelivered: boolean;
sourceReplyDelivered?: true;
sourceReplyDeliveryState?: ReplyDeliveryState;
/** Latest tool batch: "open" before any result, "progress" while it only sent source progress. */
toolBatchSourceProgress?: "open" | "progress";
/** Whether the current provider turn's finished tools were only complete source progress. */
turnToolsOnlySourceProgress?: boolean;
/** The same fact for the latest provider turn that finished any tool. */
lastToolTurnOnlySourceProgress?: boolean;
successfulCronAdds: number;
pendingToolMediaUrls: string[];
pendingToolMediaAttachments?: ReplyMediaAttachment[];
@ -393,7 +395,8 @@ type ToolHandlerState = Pick<
| "messageToolOnlySourceReplyDelivered"
| "sourceReplyDelivered"
| "sourceReplyDeliveryState"
| "toolBatchSourceProgress"
| "turnToolsOnlySourceProgress"
| "lastToolTurnOnlySourceProgress"
| "messagingToolSentTargets"
| "heartbeatToolResponse"
| "successfulCronAdds"

View file

@ -474,7 +474,7 @@ export function subscribeEmbeddedAgentSession(params: SubscribeEmbeddedAgentSess
getMessagingToolSourceReplyPayloads: () => messagingToolSourceReplyPayloads.slice(),
getSourceReplyDelivered: () => state.sourceReplyDelivered,
getSourceReplyDeliveryState: () => state.sourceReplyDeliveryState,
endsWithSourceProgress: () => state.toolBatchSourceProgress === "progress",
endsWithSourceProgress: () => state.lastToolTurnOnlySourceProgress === true,
getHeartbeatToolResponse: () =>
state.heartbeatToolResponse ? { ...state.heartbeatToolResponse } : undefined,
getPendingToolMediaReply: () => readPendingToolMediaReply(state),