diff --git a/src/agents/embedded-agent-runner/run/abortable.test.ts b/src/agents/embedded-agent-runner/run/abortable.test.ts index 2c6e78bb2495..ad3fe372d843 100644 --- a/src/agents/embedded-agent-runner/run/abortable.test.ts +++ b/src/agents/embedded-agent-runner/run/abortable.test.ts @@ -82,4 +82,19 @@ describe("joinWithRunLivenessDeadline", () => { }); expect(onTimeout).not.toHaveBeenCalled(); }); + + it("runs work without an abort signal and remains bounded", async () => { + vi.useFakeTimers(); + try { + const joinWork = vi.fn(() => new Promise(() => {})); + const onTimeout = vi.fn(); + const join = joinWithRunLivenessDeadline({ joinWork, onTimeout }); + await vi.advanceTimersByTimeAsync(RUN_LIVENESS_JOIN_TIMEOUT_MS); + await join; + expect(joinWork).toHaveBeenCalledOnce(); + expect(onTimeout).toHaveBeenCalledOnce(); + } finally { + vi.useRealTimers(); + } + }); }); diff --git a/src/agents/embedded-agent-runner/run/abortable.ts b/src/agents/embedded-agent-runner/run/abortable.ts index c48b07429817..3c8e2496cc92 100644 --- a/src/agents/embedded-agent-runner/run/abortable.ts +++ b/src/agents/embedded-agent-runner/run/abortable.ts @@ -47,7 +47,7 @@ export const RUN_LIVENESS_JOIN_TIMEOUT_MS = 120_000; */ export function joinWithRunLivenessDeadline(input: { joinWork: () => Promise | void; - runAbortSignal: AbortSignal; + runAbortSignal?: AbortSignal; timeoutMs?: number; onTimeout: () => void; }): Promise { @@ -59,7 +59,7 @@ export function joinWithRunLivenessDeadline(input: { } settled = true; clearTimeout(timer); - input.runAbortSignal.removeEventListener("abort", onAbort); + input.runAbortSignal?.removeEventListener("abort", onAbort); if (reason === "timeout") { input.onTimeout(); } @@ -71,11 +71,11 @@ export function joinWithRunLivenessDeadline(input: { input.timeoutMs ?? RUN_LIVENESS_JOIN_TIMEOUT_MS, ); timer.unref?.(); - if (input.runAbortSignal.aborted) { + if (input.runAbortSignal?.aborted) { finish("abort"); return; } - input.runAbortSignal.addEventListener("abort", onAbort, { once: true }); + input.runAbortSignal?.addEventListener("abort", onAbort, { once: true }); Promise.resolve() .then(() => input.joinWork()) .then( diff --git a/src/agents/embedded-agent-runner/run/attempt-prompt-phase.test.ts b/src/agents/embedded-agent-runner/run/attempt-prompt-phase.test.ts index fe21d2cdc57b..9f4caa93a4ff 100644 --- a/src/agents/embedded-agent-runner/run/attempt-prompt-phase.test.ts +++ b/src/agents/embedded-agent-runner/run/attempt-prompt-phase.test.ts @@ -261,6 +261,7 @@ function createFixture() { setFinalPromptText, markBeforeAgentRunBlocked, markYieldAborted, + isRunBudgetTimeoutAbort: () => false, readYieldState: () => yieldState, stopAcceptingSteerMessages, takePendingMidTurnPrecheckRequest: () => undefined, @@ -409,6 +410,36 @@ describe("runEmbeddedAttemptPromptPhase", () => { ); }); + it("keeps a run-budget timeout failure-free for partial-output salvage", async () => { + const fixture = createFixture(); + const timeoutAbort = new Error("request timed out"); + mocks.submitPrompt.mockRejectedValueOnce(timeoutAbort); + mocks.handlePromptError.mockResolvedValueOnce({ + promptFailure: { error: timeoutAbort, source: "prompt" }, + }); + fixture.input.lifecycle.isRunBudgetTimeoutAbort = (error) => error === timeoutAbort; + + await runEmbeddedAttemptPromptPhase(fixture.input); + + expect(fixture.state.promptError).toBeNull(); + expect(fixture.state.promptErrorSource).toBeNull(); + }); + + it("records a provider failure that races a run-budget timeout", async () => { + const fixture = createFixture(); + const providerError = new Error("provider failed"); + mocks.submitPrompt.mockRejectedValueOnce(providerError); + mocks.handlePromptError.mockResolvedValueOnce({ + promptFailure: { error: providerError, source: "prompt" }, + }); + fixture.input.lifecycle.isRunBudgetTimeoutAbort = () => false; + + await runEmbeddedAttemptPromptPhase(fixture.input); + + expect(fixture.state.promptError).toBe(providerError); + expect(fixture.state.promptErrorSource).toBe("prompt"); + }); + it("releases steering when preflight skips provider submission", async () => { const fixture = createFixture(); const promptError = new Error("preflight rejected"); diff --git a/src/agents/embedded-agent-runner/run/attempt-prompt-phase.ts b/src/agents/embedded-agent-runner/run/attempt-prompt-phase.ts index a32982584bcf..eb77926f14fe 100644 --- a/src/agents/embedded-agent-runner/run/attempt-prompt-phase.ts +++ b/src/agents/embedded-agent-runner/run/attempt-prompt-phase.ts @@ -140,6 +140,7 @@ export async function runEmbeddedAttemptPromptPhase(input: { setFinalPromptText: (prompt: string) => void; markBeforeAgentRunBlocked: (outcome: BeforeAgentRunOutcome) => void; markYieldAborted: () => void; + isRunBudgetTimeoutAbort: (error: unknown) => boolean; readYieldState: () => Pick< PromptErrorInput, "yieldAbortSettled" | "yieldDetected" | "yieldMessage" @@ -398,7 +399,12 @@ export async function runEmbeddedAttemptPromptPhase(input: { withOwnedTranscriptWrite: input.withOwnedTranscriptWrite, ...input.lifecycle.readYieldState(), }); - if (promptErrorOutcome.promptFailure) { + // The timeout owner records its terminal before aborting the prompt. That + // abort is not a provider failure and must leave timeout salvage eligible. + if ( + promptErrorOutcome.promptFailure && + !input.lifecycle.isRunBudgetTimeoutAbort(promptErrorOutcome.promptFailure.error) + ) { patchState({ promptError: promptErrorOutcome.promptFailure.error, promptErrorSource: promptErrorOutcome.promptFailure.source, diff --git a/src/agents/embedded-agent-runner/run/attempt-settle.ts b/src/agents/embedded-agent-runner/run/attempt-settle.ts index 2e9310a0efaa..c4c9fad64d96 100644 --- a/src/agents/embedded-agent-runner/run/attempt-settle.ts +++ b/src/agents/embedded-agent-runner/run/attempt-settle.ts @@ -18,7 +18,11 @@ import type { NormalizedUsage } from "../../usage.js"; import { log } from "../logger.js"; import type { PromptCacheBreak, PromptCacheChange } from "../prompt-cache-observability.js"; import { clearActiveEmbeddedRun } from "../runs.js"; -import { joinWithRunLivenessDeadline, RUN_LIVENESS_JOIN_TIMEOUT_MS } from "./abortable.js"; +import { + isOpenClawAbortableWrapper, + joinWithRunLivenessDeadline, + RUN_LIVENESS_JOIN_TIMEOUT_MS, +} from "./abortable.js"; import type { EmbeddedAttemptExecutionPhaseInput, EmbeddedAttemptExecutionState, @@ -335,28 +339,60 @@ export async function runEmbeddedAttemptSettledPhase( source: "yield_cleanup", }); }, + isRunBudgetTimeoutAbort: (error) => + readTerminal().timedOutByRunBudget && + isOpenClawAbortableWrapper(error) && + error instanceof Error && + error.cause === input.runAbortController.signal.reason, readYieldState: input.lifecycle.readYieldState, stopAcceptingSteerMessages, takePendingMidTurnPrecheckRequest: contextGuards.takePendingMidTurnPrecheckRequest, }, }); - // Queued subscription handlers (block-reply delivery, tool events) are - // fire-and-forget during the turn; the pending-events join below is the only - // place the run waits for them. One hung handler (e.g. a stuck delivery - // dispatch lane) must not dead-end the turn until the run budget — 48h by - // default — so the join is bounded and settlement proceeds with a recorded - // warning instead of producing no visible outcome at all. - await joinWithRunLivenessDeadline({ - joinWork: waitForPendingEvents, - runAbortSignal: input.runAbortController.signal, - onTimeout: () => { - log.warn( - `pending subscription events did not settle within ${RUN_LIVENESS_JOIN_TIMEOUT_MS}ms; ` + - `proceeding to stream settlement: runId=${attempt.runId}`, - ); - }, - }); + // Only a failure-free run-budget terminal may publish buffered text. + const isFailureFreeRunBudgetTimeout = (): boolean => { + const terminal = readTerminal(); + return terminal.timedOutByRunBudget && !terminal.failed; + }; + const runBudgetTimeoutTerminal = isFailureFreeRunBudgetTimeout(); + const drainPendingEventsBounded = () => + joinWithRunLivenessDeadline({ + // Partial-reply callbacks cannot mutate the buffer and may be stalled + // on transport; timeout salvage needs only the serialized event chain. + joinWork: () => waitForPendingEvents({ includePartialReplies: false }), + onTimeout: () => { + log.warn( + `pending subscription events did not settle within ${RUN_LIVENESS_JOIN_TIMEOUT_MS}ms; ` + + `proceeding to stream settlement: runId=${attempt.runId}`, + ); + }, + }); + if (runBudgetTimeoutTerminal) { + // The timeout already aborted the signal; drain without racing it. + await drainPendingEventsBounded(); + } else { + await joinWithRunLivenessDeadline({ + joinWork: waitForPendingEvents, + runAbortSignal: input.runAbortController.signal, + onTimeout: () => { + log.warn( + `pending subscription events did not settle within ${RUN_LIVENESS_JOIN_TIMEOUT_MS}ms; ` + + `proceeding to stream settlement: runId=${attempt.runId}`, + ); + }, + }); + // A timeout can fire during the abort-aware join and resolve it before + // its queue drains. Re-read terminal ownership, then drain if eligible. + if (isFailureFreeRunBudgetTimeout()) { + await drainPendingEventsBounded(); + } + } + // Ownership can change during the drain; publish only after the final read. + const salvageTerminal = readTerminal(); + if (salvageTerminal.timedOutByRunBudget && !salvageTerminal.failed) { + subscription.flushPartialAssistantText(); + } const beforeAgentFinalizeRevisionReason = getBeforeAgentFinalizeRevisionReason(); const beforeAgentFinalizeRevisionEntryId = getBeforeAgentFinalizeRevisionEntryId(); let rewoundBeforeAgentFinalizeRevision = false; diff --git a/src/agents/embedded-agent-runner/run/attempt-spawn-workspace.test-support.ts b/src/agents/embedded-agent-runner/run/attempt-spawn-workspace.test-support.ts index 14216d1f62f2..94419231bdcb 100644 --- a/src/agents/embedded-agent-runner/run/attempt-spawn-workspace.test-support.ts +++ b/src/agents/embedded-agent-runner/run/attempt-spawn-workspace.test-support.ts @@ -138,6 +138,7 @@ function createSubscriptionMock(): SubscriptionMock { setTerminalLifecycleMeta: () => {}, waitForCompactionRetry: async () => {}, waitForPendingEvents: async () => {}, + flushPartialAssistantText: () => {}, getAcceptedSessionSpawns: () => [], getMessagingToolSentTexts: () => [] as string[], getMessagingToolSentMediaUrls: () => [] as string[], diff --git a/src/agents/embedded-agent-runner/run/attempt-stream-finalize.test.ts b/src/agents/embedded-agent-runner/run/attempt-stream-finalize.test.ts index b9ae3731bca9..5177899baad5 100644 --- a/src/agents/embedded-agent-runner/run/attempt-stream-finalize.test.ts +++ b/src/agents/embedded-agent-runner/run/attempt-stream-finalize.test.ts @@ -41,12 +41,13 @@ type SettleMockInput = { }; type FixtureOverrides = { activeSession?: SettledInput["prepared"]["sessionRuntime"]["agentSession"]["activeSession"]; + flushPartialAssistantText?: () => void; getBeforeAgentFinalizeRevisionEntryId?: () => string | undefined; getBeforeAgentFinalizeRevisionReason?: () => string | undefined; repairedRejectedProviderReplay?: boolean; runAbortController?: AbortController; sessionManager?: SettledInput["prepared"]["sessionRuntime"]["sessionManager"]; - waitForPendingEvents?: () => Promise; + waitForPendingEvents?: (options?: { includePartialReplies?: boolean }) => Promise; }; function createFixture(overrides: FixtureOverrides = {}) { @@ -68,15 +69,19 @@ function createFixture(overrides: FixtureOverrides = {}) { } as never); const waitForPendingEvents = overrides.waitForPendingEvents ?? - vi.fn(async () => { - order.push("pending-events"); + vi.fn(async (options?: { includePartialReplies?: boolean }) => { + order.push( + options?.includePartialReplies === false ? "pending-event-chain" : "pending-events", + ); }); const getBeforeAgentFinalizeRevisionReason = overrides.getBeforeAgentFinalizeRevisionReason ?? (() => "revision changed"); const getBeforeAgentFinalizeRevisionEntryId = overrides.getBeforeAgentFinalizeRevisionEntryId ?? (() => undefined); + const flushPartialAssistantText = overrides.flushPartialAssistantText ?? vi.fn(); const unsubscribe = vi.fn(); const subscription = { + flushPartialAssistantText, isCompacting: vi.fn(() => false), unsubscribe, waitForPendingEvents, @@ -205,6 +210,25 @@ function createFixture(overrides: FixtureOverrides = {}) { markYieldAborted = promptInput.lifecycle.markYieldAborted; return { promptStartedAt: 100 }; }); + mocks.settleStream.mockResolvedValue({ + promptError: null, + promptErrorSource: null, + timedOutDuringCompaction: false, + compactionOccurredThisAttempt: false, + messagesSnapshot: [], + sessionIdUsed: "session-1", + lastAssistant: undefined, + currentAttemptAssistant: undefined, + currentAttemptCompletedAssistant: undefined, + attemptUsage: undefined, + cacheBreak: null, + lastCallUsage: undefined, + promptCache: undefined, + }); + mocks.completeAfterTurn.mockResolvedValue({ + sessionIdUsed: "session-1", + sessionFileUsed: "session.jsonl", + }); mocks.completeResult.mockImplementation((resultInput) => ({ sessionIdUsed: resultInput.state.sessionIdUsed, sessionFileUsed: resultInput.state.sessionFileUsed, @@ -213,6 +237,7 @@ function createFixture(overrides: FixtureOverrides = {}) { return { activeSession, + flushPartialAssistantText, input, markYieldAborted: () => markYieldAborted?.(), order, @@ -605,4 +630,251 @@ describe("runEmbeddedAttemptSettledPhase stream finalization", () => { expect(mocks.settleStream).toHaveBeenCalledOnce(); expect(mocks.completeAfterTurn).not.toHaveBeenCalled(); }); + + it("drains queued events after a run-budget abort before re-flushing partial assistant text", async () => { + // abortRun(true) aborts the run signal synchronously before settlement, so the abort-aware join returns + // without draining. The run-budget terminal must still drain the serialized + // event chain (bounded) so a message_update queued behind the abort commits + // before the re-flush. + const abortController = new AbortController(); + abortController.abort(new Error("run budget exceeded")); + const fixture = createFixture({ + runAbortController: abortController, + waitForPendingEvents: vi.fn(async () => { + fixture.order.push("pending-event-chain"); + }), + flushPartialAssistantText: vi.fn(() => { + fixture.order.push("flush-partial"); + }), + }); + fixture.state.terminal = { kind: "timeout", phase: "prompt", source: "run_budget" }; + + await expect(runEmbeddedAttemptSettledPhase(fixture.input)).resolves.toEqual({ + sessionIdUsed: "session-1", + sessionFileUsed: "session.jsonl", + }); + + // The queued-event drain must run (and complete) BEFORE the re-flush reads + // the buffer; with the abort-aware join this ordering was unreachable. The + // timeout salvage path drains only the serialized event chain (queue-only), + // not partial-reply fan-out callbacks. + expect(fixture.order).toEqual(["pending-event-chain", "flush-partial"]); + expect(mocks.settleStream).toHaveBeenCalledOnce(); + }); + + it("discards buffered partial text when an external abort supersedes the run-budget timeout during the drain", async () => { + // Partial output must be committed only after terminal ownership is final. The drain is awaited, then the + // terminal is re-read: if an external abort wins while the queued chain + // drains, the run-budget timeout no longer owns the terminal and the + // buffered text must NOT be published. + const abortController = new AbortController(); + abortController.abort(new Error("run budget exceeded")); + const fixture = createFixture({ + runAbortController: abortController, + waitForPendingEvents: vi.fn(async () => { + fixture.order.push("pending-event-chain"); + // External abort lands while the queued chain drains. + fixture.state.terminal = { kind: "aborted", source: "external" }; + }), + flushPartialAssistantText: vi.fn(() => { + fixture.order.push("flush-partial"); + }), + }); + fixture.state.terminal = { kind: "timeout", phase: "prompt", source: "run_budget" }; + + await expect(runEmbeddedAttemptSettledPhase(fixture.input)).resolves.toEqual({ + sessionIdUsed: "session-1", + sessionFileUsed: "session.jsonl", + }); + + // The drain still ran (bounded, abort-independent), but the superseded + // terminal discards the buffered text: no flush, no partial output. + expect(fixture.order).toEqual(["pending-event-chain"]); + expect(fixture.flushPartialAssistantText).not.toHaveBeenCalled(); + expect(mocks.settleStream).toHaveBeenCalledOnce(); + }); + + it("discards buffered partial text when a provider failure is attached before the terminal-owned flush", async () => { + // A provider error queued behind the run-budget abort is merged into the terminal before the + // post-drain flush decision. The failure-terminal invariant must suppress + // the salvage — a timed-out run that also failed must not publish partial + // output. + const abortController = new AbortController(); + abortController.abort(new Error("run budget exceeded")); + const fixture = createFixture({ + runAbortController: abortController, + waitForPendingEvents: vi.fn(async () => { + fixture.order.push("pending-event-chain"); + // Provider failure is recorded while the queued chain drains. + fixture.state.terminal = { + kind: "timeout", + phase: "prompt", + source: "run_budget", + failure: { source: "prompt", error: new Error("provider stream failed") }, + }; + }), + flushPartialAssistantText: vi.fn(() => { + fixture.order.push("flush-partial"); + }), + }); + fixture.state.terminal = { kind: "timeout", phase: "prompt", source: "run_budget" }; + + await expect(runEmbeddedAttemptSettledPhase(fixture.input)).resolves.toEqual({ + sessionIdUsed: "session-1", + sessionFileUsed: "session.jsonl", + }); + + // The drain still ran (bounded, abort-independent), but the attached + // provider failure discards the buffered text: no flush, no partial output. + expect(fixture.order).toEqual(["pending-event-chain"]); + expect(fixture.flushPartialAssistantText).not.toHaveBeenCalled(); + expect(mocks.settleStream).toHaveBeenCalledOnce(); + }); + + it("skips the bounded drain when a provider failure is already attached before settlement", async () => { + // A failure already attached to the run-budget terminal must prevent the drain. The bounded drain can only stall settlement + // for the full 120s liveness deadline when a serialized handler is wedged, + // and partial output would be discarded by the flush gate anyway. A failed + // run-budget terminal must skip the drain entirely so a wedged event chain + // cannot delay a failed run. + const abortController = new AbortController(); + abortController.abort(new Error("run budget exceeded")); + // Wedged serialized handler: never resolves. Pre-fix the bounded drain + // would wait the full RUN_LIVENESS_JOIN_TIMEOUT_MS (120s) before settling. + const chainGate = new Promise(() => { + // Intentionally never resolved; reaching the await means the drain ran, + // which would be a regression. + }); + const waitForPendingEvents = vi.fn(async () => { + fixture.order.push("pending-event-chain"); + await chainGate; + }); + const fixture = createFixture({ + runAbortController: abortController, + waitForPendingEvents, + flushPartialAssistantText: vi.fn(() => { + fixture.order.push("flush-partial"); + }), + }); + // Failure is attached before settlement chooses whether to drain. + fixture.state.terminal = { + kind: "timeout", + phase: "prompt", + source: "run_budget", + failure: { source: "prompt", error: new Error("provider stream failed") }, + }; + mocks.settleStream.mockResolvedValue({ + promptError: new Error("provider stream failed"), + promptErrorSource: "prompt", + timedOutDuringCompaction: false, + compactionOccurredThisAttempt: false, + messagesSnapshot: [], + sessionIdUsed: "session-1", + lastAssistant: undefined, + currentAttemptAssistant: undefined, + currentAttemptCompletedAssistant: undefined, + attemptUsage: undefined, + cacheBreak: null, + lastCallUsage: undefined, + promptCache: undefined, + }); + mocks.completeAfterTurn.mockResolvedValue({ + sessionIdUsed: "session-1", + sessionFileUsed: "session.jsonl", + }); + + // Settlement must resolve immediately without entering the bounded drain. + // A 120s wall-clock guard ensures pre-fix (which would drain the wedged + // chain) fails fast rather than hanging the suite. + await expect( + Promise.race([ + runEmbeddedAttemptSettledPhase(fixture.input), + new Promise((_, reject) => { + setTimeout( + () => reject(new Error("settlement stalled on the bounded drain")), + 5_000, + ).unref?.(); + }), + ]), + ).resolves.toEqual({ + sessionIdUsed: "session-1", + sessionFileUsed: "session.jsonl", + }); + + // The bounded drain was skipped: the wedged serialized chain was never + // awaited, and no partial text was flushed (failed terminal). + expect(waitForPendingEvents).not.toHaveBeenCalled(); + expect(fixture.flushPartialAssistantText).not.toHaveBeenCalled(); + expect(fixture.order).toEqual([]); + expect(mocks.settleStream).toHaveBeenCalledOnce(); + }); + + it("re-drains queued events when the run-budget timeout fires during the abort-aware join", async () => { + // Settlement starts with a non-budget terminal and waits on the abort-aware join. If the run-budget + // timer fires while that join is pending, the abort resolves the join + // immediately WITHOUT draining; the salvage must then run the bounded + // drain before flushing so a queued suffix is not lost. + const abortController = new AbortController(); + let releaseJoin!: () => void; + const joinGate = new Promise((resolve) => { + releaseJoin = resolve; + }); + const fixture = createFixture({ + runAbortController: abortController, + waitForPendingEvents: vi.fn(async (options) => { + if (options?.includePartialReplies === false) { + fixture.order.push("pending-event-chain"); + return; + } + fixture.order.push("pending-events"); + await joinGate; + }), + flushPartialAssistantText: vi.fn(() => { + fixture.order.push("flush-partial"); + }), + }); + fixture.state.terminal = { kind: "ok" }; + + const settlePromise = runEmbeddedAttemptSettledPhase(fixture.input); + // Let the abort-aware join reach waitForPendingEvents, then fire the + // run-budget timeout while the join is pending. + await vi.waitFor(() => { + expect(fixture.order).toContain("pending-events"); + }); + fixture.state.terminal = { kind: "timeout", phase: "prompt", source: "run_budget" }; + abortController.abort(new Error("run budget exceeded")); + // Release the join gate so the bounded re-drain can complete too. + releaseJoin(); + + await expect(settlePromise).resolves.toEqual({ + sessionIdUsed: "session-1", + sessionFileUsed: "session.jsonl", + }); + + // The abort-aware join resolved on abort without draining; the bounded + // re-drain must run BEFORE the salvage flush (pre-fix code flushed + // without it, losing the queued suffix). The re-drain is queue-only + // (pending-event-chain), separate from the abort-aware join's + // pending-events wait. + expect(fixture.order.filter((entry) => entry === "pending-events")).toHaveLength(1); + expect(fixture.order.filter((entry) => entry === "pending-event-chain")).toHaveLength(1); + expect(fixture.order).toEqual(["pending-events", "pending-event-chain", "flush-partial"]); + expect(mocks.settleStream).toHaveBeenCalledOnce(); + }); + + it("does not re-flush partial assistant text on non-run-budget terminals", async () => { + // Cancellation and provider-failure aborts must not publish partial output through settlement. + const abortController = new AbortController(); + abortController.abort(new Error("operator cancel")); + const fixture = createFixture({ runAbortController: abortController }); + fixture.state.terminal = { kind: "aborted", source: "external" }; + + await expect(runEmbeddedAttemptSettledPhase(fixture.input)).resolves.toEqual({ + sessionIdUsed: "session-1", + sessionFileUsed: "session.jsonl", + }); + + expect(fixture.flushPartialAssistantText).not.toHaveBeenCalled(); + expect(mocks.settleStream).toHaveBeenCalledOnce(); + }); }); diff --git a/src/agents/embedded-agent-runner/run/attempt-timeout-prepare.test.ts b/src/agents/embedded-agent-runner/run/attempt-timeout-prepare.test.ts index ef9902f30249..8fdaa75428de 100644 --- a/src/agents/embedded-agent-runner/run/attempt-timeout-prepare.test.ts +++ b/src/agents/embedded-agent-runner/run/attempt-timeout-prepare.test.ts @@ -63,6 +63,14 @@ describe("prepareEmbeddedAttemptTimeout", () => { expect(harness.markTimedOutByRunBudget).toHaveBeenCalledOnce(); expect(harness.abortRun).toHaveBeenCalledWith(true); + // The run-budget marker must be recorded before the abort so settlement + // can re-confirm terminal ownership before committing partial output; the + // timeout callback itself never commits buffered text. + const markOrder = harness.markTimedOutByRunBudget.mock.invocationCallOrder[0]; + const abortOrder = harness.abortRun.mock.invocationCallOrder[0]; + expect(markOrder).toBeDefined(); + expect(abortOrder).toBeDefined(); + expect(markOrder ?? -1).toBeLessThan(abortOrder ?? -1); harness.timeout.clearTimers(); }); diff --git a/src/agents/embedded-agent-runner/run/attempt-timeout-prepare.ts b/src/agents/embedded-agent-runner/run/attempt-timeout-prepare.ts index 77a056805482..f605e85c7293 100644 --- a/src/agents/embedded-agent-runner/run/attempt-timeout-prepare.ts +++ b/src/agents/embedded-agent-runner/run/attempt-timeout-prepare.ts @@ -71,6 +71,8 @@ export function prepareEmbeddedAttemptTimeout(input: { ) { input.markTimedOutDuringCompaction(); } + // Settlement owns partial-output publication because abort or failure + // can still supersede this timeout while queued events drain. input.markTimedOutByRunBudget(); input.abortRun(true); if (!abortWarnTimer) { diff --git a/src/agents/embedded-agent-subscribe.handlers.messages.lifecycle.ts b/src/agents/embedded-agent-subscribe.handlers.messages.lifecycle.ts index cd115c1203f3..06eb825962c8 100644 --- a/src/agents/embedded-agent-subscribe.handlers.messages.lifecycle.ts +++ b/src/agents/embedded-agent-subscribe.handlers.messages.lifecycle.ts @@ -247,6 +247,8 @@ export function handleMessageEnd( const finalizeMessageEnd = () => { ctx.state.deltaBuffer = ""; ctx.state.thinkingTagStream = createThinkingTagStreamState(); + ctx.state.deltaBufferIsCommentary = false; + ctx.state.hasFlushedPartialText = false; ctx.state.blockBuffer = ""; ctx.blockChunker?.reset(); ctx.state.blockState.thinking = false; diff --git a/src/agents/embedded-agent-subscribe.handlers.messages.update.ts b/src/agents/embedded-agent-subscribe.handlers.messages.update.ts index 83f1518811e3..0916b74fb0dd 100644 --- a/src/agents/embedded-agent-subscribe.handlers.messages.update.ts +++ b/src/agents/embedded-agent-subscribe.handlers.messages.update.ts @@ -252,6 +252,7 @@ export function handleMessageUpdate( if (isResponsesCommentary && chunk) { // Keep cumulative end events monotonic without feeding commentary into reply buffers. ctx.state.deltaBuffer += chunk; + ctx.state.deltaBufferIsCommentary = true; } const commentaryText = !chunk && (!isResponsesCommentary || !hadResponsesCommentaryText) @@ -289,6 +290,7 @@ export function handleMessageUpdate( if (chunk) { ctx.state.deltaBuffer += chunk; + ctx.state.deltaBufferIsCommentary = false; if (!skipLiveStream && !shouldUsePhaseAwareBlockReply) { if (!isPhasePendingAnthropicText && !isPhasePendingCompletionsText) { appendBlockReplyChunk(ctx, chunk); diff --git a/src/agents/embedded-agent-subscribe.handlers.types.ts b/src/agents/embedded-agent-subscribe.handlers.types.ts index 68840dff92f0..24d1353c882f 100644 --- a/src/agents/embedded-agent-subscribe.handlers.types.ts +++ b/src/agents/embedded-agent-subscribe.handlers.types.ts @@ -117,6 +117,16 @@ export type EmbeddedAgentSubscribeState = { deltaBuffer: string; /** Scanner state shares deltaBuffer's lifecycle so each provider byte is parsed once. */ thinkingTagStream: ThinkingTagStreamState; + /** + * True while the buffered stream text belongs to an explicit commentary + * item (e.g. the Responses API "commentary" phase). Commentary is routed to + * a separate lane by the normal stream path, so the run-budget timeout + * flush must skip it too: flushing the raw deltaBuffer without this marker + * would publish reasoning/commentary bytes as assistant text. + */ + deltaBufferIsCommentary: boolean; + /** Whether timeout settlement committed visible text for this message. */ + hasFlushedPartialText: boolean; blockBuffer: string; blockState: { thinking: boolean; diff --git a/src/agents/embedded-agent-subscribe.partial-reply-lifecycle.test.ts b/src/agents/embedded-agent-subscribe.partial-reply-lifecycle.test.ts index e4106da1bd0d..a017c4c9ff3a 100644 --- a/src/agents/embedded-agent-subscribe.partial-reply-lifecycle.test.ts +++ b/src/agents/embedded-agent-subscribe.partial-reply-lifecycle.test.ts @@ -96,4 +96,57 @@ describe("subscribeEmbeddedAgentSession partial reply lifecycle", () => { `assistant partial reply callback failed: ${String(callbackError)}`, ); }); + + it("queue-only drain is not blocked by a stalled partial reply callback", async () => { + // Timeout salvage drains only the serialized event chain — the queue whose handlers mutate the assistant + // text buffer. A stalled onPartialReply transport callback is external + // fan-out and cannot change the buffered text, so it must not hold an + // already-aborted run in settlement. Pre-fix, the salvage path used + // waitForPendingEvents, which also awaits pendingPartialReplyTasks; a + // stalled callback would block the drain until the bounded liveness + // deadline (120s) elapsed. + const onPartialReply = vi.fn(() => new Promise(() => {})); + const { emit, subscription } = createSubscribedSessionHarness({ + runId: "run-stalled-partial-callback", + onPartialReply, + }); + + // First delta fires the partial-reply callback (stalled, never resolves). + emit({ + type: "message_update", + message: { role: "assistant" }, + assistantMessageEvent: { type: "text_delta", delta: "partial " }, + }); + await vi.waitFor(() => expect(onPartialReply).toHaveBeenCalledOnce()); + + // A second delta queues behind the first on the serialized event chain. + emit({ + type: "message_update", + message: { role: "assistant" }, + assistantMessageEvent: { type: "text_delta", delta: "answer" }, + }); + + // Queue-only drain completes promptly: the event chain is null once both + // deltas are processed, regardless of the stalled fan-out callback. + await expect( + Promise.race([ + subscription.waitForPendingEvents({ includePartialReplies: false }).then(() => "drained"), + new Promise<"timeout">((resolve) => { + setTimeout(() => resolve("timeout"), 1000); + }), + ]), + ).resolves.toBe("drained"); + + // The broad join still observes the stalled callback (proving the two + // drains are genuinely distinct and the queue-only split is meaningful). + const broad = Promise.race([ + subscription.waitForPendingEvents().then(() => "drained"), + new Promise<"timeout">((resolve) => { + setTimeout(() => resolve("timeout"), 1000); + }), + ]); + await expect(broad).resolves.toBe("timeout"); + + subscription.unsubscribe(); + }); }); diff --git a/src/agents/embedded-agent-subscribe.reply-delivery.ts b/src/agents/embedded-agent-subscribe.reply-delivery.ts index 8a80abe08caf..25f733bf2638 100644 --- a/src/agents/embedded-agent-subscribe.reply-delivery.ts +++ b/src/agents/embedded-agent-subscribe.reply-delivery.ts @@ -215,6 +215,18 @@ export function createReplyDelivery({ params, state, log }: ReplyDeliveryParams) rememberAssistantText(text); }; + const replaceCurrentAssistantText = (text: string) => { + const count = assistantTexts.length - state.assistantTextBaseline; + if (!text) { + assistantTexts.splice(state.assistantTextBaseline, count); + } else if (count > 0) { + assistantTexts.splice(state.assistantTextBaseline, count, text); + rememberAssistantText(text); + } else { + pushAssistantText(text); + } + }; + const finalizeAssistantTexts = (args: { text: string; addedDuringMessage: boolean; @@ -222,19 +234,22 @@ export function createReplyDelivery({ params, state, log }: ReplyDeliveryParams) }) => { const { text, addedDuringMessage, chunkerHasBuffered } = args; + // A run-budget timeout flush may already have committed partial text for + // this message. When message_end later finalizes the complete text, replace + // the flushed partial instead of appending a duplicate. The partial stays + // when message_end never arrives (hard run-budget abort) — that is the + // salvage the timeout flush exists for. + if (state.hasFlushedPartialText && text) { + replaceCurrentAssistantText(text); + state.hasFlushedPartialText = false; + state.assistantTextBaseline = assistantTexts.length; + return; + } + // If we're not streaming block replies, ensure the final payload includes // the final text even when interim streaming was enabled. if (state.includeReasoning && text && !params.onBlockReply) { - if (assistantTexts.length > state.assistantTextBaseline) { - assistantTexts.splice( - state.assistantTextBaseline, - assistantTexts.length - state.assistantTextBaseline, - text, - ); - rememberAssistantText(text); - } else { - pushAssistantText(text); - } + replaceCurrentAssistantText(text); state.suppressBlockChunks = true; } else if (!addedDuringMessage && !chunkerHasBuffered && text) { // Non-streaming models (no text_delta): ensure assistantTexts gets the final @@ -245,14 +260,17 @@ export function createReplyDelivery({ params, state, log }: ReplyDeliveryParams) state.assistantTextBaseline = assistantTexts.length; }; - const waitForPendingEvents = async () => { + const waitForPendingEvents = async (options?: { includePartialReplies?: boolean }) => { // Partial presentation stays concurrent with provider events, but terminal // settlement must observe callbacks launched while the event chain drains. - while (state.pendingEventChain || pendingPartialReplyTasks.size > 0) { - await Promise.allSettled([ - ...(state.pendingEventChain ? [state.pendingEventChain] : []), - ...pendingPartialReplyTasks, - ]); + const includePartialReplies = options?.includePartialReplies !== false; + while (true) { + const eventChain = state.pendingEventChain; + const partialReplyTasks = includePartialReplies ? [...pendingPartialReplyTasks] : []; + if (!eventChain && partialReplyTasks.length === 0) { + return; + } + await Promise.allSettled([...(eventChain ? [eventChain] : []), ...partialReplyTasks]); } }; @@ -267,6 +285,7 @@ export function createReplyDelivery({ params, state, log }: ReplyDeliveryParams) flushDeferredBlockReplies, pendingBlockReplyTasks, pushAssistantText, + replaceCurrentAssistantText, shouldSkipAssistantText, waitForPendingEvents, }; diff --git a/src/agents/embedded-agent-subscribe.run-state.ts b/src/agents/embedded-agent-subscribe.run-state.ts index a5d0c74d6596..32888627f934 100644 --- a/src/agents/embedded-agent-subscribe.run-state.ts +++ b/src/agents/embedded-agent-subscribe.run-state.ts @@ -35,6 +35,8 @@ export function createEmbeddedAgentSubscribeState( typeof params.onReasoningStream === "function", deltaBuffer: "", thinkingTagStream: createThinkingTagStreamState(), + deltaBufferIsCommentary: false, + hasFlushedPartialText: false, blockBuffer: "", // Track if a streamed chunk opened a block (stateful across chunks). blockState: { thinking: false, final: false, inlineCode: createInlineCodeState() }, diff --git a/src/agents/embedded-agent-subscribe.stream-rendering.ts b/src/agents/embedded-agent-subscribe.stream-rendering.ts index 599002062fc2..be2b4f7aa424 100644 --- a/src/agents/embedded-agent-subscribe.stream-rendering.ts +++ b/src/agents/embedded-agent-subscribe.stream-rendering.ts @@ -581,6 +581,8 @@ export function createStreamRendering({ const resetAssistantMessageState = (nextAssistantTextBaseline: number) => { state.deltaBuffer = ""; state.thinkingTagStream = createThinkingTagStreamState(); + state.deltaBufferIsCommentary = false; + state.hasFlushedPartialText = false; state.blockBuffer = ""; blockChunker?.reset(); replyDirectiveAccumulator.reset(); diff --git a/src/agents/embedded-agent-subscribe.subscribe-embedded-agent-session.subscribeembeddedagentsession.test.ts b/src/agents/embedded-agent-subscribe.subscribe-embedded-agent-session.subscribeembeddedagentsession.test.ts index 131c79169a65..f11f7f7cd8dd 100644 --- a/src/agents/embedded-agent-subscribe.subscribe-embedded-agent-session.subscribeembeddedagentsession.test.ts +++ b/src/agents/embedded-agent-subscribe.subscribe-embedded-agent-session.subscribeembeddedagentsession.test.ts @@ -20,6 +20,7 @@ import { findLifecycleErrorAgentEvent, } from "./embedded-agent-subscribe.e2e-harness.js"; import { subscribeEmbeddedAgentSession } from "./embedded-agent-subscribe.js"; +import { createOpenAiResponsesTextEvent } from "./embedded-agent-subscribe.openai-responses.test-helpers.js"; import { makeZeroUsageSnapshot } from "./usage.js"; const retryingCompactionEnd = () => @@ -1718,5 +1719,354 @@ describe("subscribeEmbeddedAgentSession", () => { summary: "Nothing needs attention.", }); }); + + describe("flushPartialAssistantText", () => { + it("does not commit commentary-phase text on timeout flush", () => { + const { emit, subscription } = createSubscribedHarness({ + runId: "run", + }); + + emit({ type: "message_start", message: { role: "assistant" } }); + // OpenAI Responses commentary items stream text_delta events that the + // normal path deliberately keeps out of reply buffers. The timeout flush + // must preserve that boundary: commentary must not become assistantTexts. + emit( + createOpenAiResponsesTextEvent({ + type: "text_delta", + text: "Working...", + delta: "Working...", + id: "item-commentary", + signaturePhase: "commentary", + partialPhase: "commentary", + }), + ); + + subscription.flushPartialAssistantText(); + + expect(subscription.assistantTexts).toEqual([]); + }); + + it("commits final-answer text that follows a commentary item", () => { + const { emit, subscription } = createSubscribedHarness({ + runId: "run", + }); + + emit({ type: "message_start", message: { role: "assistant" } }); + emit( + createOpenAiResponsesTextEvent({ + type: "text_delta", + text: "Working...", + delta: "Working...", + id: "item-commentary", + signaturePhase: "commentary", + partialPhase: "commentary", + }), + ); + // A later final-answer item resets the buffered item boundary, so the + // timeout flush must preserve the visible final text while dropping the + // preceding commentary bytes. + emit( + createOpenAiResponsesTextEvent({ + type: "text_delta", + text: "Final answer", + delta: "Final answer", + id: "item-final", + signaturePhase: "final_answer", + partialPhase: "final_answer", + }), + ); + + subscription.flushPartialAssistantText(); + + expect(subscription.assistantTexts).toEqual(["Final answer"]); + }); + + it("preserves normal visible text", () => { + const { emit, subscription } = createSubscribedHarness({ + runId: "run", + }); + + emit({ type: "message_start", message: { role: "assistant" } }); + emitAssistantTextDelta(emit, "Hello "); + emitAssistantTextDelta(emit, "world"); + + subscription.flushPartialAssistantText(); + + expect(subscription.assistantTexts).toEqual(["Hello world"]); + }); + + it("strips think tags before committing text", () => { + const { emit, subscription } = createSubscribedHarness({ + runId: "run", + }); + + emit({ type: "message_start", message: { role: "assistant" } }); + emitAssistantTextDelta(emit, "Before"); + emitAssistantTextDelta(emit, " secret"); + emitAssistantTextDelta(emit, "After"); + + subscription.flushPartialAssistantText(); + + expect(subscription.assistantTexts).toEqual(["BeforeAfter"]); + }); + + it("handles final tags matching enforceFinalTag param", () => { + const { emit, subscription } = createSubscribedHarness({ + runId: "run", + enforceFinalTag: true, + }); + + emit({ type: "message_start", message: { role: "assistant" } }); + emitAssistantTextDelta(emit, "Discarded "); + emitAssistantTextDelta(emit, "preserved"); + emitAssistantTextDelta(emit, " also discarded"); + + subscription.flushPartialAssistantText(); + + expect(subscription.assistantTexts).toEqual(["preserved"]); + }); + + it("strips final tags but preserves visible text when enforceFinalTag is disabled", () => { + const { emit, subscription } = createSubscribedHarness({ + runId: "run", + // Default policy: final-tag enforcement is off, so the timeout flush + // must keep the same visible text the normal path would retain and + // only strip the markers themselves. + enforceFinalTag: false, + }); + + emit({ type: "message_start", message: { role: "assistant" } }); + emitAssistantTextDelta(emit, "Discarded "); + emitAssistantTextDelta(emit, "preserved"); + emitAssistantTextDelta(emit, " also kept"); + + subscription.flushPartialAssistantText(); + + // Same normalization as normal completion with enforceFinalTag=false: + // the final-tag markers are stripped, no surrounding visible text is lost. + expect(subscription.assistantTexts).toEqual(["Discarded preserved also kept"]); + }); + + it("strips downgraded tool call text", () => { + const { emit, subscription } = createSubscribedHarness({ + runId: "run", + }); + + emit({ type: "message_start", message: { role: "assistant" } }); + emitAssistantTextDelta(emit, "Visible answer"); + emitAssistantTextDelta(emit, " [Tool Call: some_fn]"); + + subscription.flushPartialAssistantText(); + + expect(subscription.assistantTexts).toEqual(["Visible answer"]); + }); + + it("is a no-op when deltaBuffer is empty", () => { + const { subscription } = createSubscribedHarness({ + runId: "run", + }); + + subscription.flushPartialAssistantText(); + + expect(subscription.assistantTexts).toEqual([]); + }); + + it("preserves visible prefix before unclosed think tag on flush", () => { + const { emit, subscription } = createSubscribedHarness({ + runId: "run", + }); + + emit({ type: "message_start", message: { role: "assistant" } }); + // Streaming path advances state.blockState.thinking to true on , + // then a timeout fires before . flushPartialAssistantText must + // use fresh filter state so "Before " is not treated as hidden content. + emitAssistantTextDelta(emit, "Before "); + emitAssistantTextDelta(emit, " reasoning without close"); + + subscription.flushPartialAssistantText(); + + // The visible prefix is preserved (trimEnd removes trailing space). + expect(subscription.assistantTexts).toEqual(["Before"]); + }); + + it("preserves visible prefix before unclosed final tag on flush", () => { + const { emit, subscription } = createSubscribedHarness({ + runId: "run", + enforceFinalTag: true, + }); + + emit({ type: "message_start", message: { role: "assistant" } }); + // Same boundary: streaming advances state.blockState.final to true + // on , then timeout fires. Flush must preserve text inside + // the unclosed final block and hide text that appeared before . + emitAssistantTextDelta(emit, "Before "); + emitAssistantTextDelta(emit, " content without close"); + + subscription.flushPartialAssistantText(); + + // enforceFinalTag hides text before ; text inside the + // unclosed final block is preserved. + expect(subscription.assistantTexts).toEqual([" content without close"]); + }); + + it("does not re-append text already committed by an earlier flush", () => { + const { emit, subscription } = createSubscribedHarness({ + runId: "run", + }); + + emit({ type: "message_start", message: { role: "assistant" } }); + emitAssistantTextDelta(emit, "Hello world"); + + // Pre-abort flush commits the buffered text. + subscription.flushPartialAssistantText(); + // Post-drain re-flush sees the same buffer (a queued suffix may or may + // not have landed); it must not append the cumulative text again. + subscription.flushPartialAssistantText(); + + expect(subscription.assistantTexts).toEqual(["Hello world"]); + }); + + it("commits only the queued suffix on a second flush", () => { + const { emit, subscription } = createSubscribedHarness({ + runId: "run", + }); + + emit({ type: "message_start", message: { role: "assistant" } }); + emitAssistantTextDelta(emit, "Hello "); + + subscription.flushPartialAssistantText(); + // A message_update serialized behind the abort lands after the first + // flush; the re-flush must append only the new suffix to the same entry + // (never re-append the already-committed prefix). + emitAssistantTextDelta(emit, "world"); + subscription.flushPartialAssistantText(); + + expect(subscription.assistantTexts).toEqual(["Hello world"]); + }); + + it("replaces already-delivered live block chunks with the cumulative text instead of duplicating them", () => { + const onBlockReply = vi.fn(); + const { emit, subscription } = createSubscribedHarness({ + runId: "run", + onBlockReply, + blockReplyChunking: { + minChars: 8, + maxChars: 200, + breakPreference: "sentence", + }, + }); + + emit({ type: "message_start", message: { role: "assistant" } }); + emitAssistantTextDelta(emit, "Hello world. "); + emitAssistantTextDelta(emit, "Next sentence. "); + + // Normal live block streaming already committed each chunk into + // assistantTexts before the deadline; the timeout flush must not append + // the cumulative buffer on top of them (P1: avoid duplicating live block + // chunks during timeout flushing). + expect(subscription.assistantTexts).toEqual(["Hello world.", "Next sentence."]); + + subscription.flushPartialAssistantText(); + + expect(subscription.assistantTexts).toEqual(["Hello world. Next sentence."]); + expect(onBlockReply).toHaveBeenCalled(); + }); + + it("folds a queued suffix into the already-committed live projection without duplicating it", () => { + const onBlockReply = vi.fn(); + const { emit, subscription } = createSubscribedHarness({ + runId: "run", + onBlockReply, + blockReplyChunking: { + minChars: 8, + maxChars: 200, + breakPreference: "sentence", + }, + }); + + emit({ type: "message_start", message: { role: "assistant" } }); + emitAssistantTextDelta(emit, "Hello world. "); + + // Pre-abort flush replaces the live chunk with the buffered projection. + subscription.flushPartialAssistantText(); + expect(subscription.assistantTexts).toEqual(["Hello world."]); + + // A message_update serialized behind the abort lands after the first + // flush; the live path also commits the new chunk. The re-flush must + // reconcile the whole segment instead of appending the suffix twice. + emitAssistantTextDelta(emit, "Next sentence. "); + expect(subscription.assistantTexts).toEqual(["Hello world.", "Next sentence."]); + + subscription.flushPartialAssistantText(); + + expect(subscription.assistantTexts).toEqual(["Hello world. Next sentence."]); + }); + + it("retains hidden-tag context across flushes so a queued suffix inside an unclosed think tag never leaks", () => { + const { emit, subscription } = createSubscribedHarness({ + runId: "run", + }); + + emit({ type: "message_start", message: { role: "assistant" } }); + emitAssistantTextDelta(emit, "Before "); + emitAssistantTextDelta(emit, " reasoning without close"); + + // First flush commits the visible prefix and would have cleared the + // buffer under the previous implementation, losing the opening . + subscription.flushPartialAssistantText(); + expect(subscription.assistantTexts).toEqual(["Before"]); + + // A queued suffix inside the still-open hidden block must stay hidden: + // the retained buffer keeps the opening tag visible to the filter. + emitAssistantTextDelta(emit, "secret continuation"); + subscription.flushPartialAssistantText(); + + expect(subscription.assistantTexts).toEqual(["Before"]); + }); + + it("replaces flushed partial text with the complete text when message_end arrives", () => { + const { emit, subscription } = createSubscribedHarness({ + runId: "run", + }); + + emit({ type: "message_start", message: { role: "assistant" } }); + emitAssistantTextDelta(emit, "Hello"); + subscription.flushPartialAssistantText(); + expect(subscription.assistantTexts).toEqual(["Hello"]); + + // The abort raced a clean completion: message_end finalizes the complete + // text. The flushed partial must be replaced, not duplicated. + emit({ + type: "message_end", + message: { + role: "assistant", + content: [{ type: "text", text: "Hello world" }], + }, + }); + + expect(subscription.assistantTexts).toEqual(["Hello world"]); + }); + + it("replaces a flushed entry when a queued orphan reasoning close retracts the prefix", () => { + const { emit, subscription } = createSubscribedHarness({ + runId: "run", + }); + + emit({ type: "message_start", message: { role: "assistant" } }); + // First flush commits text that the sanitizer still treats as visible: + // the opening reasoning tag has not arrived yet. + emitAssistantTextDelta(emit, "private chain"); + subscription.flushPartialAssistantText(); + expect(subscription.assistantTexts).toEqual(["private chain"]); + + // A queued delta delivers the orphan close plus the real answer. The + // full-buffer re-filter retracts the leaked prefix; the flush must + // REPLACE the stored entry, not extend it (P1: reconcile retractions). + emitAssistantTextDelta(emit, "Visible answer"); + subscription.flushPartialAssistantText(); + + expect(subscription.assistantTexts).toEqual(["Visible answer"]); + }); + }); }); /* oxlint-disable max-lines -- TODO: split this grandfathered oversized file. */ diff --git a/src/agents/embedded-agent-subscribe.ts b/src/agents/embedded-agent-subscribe.ts index 60c8d0b3f414..088b01f20ac6 100644 --- a/src/agents/embedded-agent-subscribe.ts +++ b/src/agents/embedded-agent-subscribe.ts @@ -1,3 +1,4 @@ +import { createInlineCodeState } from "../../packages/markdown-core/src/code-spans.js"; /** * Subscribes to embedded-agent sessions and streams formatted replies/events. */ @@ -34,6 +35,7 @@ import { filterToolResultMediaUrls, } from "./embedded-agent-tool-media.js"; import { buildToolLifecycleErrorResult } from "./embedded-agent-tool-results.js"; +import { stripDowngradedToolCallText } from "./embedded-agent-utils.js"; import type { AgentRunTimeoutPhase } from "./run-timeout-attribution.js"; import type { AgentMessage } from "./runtime/index.js"; import { hasNonzeroUsage, normalizeUsage, type UsageLike } from "./usage.js"; @@ -467,6 +469,42 @@ export function subscribeEmbeddedAgentSession(params: SubscribeEmbeddedAgentSess } }; + // Re-filter the full raw buffer. Reusing live scanner state would hide the + // visible prefix when timeout interrupts an open or block. + const finalizeFlushedAssistantText = (text: string) => + stripDowngradedToolCallText( + stripBlockTags( + text, + { + thinking: false, + final: false, + inlineCode: createInlineCodeState(), + }, + { final: true }, + ), + ).trimEnd(); + + // Settlement calls this only for the final, failure-free run-budget terminal. + // Retain and re-filter the full buffer so queued suffixes keep hidden-tag + // context; replace live chunks instead of appending cumulative text twice. + const flushPartialAssistantText = () => { + const text = state.deltaBuffer; + if (!text) { + return; + } + if (state.deltaBufferIsCommentary) { + state.hasFlushedPartialText = false; + return; + } + const visibleText = finalizeFlushedAssistantText(text); + if (assistantTexts.length > state.assistantTextBaseline || state.hasFlushedPartialText) { + replyDelivery.replaceCurrentAssistantText(visibleText); + } else if (visibleText) { + replyDelivery.pushAssistantText(visibleText); + } + state.hasFlushedPartialText = Boolean(visibleText); + }; + const ctx: EmbeddedAgentSubscribeContext = { params, state, @@ -666,6 +704,7 @@ export function subscribeEmbeddedAgentSession(params: SubscribeEmbeddedAgentSess getLastCompactionTokensAfter: () => state.lastCompactionTokensAfter, getAssistantTurnCount: () => state.assistantTurnCount, waitForPendingEvents: replyDelivery.waitForPendingEvents, + flushPartialAssistantText, getItemLifecycle: () => ({ startedCount: state.itemStartedCount, completedCount: state.itemCompletedCount,