From e082acd0776293f23920b72f7f9f772e12e5d3f5 Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Tue, 11 Aug 2026 03:46:29 -0700 Subject: [PATCH] refactor(agents): collapse stream relay layers (#121978) * refactor(agents): inline stream runtime preparation * refactor(agents): inline stream finalization --- .../run/attempt-execution-phase.test.ts | 259 +++++++++-- .../run/attempt-execution-phase.ts | 227 +++++++--- .../run/attempt-execution-settle.test.ts | 47 +- .../run/attempt-execution-types.ts | 18 +- .../run/attempt-settle.ts | 275 ++++++++---- .../run/attempt-stream-finalize.test.ts | 401 ++++++++++++------ .../run/attempt-stream-finalize.ts | 149 ------- .../attempt-stream-runtime-prepare.test.ts | 267 ------------ .../run/attempt-stream-runtime-prepare.ts | 192 --------- 9 files changed, 876 insertions(+), 959 deletions(-) delete mode 100644 src/agents/embedded-agent-runner/run/attempt-stream-finalize.ts delete mode 100644 src/agents/embedded-agent-runner/run/attempt-stream-runtime-prepare.test.ts delete mode 100644 src/agents/embedded-agent-runner/run/attempt-stream-runtime-prepare.ts diff --git a/src/agents/embedded-agent-runner/run/attempt-execution-phase.test.ts b/src/agents/embedded-agent-runner/run/attempt-execution-phase.test.ts index 32709516c50a..bc1b49466b36 100644 --- a/src/agents/embedded-agent-runner/run/attempt-execution-phase.test.ts +++ b/src/agents/embedded-agent-runner/run/attempt-execution-phase.test.ts @@ -1,39 +1,97 @@ import { beforeEach, describe, expect, it, vi } from "vitest"; const mocks = vi.hoisted(() => ({ - prepareStreamRuntime: vi.fn(), + abortable: vi.fn(), + bindOwnedSessionTranscriptWrites: vi.fn(), + createRunAbort: vi.fn(), + flushPendingToolResultsAfterIdle: vi.fn(), + installStreamGuards: vi.fn(), + prepareHistory: vi.fn(), + prepareStream: vi.fn(), + prepareTimeout: vi.fn(), runSettledPhase: vi.fn(), + withOwnedSessionTranscriptWrites: vi.fn(), })); -vi.mock("./attempt-stream-runtime-prepare.js", () => ({ - prepareEmbeddedAttemptStreamRuntime: mocks.prepareStreamRuntime, +vi.mock("../../../config/sessions/transcript-write-context.js", () => ({ + bindOwnedSessionTranscriptWrites: mocks.bindOwnedSessionTranscriptWrites, + withOwnedSessionTranscriptWrites: mocks.withOwnedSessionTranscriptWrites, +})); +vi.mock("../wait-for-idle-before-flush.js", () => ({ + flushPendingToolResultsAfterIdle: mocks.flushPendingToolResultsAfterIdle, +})); +vi.mock("./abortable.js", () => ({ abortable: mocks.abortable })); +vi.mock("./attempt-finalize.js", () => ({ + createEmbeddedAttemptRunAbort: mocks.createRunAbort, +})); +vi.mock("./attempt-history.js", () => ({ + prepareEmbeddedAttemptHistory: mocks.prepareHistory, })); vi.mock("./attempt-settle.js", () => ({ runEmbeddedAttemptSettledPhase: mocks.runSettledPhase, })); +vi.mock("./attempt-stream-prepare.js", () => ({ + prepareEmbeddedAttemptStream: mocks.prepareStream, +})); +vi.mock("./attempt-stream.js", () => ({ + installEmbeddedAttemptStreamGuards: mocks.installStreamGuards, +})); +vi.mock("./attempt-timeout-prepare.js", () => ({ + prepareEmbeddedAttemptTimeout: mocks.prepareTimeout, +})); import { runEmbeddedAttemptExecutionPhase } from "./attempt-execution-phase.js"; type ExecutionInput = Parameters[0]; -function createFixture() { +function createFixture(options: { aborted?: boolean } = {}) { const order: string[] = []; - const activeSession = { sessionId: "active-session" }; - const sessionManager = { kind: "session-manager" }; + const attemptAbortController = new AbortController(); + if (options.aborted) { + attemptAbortController.abort(new Error("already aborted")); + } + const runAbort = vi.fn(); + const toolSearchCatalogExecutor = vi.fn(); + const subscription = { + isCompacting: vi.fn(() => false), + }; + const queueHandle = { kind: "embedded", runId: "run-1" }; + const streamResult = { + subscription, + queueHandle, + toolSearchCatalogExecutor, + getBeforeAgentFinalizeRevisionReason: vi.fn(), + stopAcceptingSteerMessages: vi.fn(), + }; + const timeoutResult = { + getRunAbortDeadlineAtMs: vi.fn(() => 123), + clearTimers: vi.fn(), + }; + const activeSession = { + agent: { streamFn: vi.fn() }, + dispose: vi.fn(), + isCompacting: false, + messages: [], + prompt: vi.fn(async () => undefined), + sessionId: "active-session", + }; + const sessionManager = {}; const abortActiveSession = vi.fn(async () => undefined); const trackPromptSettlePromise = vi.fn((promise: Promise) => promise); - const toolSearchCatalogExecutor = vi.fn(); + const externalAbortController = { + setRunAbort: vi.fn(() => order.push("set-run-abort")), + setCompactionState: vi.fn(() => order.push("set-compaction-state")), + }; + const prepStages = { mark: vi.fn(() => order.push("stream-ready")) }; + const emitPrepStageSummary = vi.fn(); + const setToolSearchCatalogExecutor = vi.fn(() => order.push("set-catalog")); + const replaySafeTool = { name: "read" }; const result = { messages: [] }; - const preparedStreamRuntime = { stream: { queueHandle: { kind: "embedded" } } }; const state = { beforeAgentRunBlockedBy: undefined, terminal: { kind: "ok" as const }, trajectoryEndRecorded: false, }; - const prepStages = { mark: vi.fn() }; - const emitPrepStageSummary = vi.fn(); - const setToolSearchCatalogExecutor = vi.fn(); - const replaySafeTool = { name: "read" }; const sessionRuntime = { agentSession: { activeSession, @@ -63,13 +121,20 @@ function createFixture() { }, }; const input = { - attempt: { runId: "run-1", sessionId: "session-1" }, + attempt: { + abortSignal: attemptAbortController.signal, + onBlockReply: vi.fn(), + onBlockReplyFlush: vi.fn(), + runId: "run-1", + sessionId: "session-1", + timeoutMs: 30_000, + }, activeContextEngine: { info: { id: "engine" } }, agentDir: "/agent", isRawModelRun: false, resolveActiveContextEnginePluginId: vi.fn(), runAbortController: new AbortController(), - externalAbortController: {}, + externalAbortController, abortState: {}, prepared: { bootstrap: {}, @@ -111,21 +176,46 @@ function createFixture() { }, } as unknown as ExecutionInput; - mocks.prepareStreamRuntime.mockImplementation(async (streamInput) => { - order.push("stream-runtime"); - streamInput.lifecycle.markStreamReady(); - streamInput.lifecycle.markIdleTimedOut(); - streamInput.lifecycle.markExternalAbort(); - streamInput.lifecycle.markTimedOutDuringCompaction(); - streamInput.lifecycle.markTimedOutByRunBudget(); - streamInput.lifecycle.setToolSearchCatalogExecutor(toolSearchCatalogExecutor); - return preparedStreamRuntime; + mocks.abortable.mockImplementation((_signal, promise) => promise); + mocks.bindOwnedSessionTranscriptWrites.mockImplementation((_context, operation) => operation); + mocks.withOwnedSessionTranscriptWrites.mockImplementation( + async (_context, operation) => await operation(), + ); + mocks.installStreamGuards.mockImplementation(() => { + order.push("guards"); + return { + cacheObservabilityEnabled: true, + promptCacheTools: [{ name: "read" }], + }; + }); + mocks.prepareHistory.mockImplementation(async () => { + order.push("history"); + return { + contextEnginePromptAuthority: "assembled", + contextEngineAssemblySucceeded: true, + }; + }); + mocks.createRunAbort.mockImplementation(() => { + order.push("abort"); + return runAbort; + }); + mocks.prepareStream.mockImplementation((streamInput) => { + order.push("stream"); + const idleError = new Error("idle timeout"); + mocks.installStreamGuards.mock.calls[0]?.[0].onIdleTimeout(idleError); + streamInput.markExternalAbort(); + return streamResult; + }); + mocks.prepareTimeout.mockImplementation((timeoutInput) => { + order.push("timeout"); + timeoutInput.markTimedOutDuringCompaction(); + timeoutInput.markTimedOutByRunBudget(); + return timeoutResult; }); mocks.runSettledPhase.mockImplementation(async (settledInput) => { order.push("settled-phase"); expect(settledInput.getRepairedRejectedThinkingReplay()).toBe(false); - const streamInput = mocks.prepareStreamRuntime.mock.calls[0]?.[0]; - streamInput.lifecycle.markRejectedThinkingReplayRepaired(); + mocks.installStreamGuards.mock.calls[0]?.[0].onRejectedThinkingReplayRepaired(); expect(settledInput.getRepairedRejectedThinkingReplay()).toBe(true); return result; }); @@ -134,15 +224,19 @@ function createFixture() { abortActiveSession, activeSession, emitPrepStageSummary, + externalAbortController, input, order, prepStages, - preparedStreamRuntime, replaySafeTool, result, + runAbort, sessionManager, setToolSearchCatalogExecutor, state, + streamResult, + subscription, + timeoutResult, toolSearchCatalogExecutor, trackPromptSettlePromise, }; @@ -153,13 +247,24 @@ beforeEach(() => { }); describe("runEmbeddedAttemptExecutionPhase", () => { - it("prepares the guarded stream and delegates settlement with live lifecycle state", async () => { + it("prepares guarded history, stream handling, deadlines, and settlement in order", async () => { const fixture = createFixture(); const result = await runEmbeddedAttemptExecutionPhase(fixture.input); expect(result).toBe(fixture.result); - expect(fixture.order).toEqual(["stream-runtime", "settled-phase"]); + expect(fixture.order).toEqual([ + "guards", + "stream-ready", + "history", + "abort", + "set-run-abort", + "stream", + "set-catalog", + "set-compaction-state", + "timeout", + "settled-phase", + ]); expect(fixture.state).toEqual( expect.objectContaining({ terminal: { @@ -175,37 +280,111 @@ describe("runEmbeddedAttemptExecutionPhase", () => { expect(fixture.setToolSearchCatalogExecutor).toHaveBeenCalledWith( fixture.toolSearchCatalogExecutor, ); - expect(mocks.runSettledPhase).toHaveBeenCalledWith( + + const settledInput = mocks.runSettledPhase.mock.calls[0]?.[0]; + expect(settledInput).toEqual( expect.objectContaining({ getRepairedRejectedThinkingReplay: expect.any(Function), - preparedStreamRuntime: fixture.preparedStreamRuntime, + preparedStreamRuntime: expect.objectContaining({ + cache: { + observabilityEnabled: true, + promptTools: [{ name: "read" }], + }, + history: expect.objectContaining({ contextEngineAssemblySucceeded: true }), + isProbeSession: false, + stream: fixture.streamResult, + timeout: fixture.timeoutResult, + }), }), ); - const settledInput = mocks.runSettledPhase.mock.calls[0]?.[0]; expect(settledInput.getRepairedRejectedThinkingReplay()).toBe(true); - const streamInput = mocks.prepareStreamRuntime.mock.calls[0]?.[0]; - expect(streamInput).toEqual( + const guardInput = mocks.installStreamGuards.mock.calls[0]?.[0]; + expect(guardInput).toEqual( expect.objectContaining({ - activeSession: fixture.activeSession, + attempt: fixture.input.attempt, + session: fixture.activeSession, sessionManager: fixture.sessionManager, - abortActiveSession: fixture.abortActiveSession, - trackPromptSettlePromise: fixture.trackPromptSettlePromise, }), ); - expect(streamInput.lifecycle.isYieldDetected()).toBe(true); - expect(streamInput.lifecycle.readRunState()).toEqual({ + expect(guardInput.isYieldDetected()).toBe(true); + expect(fixture.runAbort).toHaveBeenCalledWith(true, expect.any(Error)); + + const abortInput = mocks.createRunAbort.mock.calls[0]?.[0]; + expect(abortInput.abortActiveSession).toBe(fixture.abortActiveSession); + const streamInput = mocks.prepareStream.mock.calls[0]?.[0]; + expect(streamInput.activeSession).toBe(fixture.activeSession); + expect(streamInput.getRunState()).toEqual({ aborted: true, promptError: null, timedOut: true, yieldDetected: true, }); - expect(streamInput.stream.isReplaySafeTool(fixture.replaySafeTool)).toBe(true); + expect(streamInput.isReplaySafeTool(fixture.replaySafeTool)).toBe(true); + expect(fixture.externalAbortController.setCompactionState).toHaveBeenCalledWith({ + isPendingOrRetrying: fixture.subscription.isCompacting, + isInFlight: expect.any(Function), + }); + expect(mocks.prepareTimeout).toHaveBeenCalledWith( + expect.objectContaining({ + abortRun: fixture.runAbort, + compactionState: fixture.subscription, + }), + ); + + await settledInput.preparedStreamRuntime.promptActiveSession("hello"); + expect(fixture.activeSession.prompt).toHaveBeenCalledWith("hello", undefined); + expect(fixture.trackPromptSettlePromise).toHaveBeenCalledOnce(); + expect(mocks.withOwnedSessionTranscriptWrites).toHaveBeenCalledOnce(); + }); + + it.each([ + { label: "external cancellation", message: "run cancelled" }, + { label: "run timeout", message: "run timed out" }, + ])("does not start a prompt after $label", async ({ message }) => { + const fixture = createFixture(); + await runEmbeddedAttemptExecutionPhase(fixture.input); + const reason = new Error(message); + const abortError = new Error(message, { cause: reason }); + abortError.name = "AbortError"; + fixture.input.runAbortController.abort(reason); + mocks.abortable.mockImplementationOnce((_signal, _promise) => Promise.reject(abortError)); + const settledInput = mocks.runSettledPhase.mock.calls[0]?.[0]; + + await expect( + settledInput.preparedStreamRuntime.promptActiveSession("must not start"), + ).rejects.toBe(abortError); + + expect(fixture.activeSession.prompt).not.toHaveBeenCalled(); + expect(fixture.trackPromptSettlePromise).not.toHaveBeenCalled(); + expect(mocks.abortable).toHaveBeenCalledOnce(); + }); + + it("flushes pending tool results and disposes the session when history preparation fails", async () => { + const fixture = createFixture({ aborted: true }); + const failure = new Error("history failed"); + mocks.prepareHistory.mockRejectedValueOnce(failure); + mocks.flushPendingToolResultsAfterIdle.mockResolvedValue(undefined); + + await expect(runEmbeddedAttemptExecutionPhase(fixture.input)).rejects.toBe(failure); + + expect(mocks.flushPendingToolResultsAfterIdle).toHaveBeenCalledWith({ + agent: fixture.activeSession.agent, + sessionManager: fixture.sessionManager, + timeoutMs: 0, + }); + expect(fixture.activeSession.dispose).toHaveBeenCalledOnce(); + expect(mocks.createRunAbort).not.toHaveBeenCalled(); + expect(mocks.prepareStream).not.toHaveBeenCalled(); + expect(mocks.prepareTimeout).not.toHaveBeenCalled(); + expect(mocks.runSettledPhase).not.toHaveBeenCalled(); }); it("does not enter settlement when stream preparation fails", async () => { const fixture = createFixture(); - mocks.prepareStreamRuntime.mockRejectedValueOnce(new Error("stream setup failed")); + mocks.prepareStream.mockImplementationOnce(() => { + throw new Error("stream setup failed"); + }); await expect(runEmbeddedAttemptExecutionPhase(fixture.input)).rejects.toThrow( "stream setup failed", diff --git a/src/agents/embedded-agent-runner/run/attempt-execution-phase.ts b/src/agents/embedded-agent-runner/run/attempt-execution-phase.ts index 2b74f2acae6f..656c9380381b 100644 --- a/src/agents/embedded-agent-runner/run/attempt-execution-phase.ts +++ b/src/agents/embedded-agent-runner/run/attempt-execution-phase.ts @@ -1,12 +1,24 @@ /** Prepares the guarded stream runtime before prompt execution and settlement. */ +import { + bindOwnedSessionTranscriptWrites, + withOwnedSessionTranscriptWrites, +} from "../../../config/sessions/transcript-write-context.js"; import { mergeAgentRunAttemptTerminal, projectAgentRunAttemptTerminal, type AgentRunAttemptTerminal, } from "../../agent-run-terminal-outcome.js"; +import { log } from "../logger.js"; +import type { EmbeddedAgentQueueHandle } from "../runs.js"; +import { flushPendingToolResultsAfterIdle } from "../wait-for-idle-before-flush.js"; +import { abortable as abortableWithSignal } from "./abortable.js"; import type { EmbeddedAttemptExecutionPhaseInput } from "./attempt-execution-types.js"; +import { createEmbeddedAttemptRunAbort } from "./attempt-finalize.js"; +import { prepareEmbeddedAttemptHistory } from "./attempt-history.js"; import { runEmbeddedAttemptSettledPhase } from "./attempt-settle.js"; -import { prepareEmbeddedAttemptStreamRuntime } from "./attempt-stream-runtime-prepare.js"; +import { prepareEmbeddedAttemptStream } from "./attempt-stream-prepare.js"; +import { installEmbeddedAttemptStreamGuards } from "./attempt-stream.js"; +import { prepareEmbeddedAttemptTimeout } from "./attempt-timeout-prepare.js"; import type { EmbeddedRunAttemptResult } from "./types.js"; export async function runEmbeddedAttemptExecutionPhase( @@ -49,33 +61,40 @@ export async function runEmbeddedAttemptExecutionPhase( state.terminal = mergeAgentRunAttemptTerminal(state.terminal, incoming); }; - const preparedStreamRuntime = await prepareEmbeddedAttemptStreamRuntime({ + const idleTimeoutTriggerRef: { current?: (error: Error) => void } = {}; + const { cacheObservabilityEnabled, promptCacheTools } = installEmbeddedAttemptStreamGuards({ attempt, - activeSession, + session: activeSession, sessionManager, - ownedTranscriptWriteContext: input.sessionLock.ownedTranscriptWriteContext, - runAbortController: input.runAbortController, - externalAbortController: input.externalAbortController, - abortActiveSession, - abortState: input.abortState, - trackPromptSettlePromise, - compactionTimeoutMs: input.sessionLock.compactionTimeoutMs, - guards: { - sessionAgentId: input.setup.sessionAgentId, - cacheTrace, - allCustomTools, - systemPromptText: sessionRuntimeState.systemPromptText, - transcriptPolicy, - isOpenAIResponsesApi, - replayAllowedToolNames, - liveAllowedToolNames, - clientToolLoopDetection, - anthropicPayloadLogger, - effectiveAgentTransport, - providerTextTransforms, - runTrace: input.diagnostics.runTrace, + sessionAgentId: input.setup.sessionAgentId, + cacheTrace, + allCustomTools, + systemPromptText: sessionRuntimeState.systemPromptText, + transcriptPolicy, + isOpenAIResponsesApi, + replayAllowedToolNames, + liveAllowedToolNames, + clientToolLoopDetection, + anthropicPayloadLogger, + effectiveAgentTransport, + providerTextTransforms, + runTrace: input.diagnostics.runTrace, + isYieldDetected: () => input.lifecycle.readYieldState().yieldDetected, + onRejectedThinkingReplayRepaired: () => { + repairedRejectedThinkingReplay = true; }, - history: { + onIdleTimeout: (error) => idleTimeoutTriggerRef.current?.(error), + abortSignal: input.runAbortController.signal, + }); + input.setup.prepStages.mark("stream-setup"); + input.setup.emitPrepStageSummary("stream-ready"); + + let preparedHistory: Awaited>; + try { + preparedHistory = await prepareEmbeddedAttemptHistory({ + attempt, + activeSession, + sessionManager, ...(input.activeContextEngine ? { activeContextEngine: input.activeContextEngine } : {}), cacheTrace, capabilityToolNames, @@ -90,48 +109,124 @@ export async function runEmbeddedAttemptExecutionPhase( systemPromptText: sessionRuntimeState.systemPromptText, transcriptPolicy, setActiveSessionSystemPrompt, - }, - stream: { - runtimeChannel, - hookRunner, - hookAgentId, - diagnosticTrace: input.diagnostics.diagnosticTrace, - clientToolCallSlots, - toolSearchTargetTranscriptProjections, - isReplaySafeTool: (tool) => replaySafeTools.has(tool as never), - hasDeliveredSourceReply, - markSourceReplyDelivered, - sandboxSessionKey: input.setup.sandboxSessionKey, - builtinToolNames, - replaySafeToolNames, - }, - lifecycle: { - isYieldDetected: () => input.lifecycle.readYieldState().yieldDetected, - markRejectedThinkingReplayRepaired: () => { - repairedRejectedThinkingReplay = true; - }, - markStreamReady: () => { - input.setup.prepStages.mark("stream-setup"); - input.setup.emitPrepStageSummary("stream-ready"); - }, - markIdleTimedOut: () => mergeTerminal({ kind: "timeout", phase: "prompt", source: "idle" }), - markExternalAbort: () => mergeTerminal({ kind: "aborted", source: "external" }), - markTimedOutDuringCompaction: () => - mergeTerminal({ kind: "timeout", phase: "compaction", source: "observation" }), - markTimedOutByRunBudget: () => - mergeTerminal({ kind: "timeout", phase: "prompt", source: "run_budget" }), - readRunState: () => { - const terminal = projectAgentRunAttemptTerminal(state.terminal); - return { - aborted: terminal.aborted, - promptError: terminal.promptError, - timedOut: terminal.timedOut, - yieldDetected: input.lifecycle.readYieldState().yieldDetected, - }; - }, - setToolSearchCatalogExecutor: input.lifecycle.setToolSearchCatalogExecutor, - }, + }); + } catch (error) { + await flushPendingToolResultsAfterIdle({ + agent: activeSession.agent, + sessionManager, + // An already-aborted setup must dispose immediately without orphaning tool calls. + ...(attempt.abortSignal?.aborted ? { timeoutMs: 0 } : {}), + }); + activeSession.dispose(); + throw error; + } + + const isProbeSession = attempt.sessionId?.startsWith("probe-") ?? false; + const queueHandleRef: { current?: EmbeddedAgentQueueHandle } = {}; + const abortRun = createEmbeddedAttemptRunAbort({ + abortActiveSession, + activeSession, + attempt, + getQueueHandle: () => queueHandleRef.current, + isProbeSession, + log, + runAbortController: input.runAbortController, + state: input.abortState, }); + input.externalAbortController.setRunAbort(abortRun); + idleTimeoutTriggerRef.current = (error) => { + mergeTerminal({ kind: "timeout", phase: "prompt", source: "idle" }); + abortRun(true, error); + }; + const abortable = (promise: Promise): Promise => + abortableWithSignal(input.runAbortController.signal, promise); + const promptActiveSession = ( + prompt: string, + options?: Parameters[1], + ): Promise => + withOwnedSessionTranscriptWrites(input.sessionLock.ownedTranscriptWriteContext, async () => { + // Prompting starts its own agent loop; reject before creating a loop that + // an already-aborted attempt can no longer cancel. + if (input.runAbortController.signal.aborted) { + return abortable(Promise.resolve()); + } + return abortable(trackPromptSettlePromise(activeSession.prompt(prompt, options))); + }); + const onBlockReply = attempt.onBlockReply + ? bindOwnedSessionTranscriptWrites( + input.sessionLock.ownedTranscriptWriteContext, + attempt.onBlockReply, + ) + : undefined; + const onBlockReplyFlush = attempt.onBlockReplyFlush + ? bindOwnedSessionTranscriptWrites( + input.sessionLock.ownedTranscriptWriteContext, + attempt.onBlockReplyFlush, + ) + : undefined; + const preparedStream = prepareEmbeddedAttemptStream({ + attempt, + activeSession, + runAbortController: input.runAbortController, + abortRun, + markExternalAbort: () => mergeTerminal({ kind: "aborted", source: "external" }), + getRunState: () => { + const terminal = projectAgentRunAttemptTerminal(state.terminal); + return { + aborted: terminal.aborted, + promptError: terminal.promptError, + timedOut: terminal.timedOut, + yieldDetected: input.lifecycle.readYieldState().yieldDetected, + }; + }, + onBlockReply, + onBlockReplyFlush, + runtimeChannel, + hookRunner, + hookAgentId, + diagnosticTrace: input.diagnostics.diagnosticTrace, + clientToolCallSlots, + toolSearchTargetTranscriptProjections, + isReplaySafeTool: (tool) => replaySafeTools.has(tool as never), + hasDeliveredSourceReply, + markSourceReplyDelivered, + sandboxSessionKey: input.setup.sandboxSessionKey, + builtinToolNames, + replaySafeToolNames, + }); + input.lifecycle.setToolSearchCatalogExecutor(preparedStream.toolSearchCatalogExecutor); + input.externalAbortController.setCompactionState({ + isPendingOrRetrying: preparedStream.subscription.isCompacting, + isInFlight: () => activeSession.isCompacting, + }); + queueHandleRef.current = preparedStream.queueHandle; + + const attemptTimeout = prepareEmbeddedAttemptTimeout({ + attempt, + activeSession, + compactionState: preparedStream.subscription, + compactionTimeoutMs: input.sessionLock.compactionTimeoutMs, + isProbeSession, + abortRun, + markTimedOutDuringCompaction: () => + mergeTerminal({ kind: "timeout", phase: "compaction", source: "observation" }), + markTimedOutByRunBudget: () => + mergeTerminal({ kind: "timeout", phase: "prompt", source: "run_budget" }), + }); + + const preparedStreamRuntime = { + abortable, + cache: { + observabilityEnabled: cacheObservabilityEnabled, + promptTools: promptCacheTools, + }, + history: preparedHistory, + isProbeSession, + onBlockReplyFlush, + promptActiveSession, + stream: preparedStream, + timeout: attemptTimeout, + }; return await runEmbeddedAttemptSettledPhase({ ...input, preparedStreamRuntime, diff --git a/src/agents/embedded-agent-runner/run/attempt-execution-settle.test.ts b/src/agents/embedded-agent-runner/run/attempt-execution-settle.test.ts index 8f7a9f49cf99..ee61ac4c8893 100644 --- a/src/agents/embedded-agent-runner/run/attempt-execution-settle.test.ts +++ b/src/agents/embedded-agent-runner/run/attempt-execution-settle.test.ts @@ -2,16 +2,18 @@ import { beforeEach, describe, expect, it, vi } from "vitest"; const mocks = vi.hoisted(() => ({ clearActiveEmbeddedRun: vi.fn(), + completeAfterTurn: vi.fn(), completeResult: vi.fn(), - finalizeStream: vi.fn(), logDebug: vi.fn(), logError: vi.fn(), + logWarn: vi.fn(), settleRequesterAfterSessionSpawns: vi.fn(), + settleStream: vi.fn(), runPrompt: vi.fn(), })); vi.mock("../logger.js", () => ({ - log: { debug: mocks.logDebug, error: mocks.logError }, + log: { debug: mocks.logDebug, error: mocks.logError, warn: mocks.logWarn }, })); vi.mock("../../subagents/registry/subagent-registry.js", () => ({ settleRequesterAfterSessionSpawns: mocks.settleRequesterAfterSessionSpawns, @@ -23,8 +25,15 @@ vi.mock("./attempt-prompt-phase.js", () => ({ vi.mock("./attempt-result.js", () => ({ completeEmbeddedAttemptResult: mocks.completeResult, })); -vi.mock("./attempt-stream-finalize.js", () => ({ - finalizeEmbeddedAttemptStreamPhase: mocks.finalizeStream, +vi.mock("./attempt-finalize.js", async (importOriginal) => { + const actual = await importOriginal(); + return { + ...actual, + completeEmbeddedAttemptAfterTurn: mocks.completeAfterTurn, + }; +}); +vi.mock("./attempt-stream-settle.js", () => ({ + settleEmbeddedAttemptStream: mocks.settleStream, })); import { SESSIONS_YIELD_ABORT_REASON } from "./attempt-sessions-yield.js"; @@ -41,12 +50,17 @@ function createFixture() { const detachBackend = vi.fn(() => order.push("detach-backend")); const clearTimers = vi.fn(() => order.push("clear-timers")); const getBeforeAgentFinalizeRevisionReason = vi.fn(() => "revision"); + const getBeforeAgentFinalizeRevisionEntryId = vi.fn(() => undefined); const promptActiveSession = vi.fn(async () => undefined); const activeSession = { + agent: { state: { messages: [] } }, sessionId: "active-session", getActiveToolNames: vi.fn(() => ["read"]), }; - const sessionManager = { kind: "session-manager" }; + const sessionManager = { + kind: "session-manager", + buildSessionContext: vi.fn(() => ({ messages: [] })), + }; const hookRunner = { kind: "hook-runner" }; const cacheTrace = { kind: "cache-trace" }; const trajectoryRecorder = { kind: "trajectory" }; @@ -82,6 +96,7 @@ function createFixture() { queueHandle, stopAcceptingSteerMessages: vi.fn(), getBeforeAgentFinalizeRevisionReason, + getBeforeAgentFinalizeRevisionEntryId, }, timeout: { getRunAbortDeadlineAtMs: vi.fn(() => 123), @@ -196,9 +211,9 @@ function createFixture() { promptInput.lifecycle.markBeforeAgentRunBlocked({ blockedBy: "before_agent" }); return { promptStartedAt: 100 }; }); - mocks.finalizeStream.mockImplementation(async (finalizeInput) => { + mocks.settleStream.mockImplementation(async () => { order.push("finalize"); - finalizeInput.onSettled({ + return { promptError: null, promptErrorSource: null, timedOutDuringCompaction: false, @@ -206,10 +221,15 @@ function createFixture() { sessionIdUsed: "settled-session", lastAssistant: { role: "assistant", content: "done" }, currentAttemptAssistant: { role: "assistant", content: "done" }, + currentAttemptCompletedAssistant: undefined, attemptUsage: { input: 1, output: 2, total: 3 }, cacheBreak: null, promptCache: { cacheRead: 1 }, - }); + lastCallUsage: undefined, + compactionOccurredThisAttempt: false, + }; + }); + mocks.completeAfterTurn.mockImplementation(async () => { return { sessionIdUsed: "final-session", sessionFileUsed: "/tmp/final.jsonl" }; }); mocks.completeResult.mockImplementation(() => { @@ -318,7 +338,7 @@ describe("runEmbeddedAttemptSettledPhase", () => { await expect(runEmbeddedAttemptSettledPhase(fixture.input)).rejects.toBe(failure); - expect(mocks.finalizeStream).not.toHaveBeenCalled(); + expect(mocks.settleStream).not.toHaveBeenCalled(); expect(mocks.completeResult).not.toHaveBeenCalled(); expect(fixture.clearTimers).toHaveBeenCalledOnce(); expect(fixture.detachBackend).toHaveBeenCalledWith(fixture.queueHandle); @@ -435,8 +455,8 @@ describe("runEmbeddedAttemptSettledPhase", () => { it("defaults a source-less settlement failure without dropping it", async () => { const fixture = createFixture(); const failure = new Error("settlement failed"); - mocks.finalizeStream.mockImplementationOnce(async (finalizeInput) => { - finalizeInput.onSettled({ + mocks.settleStream.mockImplementationOnce(async () => { + return { promptError: failure, promptErrorSource: null, timedOutDuringCompaction: true, @@ -447,8 +467,9 @@ describe("runEmbeddedAttemptSettledPhase", () => { attemptUsage: undefined, cacheBreak: null, promptCache: undefined, - }); - return { sessionIdUsed: "settled-session" }; + lastCallUsage: undefined, + compactionOccurredThisAttempt: false, + }; }); await runEmbeddedAttemptSettledPhase(fixture.input); diff --git a/src/agents/embedded-agent-runner/run/attempt-execution-types.ts b/src/agents/embedded-agent-runner/run/attempt-execution-types.ts index 7ef514d24cf4..4ad0ecd39026 100644 --- a/src/agents/embedded-agent-runner/run/attempt-execution-types.ts +++ b/src/agents/embedded-agent-runner/run/attempt-execution-types.ts @@ -6,9 +6,11 @@ import type { createEmbeddedAttemptExternalAbortController, EmbeddedAttemptAbortStatePort, } from "./attempt-finalize.js"; +import type { prepareEmbeddedAttemptHistory } from "./attempt-history.js"; import type { prepareEmbeddedAttemptSessionRuntime } from "./attempt-session-runtime-prepare.js"; import type { prepareEmbeddedAttemptSetup } from "./attempt-setup.js"; -import type { prepareEmbeddedAttemptStreamRuntime } from "./attempt-stream-runtime-prepare.js"; +import type { prepareEmbeddedAttemptStream } from "./attempt-stream-prepare.js"; +import type { installEmbeddedAttemptStreamGuards } from "./attempt-stream.js"; import type { prepareEmbeddedAttemptSystemPrompt } from "./attempt-system-prompt-prepare.js"; import type { prepareEmbeddedAttemptToolCatalog } from "./attempt-tool-catalog.js"; import type { prepareEmbeddedAttemptToolBase } from "./attempt-tool-prepare.js"; @@ -18,8 +20,10 @@ import type { EmbeddedRunAttemptParams } from "./types.js"; type Prepared unknown> = Awaited>; type PreparedSetup = Prepared; type PreparedTranscriptLifecycle = Prepared; -type StreamRuntimeInput = Parameters[0]; -type AttemptContextEngine = NonNullable; +type HistoryInput = Parameters[0]; +type StreamInput = Parameters[0]; +type StreamGuardInput = Parameters[0]; +type AttemptContextEngine = NonNullable; export type EmbeddedAttemptExecutionState = { beforeAgentRunBlockedBy: string | undefined; @@ -62,8 +66,8 @@ export type EmbeddedAttemptExecutionPhaseInput = { | "sessionAgentId" >; diagnostics: { - diagnosticTrace: StreamRuntimeInput["stream"]["diagnosticTrace"]; - runTrace: StreamRuntimeInput["guards"]["runTrace"]; + diagnosticTrace: StreamInput["diagnosticTrace"]; + runTrace: StreamGuardInput["runTrace"]; }; state: EmbeddedAttemptExecutionState; lifecycle: { @@ -72,6 +76,8 @@ export type EmbeddedAttemptExecutionPhaseInput = { yieldDetected: boolean; yieldMessage: string | null; }; - setToolSearchCatalogExecutor: StreamRuntimeInput["lifecycle"]["setToolSearchCatalogExecutor"]; + setToolSearchCatalogExecutor: ( + executor: ReturnType["toolSearchCatalogExecutor"], + ) => void; }; }; diff --git a/src/agents/embedded-agent-runner/run/attempt-settle.ts b/src/agents/embedded-agent-runner/run/attempt-settle.ts index 010aac604a43..0b30cb9fa2ac 100644 --- a/src/agents/embedded-agent-runner/run/attempt-settle.ts +++ b/src/agents/embedded-agent-runner/run/attempt-settle.ts @@ -15,19 +15,41 @@ 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 type { EmbeddedAttemptExecutionPhaseInput, EmbeddedAttemptExecutionState, } from "./attempt-execution-types.js"; +import { completeEmbeddedAttemptAfterTurn } from "./attempt-finalize.js"; +import type { prepareEmbeddedAttemptHistory } from "./attempt-history.js"; import { runEmbeddedAttemptPromptPhase } from "./attempt-prompt-phase.js"; import { completeEmbeddedAttemptResult } from "./attempt-result.js"; -import { finalizeEmbeddedAttemptStreamPhase } from "./attempt-stream-finalize.js"; -import type { prepareEmbeddedAttemptStreamRuntime } from "./attempt-stream-runtime-prepare.js"; +import type { prepareEmbeddedAttemptStream } from "./attempt-stream-prepare.js"; +import { settleEmbeddedAttemptStream } from "./attempt-stream-settle.js"; +import type { installEmbeddedAttemptStreamGuards } from "./attempt-stream.js"; +import type { prepareEmbeddedAttemptTimeout } from "./attempt-timeout-prepare.js"; import type { EmbeddedRunAttemptParams, EmbeddedRunAttemptResult } from "./types.js"; /** Runs prompt dispatch, stream settlement, cleanup, and result projection. */ -type PreparedStreamRuntime = Awaited>; +type PreparedStreamRuntime = { + abortable: (promise: Promise) => Promise; + cache: { + observabilityEnabled: boolean; + promptTools: ReturnType["promptCacheTools"]; + }; + history: Awaited>; + isProbeSession: boolean; + onBlockReplyFlush: Parameters[0]["onBlockReplyFlush"]; + promptActiveSession: ( + prompt: string, + options?: Parameters< + Parameters[0]["activeSession"]["prompt"] + >[1], + ) => Promise; + stream: ReturnType; + timeout: ReturnType; +}; type StreamCleanupInput = { attempt: EmbeddedRunAttemptParams; @@ -307,106 +329,171 @@ export async function runEmbeddedAttemptSettledPhase( }, }); - const afterTurn = await finalizeEmbeddedAttemptStreamPhase({ + // 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}`, + ); + }, + }); + const beforeAgentFinalizeRevisionReason = getBeforeAgentFinalizeRevisionReason(); + const beforeAgentFinalizeRevisionEntryId = getBeforeAgentFinalizeRevisionEntryId(); + let rewoundBeforeAgentFinalizeRevision = false; + if (beforeAgentFinalizeRevisionReason && beforeAgentFinalizeRevisionEntryId) { + await input.sessionLock.withOwnedTranscriptWrite(() => { + const rejectedEntry = sessionManager.getEntry(beforeAgentFinalizeRevisionEntryId); + if (rejectedEntry?.type !== "message" || rejectedEntry.message.role !== "assistant") { + throw new Error( + `before_agent_finalize persisted assistant entry is missing or invalid ` + + `(entry=${beforeAgentFinalizeRevisionEntryId})`, + ); + } + // Keep persistence append-only while excluding the rejected draft and + // every trailing descendant from the hidden retry's active branch. + sessionManager.appendLeafControl({ + targetId: rejectedEntry.parentId, + appendParentId: rejectedEntry.parentId, + }); + rewoundBeforeAgentFinalizeRevision = true; + }); + } + let settledStream: Awaited>; + try { + if (input.getRepairedRejectedThinkingReplay() && !rewoundBeforeAgentFinalizeRevision) { + activeSession.agent.state.messages = sessionManager.buildSessionContext().messages; + } + const settleTerminal = readTerminal(); + const streamSettleState = { + promptError: settleTerminal.promptError, + promptErrorSource: settleTerminal.promptErrorSource, + yieldAborted, + sessionIdUsed, + }; + try { + settledStream = await settleEmbeddedAttemptStream({ + attempt, + activeSession, + sessionManager, + withOwnedTranscriptWrite: input.sessionLock.withOwnedTranscriptWrite, + state: streamSettleState, + runAbortDeadlineAtMs: getRunAbortDeadlineAtMs(), + shouldFlushForContextEngine: Boolean( + input.activeContextEngine && !getBeforeAgentFinalizeRevisionReason(), + ), + subscription, + readLifecycleState: () => { + const terminal = readTerminal(); + return { + aborted: terminal.aborted, + timedOut: terminal.timedOut, + timedOutDuringCompaction: terminal.timedOutDuringCompaction, + }; + }, + markTimedOutDuringCompaction: () => { + state.terminal = mergeAgentRunAttemptTerminal(state.terminal, { + kind: "timeout", + phase: "compaction", + source: "observation", + }); + }, + runAbortSignal: input.runAbortController.signal, + isProbeSession, + onBlockReplyFlush, + abortable, + prePromptMessageCount: sessionRuntimeState.prePromptMessageCount, + toolSearchTargetTranscriptProjections, + cache: { + observabilityEnabled: cacheObservabilityEnabled, + changesForTurn: promptCacheChangesForTurn, + retention: effectivePromptCacheRetention, + }, + }); + } catch (error) { + // Settlement mutates this shared state before some failures. Publish it so + // outer teardown keeps the recorded prompt error and attribution. + setFailure(streamSettleState.promptError, streamSettleState.promptErrorSource); + throw error; + } + } finally { + if (rewoundBeforeAgentFinalizeRevision) { + await input.sessionLock.withOwnedTranscriptWrite(() => { + // Settlement classifies the completed attempt from its original + // in-memory messages. Later work always sees the rewound branch. + activeSession.agent.state.messages = sessionManager.buildSessionContext().messages; + }); + } + } + // Publish settled fields before after-turn hooks: those hooks may throw, and + // outer teardown still needs the completed stream snapshot and usage state. + setFailure(settledStream.promptError, settledStream.promptErrorSource); + if (settledStream.timedOutDuringCompaction) { + state.terminal = mergeAgentRunAttemptTerminal(state.terminal, { + kind: "timeout", + phase: "compaction", + source: "observation", + }); + } + messagesSnapshot = settledStream.messagesSnapshot; + sessionIdUsed = settledStream.sessionIdUsed; + lastAssistant = settledStream.lastAssistant; + currentAttemptAssistant = settledStream.currentAttemptAssistant; + currentAttemptCompletedAssistant = settledStream.currentAttemptCompletedAssistant; + attemptUsage = settledStream.attemptUsage; + cacheBreak = settledStream.cacheBreak; + sessionRuntimeState.promptCache = settledStream.promptCache; + + const afterTurn = await completeEmbeddedAttemptAfterTurn({ attempt, activeSession, sessionManager, withOwnedTranscriptWrite: input.sessionLock.withOwnedTranscriptWrite, - waitForPendingEvents, - repairedRejectedThinkingReplay: input.getRepairedRejectedThinkingReplay(), - getRunAbortDeadlineAtMs, - shouldFlushForContextEngine: () => - Boolean(input.activeContextEngine && !getBeforeAgentFinalizeRevisionReason()), - getBeforeAgentFinalizeRevisionReason, - getBeforeAgentFinalizeRevisionEntryId, - getContextEngineAfterTurnCheckpoint: contextGuards.getAfterTurnCheckpoint, - onSettleErrorState: (settleState) => { - setFailure(settleState.promptError, settleState.promptErrorSource); - }, - onSettled: (settledStream) => { - setFailure(settledStream.promptError, settledStream.promptErrorSource); - if (settledStream.timedOutDuringCompaction) { - state.terminal = mergeAgentRunAttemptTerminal(state.terminal, { - kind: "timeout", - phase: "compaction", - source: "observation", - }); - } - messagesSnapshot = settledStream.messagesSnapshot; - sessionIdUsed = settledStream.sessionIdUsed; - lastAssistant = settledStream.lastAssistant; - currentAttemptAssistant = settledStream.currentAttemptAssistant; - currentAttemptCompletedAssistant = settledStream.currentAttemptCompletedAssistant; - attemptUsage = settledStream.attemptUsage; - cacheBreak = settledStream.cacheBreak; - sessionRuntimeState.promptCache = settledStream.promptCache; - }, - getState: () => { + activeContextEngine: input.activeContextEngine, + readLifecycleState: () => { const terminal = readTerminal(); return { - promptError: terminal.promptError, - promptErrorSource: terminal.promptErrorSource, - yieldAborted, - sessionIdUsed, - sessionFileUsed, + aborted: terminal.aborted, + timedOut: terminal.timedOut, + idleTimedOut: terminal.idleTimedOut, + timedOutDuringCompaction: terminal.timedOutDuringCompaction, }; }, - settle: { - subscription, - readLifecycleState: () => { - const terminal = readTerminal(); - return { - aborted: terminal.aborted, - timedOut: terminal.timedOut, - timedOutDuringCompaction: terminal.timedOutDuringCompaction, - }; - }, - markTimedOutDuringCompaction: () => { - state.terminal = mergeAgentRunAttemptTerminal(state.terminal, { - kind: "timeout", - phase: "compaction", - source: "observation", - }); - }, - runAbortSignal: input.runAbortController.signal, - isProbeSession, - onBlockReplyFlush, - abortable, - prePromptMessageCount: sessionRuntimeState.prePromptMessageCount, - toolSearchTargetTranscriptProjections, - cache: { - observabilityEnabled: cacheObservabilityEnabled, - changesForTurn: promptCacheChangesForTurn, - retention: effectivePromptCacheRetention, - }, + runtime: { + effectiveWorkspace: input.setup.effectiveWorkspace, + agentDir: input.agentDir, + sessionAgentId: input.setup.sessionAgentId, + resolveActiveContextEnginePluginId: input.resolveActiveContextEnginePluginId, + shouldRecordCompletedBootstrapTurn, + cacheTrace, + anthropicPayloadLogger, + hookAgentId, + diagnosticTrace: input.diagnostics.diagnosticTrace, + skillWorkshopAvailable: uncompactedEffectiveTools.some( + (tool) => tool.name === "skill_workshop", + ), + hookRunner, + promptStartedAt, }, - afterTurn: { - activeContextEngine: input.activeContextEngine, - readLifecycleState: () => { - const terminal = readTerminal(); - return { - aborted: terminal.aborted, - timedOut: terminal.timedOut, - idleTimedOut: terminal.idleTimedOut, - timedOutDuringCompaction: terminal.timedOutDuringCompaction, - }; - }, - runtime: { - effectiveWorkspace: input.setup.effectiveWorkspace, - agentDir: input.agentDir, - sessionAgentId: input.setup.sessionAgentId, - resolveActiveContextEnginePluginId: input.resolveActiveContextEnginePluginId, - shouldRecordCompletedBootstrapTurn, - cacheTrace, - anthropicPayloadLogger, - hookAgentId, - diagnosticTrace: input.diagnostics.diagnosticTrace, - skillWorkshopAvailable: uncompactedEffectiveTools.some( - (tool) => tool.name === "skill_workshop", - ), - hookRunner, - promptStartedAt, - }, + state: { + promptError: settledStream.promptError, + yieldAborted, + sessionIdUsed: settledStream.sessionIdUsed, + sessionFileUsed, + messagesSnapshot: settledStream.messagesSnapshot, + prePromptMessageCount: sessionRuntimeState.prePromptMessageCount, + contextEngineAfterTurnCheckpoint: contextGuards.getAfterTurnCheckpoint(), + lastCallUsage: settledStream.lastCallUsage, + promptCache: settledStream.promptCache, + ...(beforeAgentFinalizeRevisionReason ? { beforeAgentFinalizeRevisionReason } : {}), + compactionOccurredThisAttempt: settledStream.compactionOccurredThisAttempt, }, }); sessionIdUsed = afterTurn.sessionIdUsed; 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 0feb1de486a7..3b1aee4f8f5a 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 @@ -1,102 +1,232 @@ import { beforeEach, describe, expect, it, vi } from "vitest"; const mocks = vi.hoisted(() => ({ + clearActiveEmbeddedRun: vi.fn(), completeAfterTurn: vi.fn(), + completeResult: vi.fn(), + logDebug: vi.fn(), + logError: vi.fn(), + logWarn: vi.fn(), + runPrompt: vi.fn(), settleStream: vi.fn(), })); +vi.mock("../logger.js", () => ({ + log: { debug: mocks.logDebug, error: mocks.logError, warn: mocks.logWarn }, +})); +vi.mock("../runs.js", () => ({ clearActiveEmbeddedRun: mocks.clearActiveEmbeddedRun })); vi.mock("./attempt-finalize.js", () => ({ completeEmbeddedAttemptAfterTurn: mocks.completeAfterTurn, })); +vi.mock("./attempt-prompt-phase.js", () => ({ + runEmbeddedAttemptPromptPhase: mocks.runPrompt, +})); +vi.mock("./attempt-result.js", () => ({ + completeEmbeddedAttemptResult: mocks.completeResult, +})); vi.mock("./attempt-stream-settle.js", () => ({ settleEmbeddedAttemptStream: mocks.settleStream, })); import { createSubscribedSessionHarness } from "../../embedded-agent-subscribe.e2e-harness.js"; import { SessionManager } from "../../sessions/index.js"; -import { finalizeEmbeddedAttemptStreamPhase } from "./attempt-stream-finalize.js"; +import { runEmbeddedAttemptSettledPhase } from "./attempt-settle.js"; -type FinalizeInput = Parameters[0]; +type SettledInput = Parameters[0]; type SettleMockInput = { state: { promptError: unknown; promptErrorSource: unknown; }; }; +type FixtureOverrides = { + activeSession?: SettledInput["prepared"]["sessionRuntime"]["agentSession"]["activeSession"]; + getBeforeAgentFinalizeRevisionEntryId?: () => string | undefined; + getBeforeAgentFinalizeRevisionReason?: () => string | undefined; + repairedRejectedThinkingReplay?: boolean; + runAbortController?: AbortController; + sessionManager?: SettledInput["prepared"]["sessionRuntime"]["sessionManager"]; + waitForPendingEvents?: () => Promise; +}; -function createFixture(overrides?: Partial) { +function createFixture(overrides: FixtureOverrides = {}) { const order: string[] = []; const repairedMessages = [{ role: "user", content: "repaired" }]; - const activeSession = { - agent: { state: { messages: [] } }, - }; - const phaseState: ReturnType = { - promptError: null, - promptErrorSource: null, - yieldAborted: false, - sessionIdUsed: "initial-session", - sessionFileUsed: "initial.jsonl", - }; - const input = { - attempt: { runId: "run-1" }, - activeSession, - sessionManager: { + const activeSession = + overrides.activeSession ?? + ({ + agent: { state: { messages: [] } }, + getActiveToolNames: vi.fn(() => ["read"]), + sessionId: "active-session", + } as never); + const sessionManager = + overrides.sessionManager ?? + ({ + appendLeafControl: vi.fn(), buildSessionContext: () => ({ messages: repairedMessages }), - }, - withOwnedTranscriptWrite: vi.fn(async (operation) => await operation()), - waitForPendingEvents: vi.fn(async () => { + getEntry: vi.fn(), + } as never); + const waitForPendingEvents = + overrides.waitForPendingEvents ?? + vi.fn(async () => { order.push("pending-events"); - }), - repairedRejectedThinkingReplay: true, - getRunAbortDeadlineAtMs: () => 123, - shouldFlushForContextEngine: () => true, - getBeforeAgentFinalizeRevisionReason: () => "revision changed", - getBeforeAgentFinalizeRevisionEntryId: () => undefined, - getContextEngineAfterTurnCheckpoint: () => 7, - onSettleErrorState: vi.fn(), - onSettled: vi.fn(() => { - order.push("settled-published"); - }), - getState: () => phaseState, - settle: { - subscription: {}, - readLifecycleState: () => ({ - aborted: false, - timedOut: false, - timedOutDuringCompaction: false, - }), - markTimedOutDuringCompaction: vi.fn(), - runAbortSignal: new AbortController().signal, - isProbeSession: false, - abortable: async (promise: Promise) => await promise, - prePromptMessageCount: 3, - toolSearchTargetTranscriptProjections: [], - cache: { - observabilityEnabled: false, - changesForTurn: null, - retention: undefined, + }); + const getBeforeAgentFinalizeRevisionReason = + overrides.getBeforeAgentFinalizeRevisionReason ?? (() => "revision changed"); + const getBeforeAgentFinalizeRevisionEntryId = + overrides.getBeforeAgentFinalizeRevisionEntryId ?? (() => undefined); + const unsubscribe = vi.fn(); + const subscription = { + isCompacting: vi.fn(() => false), + unsubscribe, + waitForPendingEvents, + }; + const queueHandle = { kind: "embedded", runId: "run-1" }; + const sessionRuntimeState = { + prePromptMessageCount: 3, + promptCache: undefined, + systemPromptText: "system prompt", + }; + const state: SettledInput["state"] = { + beforeAgentRunBlockedBy: undefined, + terminal: { kind: "ok" }, + trajectoryEndRecorded: false, + }; + let markYieldAborted: (() => void) | undefined; + const input = { + attempt: { + runId: "run-1", + sessionFile: "initial.jsonl", + sessionId: "session-1", + }, + activeContextEngine: { info: { id: "engine" } }, + agentDir: "/agent", + isRawModelRun: false, + resolveActiveContextEnginePluginId: vi.fn(), + runAbortController: overrides.runAbortController ?? new AbortController(), + prepared: { + bootstrap: { + bootstrapPromptWarning: undefined, + shouldRecordCompletedBootstrapTurn: false, + }, + bundleTools: { + tools: [{ name: "read" }], + uncompactedEffectiveTools: [{ name: "read" }], + }, + sessionRuntime: { + agentSession: { + activeSession, + clientToolCallSlots: [], + hasDeliveredSourceReply: vi.fn(() => false), + hookRunner: {}, + setActiveSessionSystemPrompt: vi.fn(), + settingsManager: { getCompactionReserveTokens: vi.fn(() => 1_000) }, + }, + anthropicPayloadLogger: {}, + boundary: { + boundaryTimezone: "UTC", + includeBoundaryTimestamp: true, + orphanRepair: undefined, + setCurrentUserTimestampOverride: vi.fn(), + }, + cacheTrace: {}, + contextGuards: { + getAfterTurnCheckpoint: vi.fn(() => 7), + takePendingMidTurnPrecheckRequest: vi.fn(() => null), + }, + preparedUserTurnMessage: undefined, + sessionManager, + sessionPromptState: {}, + state: sessionRuntimeState, + toolResultPromptProjectionState: {}, + trajectoryRecorder: {}, + transport: { + effectiveAgentTransport: "sse", + effectiveExtraParams: {}, + effectivePromptCacheRetention: undefined, + streamStrategy: "provider", + }, + }, + systemPrompt: { + runtimeInfo: { model: { id: "model" } }, + systemPromptReport: undefined, + }, + toolBase: { toolSearchTargetTranscriptProjections: [] }, + toolCatalog: { + effectiveTools: [{ name: "read" }], + emptyExplicitToolAllowlistError: undefined, + toolSearch: { compacted: false }, }, }, - afterTurn: { - readLifecycleState: () => ({ - aborted: false, - timedOut: false, - idleTimedOut: false, - timedOutDuringCompaction: false, - }), - runtime: {}, + sessionLock: { + withOwnedTranscriptWrite: vi.fn(async (operation) => await operation()), }, - ...overrides, - } as unknown as FinalizeInput; + setup: { + effectiveFsWorkspaceOnly: false, + effectiveWorkspace: "/workspace", + sandbox: null, + sessionAgentId: "main", + }, + diagnostics: { diagnosticTrace: {}, runTrace: {} }, + state, + lifecycle: { + readYieldState: () => ({ + yieldAbortSettled: null, + yieldDetected: false, + yieldMessage: null, + }), + }, + getRepairedRejectedThinkingReplay: () => overrides.repairedRejectedThinkingReplay ?? true, + preparedStreamRuntime: { + abortable: async (promise: Promise) => await promise, + cache: { observabilityEnabled: false, promptTools: [] }, + history: { + contextEnginePromptAuthority: "assembled", + contextEngineAssemblySucceeded: true, + }, + isProbeSession: false, + onBlockReplyFlush: undefined, + promptActiveSession: vi.fn(async () => undefined), + stream: { + subscription, + queueHandle, + stopAcceptingSteerMessages: vi.fn(), + getBeforeAgentFinalizeRevisionReason, + getBeforeAgentFinalizeRevisionEntryId, + }, + timeout: { + getRunAbortDeadlineAtMs: () => 123, + clearTimers: vi.fn(), + }, + }, + } as unknown as SettledInput; - return { activeSession, input, order, phaseState, repairedMessages }; + mocks.runPrompt.mockImplementation(async (promptInput) => { + markYieldAborted = promptInput.lifecycle.markYieldAborted; + return { promptStartedAt: 100 }; + }); + mocks.completeResult.mockImplementation((resultInput) => ({ + sessionIdUsed: resultInput.state.sessionIdUsed, + sessionFileUsed: resultInput.state.sessionFileUsed, + })); + mocks.clearActiveEmbeddedRun.mockReturnValue(undefined); + + return { + activeSession, + input, + markYieldAborted: () => markYieldAborted?.(), + order, + repairedMessages, + sessionRuntimeState, + state, + }; } beforeEach(() => { vi.clearAllMocks(); }); -describe("finalizeEmbeddedAttemptStreamPhase", () => { +describe("runEmbeddedAttemptSettledPhase stream finalization", () => { it("does not settle a provider failure before partial presentation finishes", async () => { let resolvePartial: (() => void) | undefined; const onPartialReply = vi.fn( @@ -147,7 +277,7 @@ describe("finalizeEmbeddedAttemptStreamPhase", () => { }); mocks.completeAfterTurn.mockResolvedValue({ sessionIdUsed: "session-1" }); - const finalize = finalizeEmbeddedAttemptStreamPhase(fixture.input); + const finalize = runEmbeddedAttemptSettledPhase(fixture.input); await vi.waitFor(() => expect(onPartialReply).toHaveBeenCalledOnce()); await Promise.resolve(); expect(mocks.settleStream).not.toHaveBeenCalled(); @@ -173,8 +303,13 @@ describe("finalizeEmbeddedAttemptStreamPhase", () => { sessionManager.appendCustomEntry("trailing-metadata", { source: "hook" }); sessionManager.appendCompaction("Summary including rejected answer", promptId, 100); const originalMessages = sessionManager.buildSessionContext().messages; + const activeSession = { + agent: { state: { messages: originalMessages } }, + getActiveToolNames: vi.fn(() => ["read"]), + sessionId: "active-session", + }; const fixture = createFixture({ - activeSession: { agent: { state: { messages: originalMessages } } } as never, + activeSession: activeSession as never, sessionManager: sessionManager as never, repairedRejectedThinkingReplay: false, getBeforeAgentFinalizeRevisionEntryId: () => rejectedId, @@ -195,7 +330,7 @@ describe("finalizeEmbeddedAttemptStreamPhase", () => { promptCache: undefined, }; mocks.settleStream.mockImplementation(async () => { - expect(fixture.input.activeSession.agent.state.messages).toBe(originalMessages); + expect(activeSession.agent.state.messages).toBe(originalMessages); expect(sessionManager.getLeafId()).toBe(promptId); return settledStream; }); @@ -204,7 +339,7 @@ describe("finalizeEmbeddedAttemptStreamPhase", () => { sessionFileUsed: "session.jsonl", }); - await finalizeEmbeddedAttemptStreamPhase(fixture.input); + await runEmbeddedAttemptSettledPhase(fixture.input); const retryMessages = sessionManager.buildSessionContext().messages; const retryTranscript = JSON.stringify(retryMessages); @@ -233,12 +368,17 @@ describe("finalizeEmbeddedAttemptStreamPhase", () => { }); it("settles the stream before publishing state and running after-turn work", async () => { - const fixture = createFixture(); const pendingError = new Error("pending event failed"); - fixture.input.waitForPendingEvents = vi.fn(async () => { - fixture.order.push("pending-events"); - fixture.phaseState.promptError = pendingError; - fixture.phaseState.promptErrorSource = "prompt"; + const reason = vi + .fn<() => string | undefined>() + .mockReturnValueOnce("revision changed") + .mockReturnValueOnce(undefined); + const fixture = createFixture({ + getBeforeAgentFinalizeRevisionReason: reason, + waitForPendingEvents: vi.fn(async () => { + fixture.order.push("pending-events"); + fixture.state.terminal = { kind: "failed", error: pendingError, source: "prompt" }; + }), }); const settledStream = { promptError: null, @@ -249,24 +389,26 @@ describe("finalizeEmbeddedAttemptStreamPhase", () => { sessionIdUsed: "settled-session", lastAssistant: undefined, currentAttemptAssistant: undefined, + currentAttemptCompletedAssistant: undefined, attemptUsage: undefined, cacheBreak: null, lastCallUsage: undefined, - promptCache: undefined, + promptCache: { published: true }, }; - mocks.settleStream.mockImplementation(async (input: SettleMockInput) => { + mocks.settleStream.mockImplementation(async (settleInput: SettleMockInput) => { fixture.order.push("settle"); - expect(input.state.promptError).toBe(pendingError); - expect(input.state.promptErrorSource).toBe("prompt"); - fixture.phaseState.yieldAborted = true; + expect(settleInput.state.promptError).toBe(pendingError); + expect(settleInput.state.promptErrorSource).toBe("prompt"); + fixture.markYieldAborted(); return settledStream; }); mocks.completeAfterTurn.mockImplementation(async () => { - fixture.order.push("after-turn"); + expect(fixture.sessionRuntimeState.promptCache).toEqual({ published: true }); + fixture.order.push("settled-published", "after-turn"); return { sessionIdUsed: "after-session", sessionFileUsed: "after.jsonl" }; }); - await expect(finalizeEmbeddedAttemptStreamPhase(fixture.input)).resolves.toEqual({ + await expect(runEmbeddedAttemptSettledPhase(fixture.input)).resolves.toEqual({ sessionIdUsed: "after-session", sessionFileUsed: "after.jsonl", }); @@ -279,7 +421,6 @@ describe("finalizeEmbeddedAttemptStreamPhase", () => { shouldFlushForContextEngine: true, }), ); - expect(fixture.input.onSettled).toHaveBeenCalledWith(settledStream); expect(mocks.completeAfterTurn).toHaveBeenCalledWith( expect.objectContaining({ state: expect.objectContaining({ @@ -298,10 +439,10 @@ describe("finalizeEmbeddedAttemptStreamPhase", () => { it("proceeds to settlement when pending subscription events never settle", async () => { vi.useFakeTimers(); try { - const fixture = createFixture(); - // A hung delivery handler must not dead-end the turn until the run budget. - fixture.input.waitForPendingEvents = vi.fn(() => new Promise(() => {})); - const settledStream = { + const fixture = createFixture({ + waitForPendingEvents: vi.fn(() => new Promise(() => {})), + }); + mocks.settleStream.mockResolvedValue({ promptError: null, promptErrorSource: null, timedOutDuringCompaction: false, @@ -315,14 +456,13 @@ describe("finalizeEmbeddedAttemptStreamPhase", () => { cacheBreak: null, lastCallUsage: undefined, promptCache: undefined, - }; - mocks.settleStream.mockResolvedValue(settledStream); + }); mocks.completeAfterTurn.mockResolvedValue({ sessionIdUsed: "session-1", sessionFileUsed: "session.jsonl", }); - const finalize = finalizeEmbeddedAttemptStreamPhase(fixture.input); + const finalize = runEmbeddedAttemptSettledPhase(fixture.input); await vi.advanceTimersByTimeAsync(119_999); expect(mocks.settleStream).not.toHaveBeenCalled(); await vi.advanceTimersByTimeAsync(1); @@ -339,15 +479,12 @@ describe("finalizeEmbeddedAttemptStreamPhase", () => { it("skips the pending-events join once the run abort signal fires", async () => { const abortController = new AbortController(); abortController.abort(new Error("operator cancel")); - const fixture = createFixture(); - fixture.input.settle.runAbortSignal = abortController.signal; - fixture.input.waitForPendingEvents = vi.fn(() => new Promise(() => {})); - fixture.input.settle.readLifecycleState = () => ({ - aborted: true, - timedOut: false, - timedOutDuringCompaction: false, + const fixture = createFixture({ + runAbortController: abortController, + waitForPendingEvents: vi.fn(() => new Promise(() => {})), }); - const settledStream = { + fixture.state.terminal = { kind: "aborted", source: "external" }; + mocks.settleStream.mockResolvedValue({ promptError: null, promptErrorSource: null, timedOutDuringCompaction: false, @@ -361,14 +498,13 @@ describe("finalizeEmbeddedAttemptStreamPhase", () => { cacheBreak: null, lastCallUsage: undefined, promptCache: undefined, - }; - mocks.settleStream.mockResolvedValue(settledStream); + }); mocks.completeAfterTurn.mockResolvedValue({ sessionIdUsed: "session-1", sessionFileUsed: "session.jsonl", }); - await expect(finalizeEmbeddedAttemptStreamPhase(fixture.input)).resolves.toEqual({ + await expect(runEmbeddedAttemptSettledPhase(fixture.input)).resolves.toEqual({ sessionIdUsed: "session-1", sessionFileUsed: "session.jsonl", }); @@ -378,37 +514,36 @@ describe("finalizeEmbeddedAttemptStreamPhase", () => { it("settles an aborted run with its recorded cancellation reason", async () => { const cancellationReason = new Error("cancelled by operator"); const fixture = createFixture({ repairedRejectedThinkingReplay: false }); - fixture.input.settle.readLifecycleState = () => ({ - aborted: true, - timedOut: false, - timedOutDuringCompaction: false, + fixture.state.terminal = { kind: "aborted", source: "external" }; + mocks.settleStream.mockImplementation(async (settleInput) => { + expect(settleInput.readLifecycleState()).toEqual( + expect.objectContaining({ aborted: true, timedOut: false }), + ); + return { + promptError: cancellationReason, + promptErrorSource: "prompt", + timedOutDuringCompaction: false, + compactionOccurredThisAttempt: false, + messagesSnapshot: [], + sessionIdUsed: "session-1", + lastAssistant: undefined, + currentAttemptAssistant: undefined, + currentAttemptCompletedAssistant: undefined, + attemptUsage: undefined, + cacheBreak: null, + lastCallUsage: undefined, + promptCache: undefined, + }; }); - const settledStream = { - promptError: cancellationReason, - 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.settleStream.mockResolvedValue(settledStream); mocks.completeAfterTurn.mockResolvedValue({ sessionIdUsed: "session-1", sessionFileUsed: "session.jsonl", }); - await expect(finalizeEmbeddedAttemptStreamPhase(fixture.input)).resolves.toEqual({ + await expect(runEmbeddedAttemptSettledPhase(fixture.input)).resolves.toEqual({ sessionIdUsed: "session-1", sessionFileUsed: "session.jsonl", }); - expect(mocks.settleStream).toHaveBeenCalledOnce(); expect(mocks.completeAfterTurn).toHaveBeenCalledOnce(); }); @@ -417,21 +552,19 @@ describe("finalizeEmbeddedAttemptStreamPhase", () => { const fixture = createFixture({ repairedRejectedThinkingReplay: false }); const settlementError = new Error("settlement failed"); const promptError = new Error("prompt failed"); - mocks.settleStream.mockImplementation(async (input: SettleMockInput) => { - input.state.promptError = promptError; - input.state.promptErrorSource = "compaction"; + mocks.settleStream.mockImplementation(async (settleInput: SettleMockInput) => { + settleInput.state.promptError = promptError; + settleInput.state.promptErrorSource = "compaction"; throw settlementError; }); - await expect(finalizeEmbeddedAttemptStreamPhase(fixture.input)).rejects.toBe(settlementError); + await expect(runEmbeddedAttemptSettledPhase(fixture.input)).rejects.toBe(settlementError); - expect(fixture.input.onSettleErrorState).toHaveBeenCalledWith( - expect.objectContaining({ - promptError, - promptErrorSource: "compaction", - }), - ); - expect(fixture.input.onSettled).not.toHaveBeenCalled(); + expect(fixture.state.terminal).toEqual({ + kind: "failed", + error: promptError, + source: "compaction", + }); expect(mocks.completeAfterTurn).not.toHaveBeenCalled(); }); @@ -449,7 +582,11 @@ describe("finalizeEmbeddedAttemptStreamPhase", () => { timestamp: 2, } as never); const originalMessages = sessionManager.buildSessionContext().messages; - const activeSession = { agent: { state: { messages: originalMessages } } }; + const activeSession = { + agent: { state: { messages: originalMessages } }, + getActiveToolNames: vi.fn(() => ["read"]), + sessionId: "active-session", + }; const fixture = createFixture({ activeSession: activeSession as never, sessionManager: sessionManager as never, @@ -459,13 +596,13 @@ describe("finalizeEmbeddedAttemptStreamPhase", () => { const settlementError = new Error("settlement failed"); mocks.settleStream.mockRejectedValue(settlementError); - await expect(finalizeEmbeddedAttemptStreamPhase(fixture.input)).rejects.toBe(settlementError); + await expect(runEmbeddedAttemptSettledPhase(fixture.input)).rejects.toBe(settlementError); expect(sessionManager.getLeafId()).toBe(promptId); expect(JSON.stringify(activeSession.agent.state.messages)).not.toContain( "Rejected first answer", ); - expect(fixture.input.onSettleErrorState).toHaveBeenCalledOnce(); + expect(mocks.settleStream).toHaveBeenCalledOnce(); expect(mocks.completeAfterTurn).not.toHaveBeenCalled(); }); }); diff --git a/src/agents/embedded-agent-runner/run/attempt-stream-finalize.ts b/src/agents/embedded-agent-runner/run/attempt-stream-finalize.ts deleted file mode 100644 index 7add4a101a52..000000000000 --- a/src/agents/embedded-agent-runner/run/attempt-stream-finalize.ts +++ /dev/null @@ -1,149 +0,0 @@ -/** Settles the provider stream and completes the post-turn lifecycle phase. */ -import { log } from "../logger.js"; -import { joinWithRunLivenessDeadline, RUN_LIVENESS_JOIN_TIMEOUT_MS } from "./abortable.js"; -import { completeEmbeddedAttemptAfterTurn } from "./attempt-finalize.js"; -import { settleEmbeddedAttemptStream } from "./attempt-stream-settle.js"; - -type StreamSettleInput = Parameters[0]; -type StreamSettleResult = Awaited>; -type AfterTurnInput = Parameters[0]; -type FinalizePhaseState = StreamSettleInput["state"] & { - sessionFileUsed?: string; -}; - -type SharedPhaseInputKeys = - | "attempt" - | "activeSession" - | "sessionManager" - | "withOwnedTranscriptWrite"; - -// 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. - -export async function finalizeEmbeddedAttemptStreamPhase(input: { - attempt: StreamSettleInput["attempt"]; - activeSession: StreamSettleInput["activeSession"]; - sessionManager: StreamSettleInput["sessionManager"]; - withOwnedTranscriptWrite: StreamSettleInput["withOwnedTranscriptWrite"]; - waitForPendingEvents: () => Promise; - repairedRejectedThinkingReplay: boolean; - getRunAbortDeadlineAtMs: () => number; - shouldFlushForContextEngine: () => boolean; - getBeforeAgentFinalizeRevisionReason: () => string | undefined; - getBeforeAgentFinalizeRevisionEntryId: () => string | undefined; - getContextEngineAfterTurnCheckpoint: () => number | null; - onSettleErrorState: (state: { - promptError: unknown; - promptErrorSource: StreamSettleInput["state"]["promptErrorSource"]; - }) => void; - onSettled: (result: StreamSettleResult) => void; - getState: () => FinalizePhaseState; - settle: Omit< - StreamSettleInput, - SharedPhaseInputKeys | "state" | "runAbortDeadlineAtMs" | "shouldFlushForContextEngine" - >; - afterTurn: Omit; -}): Promise<{ sessionIdUsed: string; sessionFileUsed?: string }> { - const { activeSession, sessionManager, withOwnedTranscriptWrite } = input; - - await joinWithRunLivenessDeadline({ - joinWork: input.waitForPendingEvents, - runAbortSignal: input.settle.runAbortSignal, - onTimeout: () => { - log.warn( - `pending subscription events did not settle within ${RUN_LIVENESS_JOIN_TIMEOUT_MS}ms; ` + - `proceeding to stream settlement: runId=${input.attempt.runId}`, - ); - }, - }); - const beforeAgentFinalizeRevisionReason = input.getBeforeAgentFinalizeRevisionReason(); - const beforeAgentFinalizeRevisionEntryId = input.getBeforeAgentFinalizeRevisionEntryId(); - let rewoundBeforeAgentFinalizeRevision = false; - if (beforeAgentFinalizeRevisionReason && beforeAgentFinalizeRevisionEntryId) { - await withOwnedTranscriptWrite(() => { - const rejectedEntry = sessionManager.getEntry(beforeAgentFinalizeRevisionEntryId); - if (rejectedEntry?.type !== "message" || rejectedEntry.message.role !== "assistant") { - throw new Error( - `before_agent_finalize persisted assistant entry is missing or invalid ` + - `(entry=${beforeAgentFinalizeRevisionEntryId})`, - ); - } - // Keep persistence append-only while excluding the rejected draft and - // every trailing descendant from the hidden retry's active branch. - sessionManager.appendLeafControl({ - targetId: rejectedEntry.parentId, - appendParentId: rejectedEntry.parentId, - }); - rewoundBeforeAgentFinalizeRevision = true; - }); - } - let settledStream: StreamSettleResult; - try { - if (input.repairedRejectedThinkingReplay && !rewoundBeforeAgentFinalizeRevision) { - activeSession.agent.state.messages = sessionManager.buildSessionContext().messages; - } - const currentState = input.getState(); - const streamSettleState = { - promptError: currentState.promptError, - promptErrorSource: currentState.promptErrorSource, - yieldAborted: currentState.yieldAborted, - sessionIdUsed: currentState.sessionIdUsed, - }; - try { - settledStream = await settleEmbeddedAttemptStream({ - attempt: input.attempt, - activeSession, - sessionManager, - withOwnedTranscriptWrite, - state: streamSettleState, - ...input.settle, - runAbortDeadlineAtMs: input.getRunAbortDeadlineAtMs(), - shouldFlushForContextEngine: input.shouldFlushForContextEngine(), - }); - } catch (error) { - // Settlement mutates this shared state before some failures. Publish it so - // outer teardown keeps the recorded prompt error and attribution. - input.onSettleErrorState(streamSettleState); - throw error; - } - } finally { - if (rewoundBeforeAgentFinalizeRevision) { - await withOwnedTranscriptWrite(() => { - // Settlement classifies the completed attempt from its original - // in-memory messages. Later work always sees the rewound branch. - activeSession.agent.state.messages = sessionManager.buildSessionContext().messages; - }); - } - } - // Publish settled fields before after-turn hooks: those hooks may throw, and - // outer teardown still needs the completed stream snapshot and usage state. - input.onSettled(settledStream); - - const afterSettleState = input.getState(); - const afterTurn = await completeEmbeddedAttemptAfterTurn({ - attempt: input.attempt, - activeSession, - sessionManager, - withOwnedTranscriptWrite, - ...input.afterTurn, - state: { - promptError: settledStream.promptError, - yieldAborted: afterSettleState.yieldAborted, - sessionIdUsed: settledStream.sessionIdUsed, - sessionFileUsed: afterSettleState.sessionFileUsed, - messagesSnapshot: settledStream.messagesSnapshot, - prePromptMessageCount: input.settle.prePromptMessageCount, - contextEngineAfterTurnCheckpoint: input.getContextEngineAfterTurnCheckpoint(), - lastCallUsage: settledStream.lastCallUsage, - promptCache: settledStream.promptCache, - ...(beforeAgentFinalizeRevisionReason ? { beforeAgentFinalizeRevisionReason } : {}), - compactionOccurredThisAttempt: settledStream.compactionOccurredThisAttempt, - }, - }); - - return afterTurn; -} diff --git a/src/agents/embedded-agent-runner/run/attempt-stream-runtime-prepare.test.ts b/src/agents/embedded-agent-runner/run/attempt-stream-runtime-prepare.test.ts deleted file mode 100644 index 40cb0df7681a..000000000000 --- a/src/agents/embedded-agent-runner/run/attempt-stream-runtime-prepare.test.ts +++ /dev/null @@ -1,267 +0,0 @@ -import { beforeEach, describe, expect, it, vi } from "vitest"; - -const mocks = vi.hoisted(() => ({ - abortable: vi.fn(), - bindOwnedSessionTranscriptWrites: vi.fn(), - createRunAbort: vi.fn(), - flushPendingToolResultsAfterIdle: vi.fn(), - installStreamGuards: vi.fn(), - prepareHistory: vi.fn(), - prepareStream: vi.fn(), - prepareTimeout: vi.fn(), - withOwnedSessionTranscriptWrites: vi.fn(), -})); - -vi.mock("../../../config/sessions/transcript-write-context.js", () => ({ - bindOwnedSessionTranscriptWrites: mocks.bindOwnedSessionTranscriptWrites, - withOwnedSessionTranscriptWrites: mocks.withOwnedSessionTranscriptWrites, -})); -vi.mock("../wait-for-idle-before-flush.js", () => ({ - flushPendingToolResultsAfterIdle: mocks.flushPendingToolResultsAfterIdle, -})); -vi.mock("./abortable.js", () => ({ abortable: mocks.abortable })); -vi.mock("./attempt-finalize.js", () => ({ - createEmbeddedAttemptRunAbort: mocks.createRunAbort, -})); -vi.mock("./attempt-history.js", () => ({ - prepareEmbeddedAttemptHistory: mocks.prepareHistory, -})); -vi.mock("./attempt-stream-prepare.js", () => ({ - prepareEmbeddedAttemptStream: mocks.prepareStream, -})); -vi.mock("./attempt-stream.js", () => ({ - installEmbeddedAttemptStreamGuards: mocks.installStreamGuards, -})); -vi.mock("./attempt-timeout-prepare.js", () => ({ - prepareEmbeddedAttemptTimeout: mocks.prepareTimeout, -})); - -import { prepareEmbeddedAttemptStreamRuntime } from "./attempt-stream-runtime-prepare.js"; - -type StreamRuntimeInput = Parameters[0]; - -function createFixture(options: { aborted?: boolean } = {}) { - const order: string[] = []; - const abortController = new AbortController(); - if (options.aborted) { - abortController.abort(new Error("already aborted")); - } - const runAbort = vi.fn(); - const toolSearchCatalogExecutor = vi.fn(); - const subscription = { - isCompacting: vi.fn(() => false), - }; - const queueHandle = { kind: "embedded", runId: "run-1" }; - const streamResult = { - subscription, - queueHandle, - toolSearchCatalogExecutor, - getBeforeAgentFinalizeRevisionReason: vi.fn(), - stopAcceptingSteerMessages: vi.fn(), - }; - const timeoutResult = { - getRunAbortDeadlineAtMs: vi.fn(() => 123), - clearTimers: vi.fn(), - }; - const activeSession = { - agent: { streamFn: vi.fn() }, - dispose: vi.fn(), - isCompacting: false, - messages: [], - prompt: vi.fn(async () => undefined), - }; - const sessionManager = {}; - const externalAbortController = { - setRunAbort: vi.fn(() => order.push("set-run-abort")), - setCompactionState: vi.fn(() => order.push("set-compaction-state")), - }; - const markIdleTimedOut = vi.fn(); - const markStreamReady = vi.fn(() => order.push("stream-ready")); - const setToolSearchCatalogExecutor = vi.fn(() => order.push("set-catalog")); - const trackPromptSettlePromise = vi.fn((promise: Promise) => promise); - - mocks.abortable.mockImplementation((_signal, promise) => promise); - mocks.bindOwnedSessionTranscriptWrites.mockImplementation((_context, operation) => operation); - mocks.withOwnedSessionTranscriptWrites.mockImplementation( - async (_context, operation) => await operation(), - ); - mocks.installStreamGuards.mockImplementation(() => { - order.push("guards"); - return { - cacheObservabilityEnabled: true, - promptCacheTools: [{ name: "read" }], - }; - }); - mocks.prepareHistory.mockImplementation(async () => { - order.push("history"); - return { - contextEnginePromptAuthority: "assembled", - contextEngineAssemblySucceeded: true, - }; - }); - mocks.createRunAbort.mockImplementation(() => { - order.push("abort"); - return runAbort; - }); - mocks.prepareStream.mockImplementation(() => { - order.push("stream"); - return streamResult; - }); - mocks.prepareTimeout.mockImplementation(() => { - order.push("timeout"); - return timeoutResult; - }); - - const input = { - attempt: { - abortSignal: abortController.signal, - onBlockReply: vi.fn(), - onBlockReplyFlush: vi.fn(), - runId: "run-1", - sessionId: "session-1", - timeoutMs: 30_000, - }, - activeSession, - sessionManager, - sessionLockController: {}, - ownedTranscriptWriteContext: {}, - runAbortController: new AbortController(), - externalAbortController, - abortActiveSession: vi.fn(async () => undefined), - abortState: {}, - trackPromptSettlePromise, - compactionTimeoutMs: 1_000, - guards: {}, - history: { sandboxed: false }, - stream: {}, - lifecycle: { - isYieldDetected: () => false, - markRejectedThinkingReplayRepaired: vi.fn(), - markStreamReady, - markIdleTimedOut, - markExternalAbort: vi.fn(), - markTimedOutDuringCompaction: vi.fn(), - markTimedOutByRunBudget: vi.fn(), - readRunState: () => ({ - aborted: false, - promptError: null, - timedOut: false, - yieldDetected: false, - }), - setToolSearchCatalogExecutor, - }, - } as unknown as StreamRuntimeInput; - - return { - activeSession, - externalAbortController, - input, - markIdleTimedOut, - order, - runAbort, - sessionManager, - streamResult, - subscription, - timeoutResult, - toolSearchCatalogExecutor, - trackPromptSettlePromise, - }; -} - -beforeEach(() => { - vi.clearAllMocks(); -}); - -describe("prepareEmbeddedAttemptStreamRuntime", () => { - it("prepares guarded history, abort handling, stream subscription, and timeout in order", async () => { - const fixture = createFixture(); - - const result = await prepareEmbeddedAttemptStreamRuntime(fixture.input); - - expect(fixture.order).toEqual([ - "guards", - "stream-ready", - "history", - "abort", - "set-run-abort", - "stream", - "set-catalog", - "set-compaction-state", - "timeout", - ]); - expect(result).toEqual( - expect.objectContaining({ - cache: { - observabilityEnabled: true, - promptTools: [{ name: "read" }], - }, - history: expect.objectContaining({ contextEngineAssemblySucceeded: true }), - isProbeSession: false, - stream: fixture.streamResult, - timeout: fixture.timeoutResult, - }), - ); - expect(fixture.input.lifecycle.setToolSearchCatalogExecutor).toHaveBeenCalledWith( - fixture.toolSearchCatalogExecutor, - ); - expect(fixture.externalAbortController.setCompactionState).toHaveBeenCalledWith({ - isPendingOrRetrying: fixture.subscription.isCompacting, - isInFlight: expect.any(Function), - }); - expect(mocks.prepareTimeout).toHaveBeenCalledWith( - expect.objectContaining({ - abortRun: fixture.runAbort, - compactionState: fixture.subscription, - }), - ); - - const guardInput = mocks.installStreamGuards.mock.calls[0]?.[0]; - const idleError = new Error("idle timeout"); - guardInput.onIdleTimeout(idleError); - expect(fixture.markIdleTimedOut).toHaveBeenCalledOnce(); - expect(fixture.runAbort).toHaveBeenCalledWith(true, idleError); - - await result.promptActiveSession("hello"); - expect(fixture.activeSession.prompt).toHaveBeenCalledWith("hello", undefined); - expect(fixture.trackPromptSettlePromise).toHaveBeenCalledOnce(); - expect(mocks.withOwnedSessionTranscriptWrites).toHaveBeenCalledOnce(); - }); - - it.each([ - { label: "external cancellation", message: "run cancelled" }, - { label: "run timeout", message: "run timed out" }, - ])("does not start a prompt after $label", async ({ message }) => { - const fixture = createFixture(); - const runtime = await prepareEmbeddedAttemptStreamRuntime(fixture.input); - const reason = new Error(message); - const abortError = new Error(message, { cause: reason }); - abortError.name = "AbortError"; - fixture.input.runAbortController.abort(reason); - mocks.abortable.mockImplementationOnce((_signal, _promise) => Promise.reject(abortError)); - - await expect(runtime.promptActiveSession("must not start")).rejects.toBe(abortError); - - expect(fixture.activeSession.prompt).not.toHaveBeenCalled(); - expect(fixture.trackPromptSettlePromise).not.toHaveBeenCalled(); - expect(mocks.abortable).toHaveBeenCalledOnce(); - }); - - it("flushes pending tool results and disposes the session when history preparation fails", async () => { - const fixture = createFixture({ aborted: true }); - const failure = new Error("history failed"); - mocks.prepareHistory.mockRejectedValueOnce(failure); - mocks.flushPendingToolResultsAfterIdle.mockResolvedValue(undefined); - - await expect(prepareEmbeddedAttemptStreamRuntime(fixture.input)).rejects.toBe(failure); - - expect(mocks.flushPendingToolResultsAfterIdle).toHaveBeenCalledWith({ - agent: fixture.activeSession.agent, - sessionManager: fixture.sessionManager, - timeoutMs: 0, - }); - expect(fixture.activeSession.dispose).toHaveBeenCalledOnce(); - expect(mocks.createRunAbort).not.toHaveBeenCalled(); - expect(mocks.prepareStream).not.toHaveBeenCalled(); - expect(mocks.prepareTimeout).not.toHaveBeenCalled(); - }); -}); diff --git a/src/agents/embedded-agent-runner/run/attempt-stream-runtime-prepare.ts b/src/agents/embedded-agent-runner/run/attempt-stream-runtime-prepare.ts deleted file mode 100644 index 75869ab36970..000000000000 --- a/src/agents/embedded-agent-runner/run/attempt-stream-runtime-prepare.ts +++ /dev/null @@ -1,192 +0,0 @@ -/** Prepares guarded history, abort handling, stream subscription, and run deadlines. */ -import { - bindOwnedSessionTranscriptWrites, - withOwnedSessionTranscriptWrites, -} from "../../../config/sessions/transcript-write-context.js"; -import { log } from "../logger.js"; -import type { EmbeddedAgentQueueHandle } from "../runs.js"; -import { flushPendingToolResultsAfterIdle } from "../wait-for-idle-before-flush.js"; -import { abortable as abortableWithSignal } from "./abortable.js"; -import { - type createEmbeddedAttemptExternalAbortController, - createEmbeddedAttemptRunAbort, -} from "./attempt-finalize.js"; -import { prepareEmbeddedAttemptHistory } from "./attempt-history.js"; -import { prepareEmbeddedAttemptStream } from "./attempt-stream-prepare.js"; -import { installEmbeddedAttemptStreamGuards } from "./attempt-stream.js"; -import { prepareEmbeddedAttemptTimeout } from "./attempt-timeout-prepare.js"; -import type { EmbeddedRunAttemptParams } from "./types.js"; - -type StreamGuardInput = Parameters[0]; -type HistoryInput = Parameters[0]; -type StreamInput = Parameters[0]; -type ToolResultFlushInput = Parameters[0]; -type ExternalAbortController = Pick< - ReturnType, - "setCompactionState" | "setRunAbort" ->; -type StreamGuardPhaseInput = Omit< - StreamGuardInput, - | "abortSignal" - | "attempt" - | "isYieldDetected" - | "onIdleTimeout" - | "onRejectedThinkingReplayRepaired" - | "session" - | "sessionManager" ->; -type HistoryPhaseInput = Omit; -type StreamPhaseInput = Omit< - StreamInput, - | "abortRun" - | "activeSession" - | "attempt" - | "getRunState" - | "markExternalAbort" - | "onBlockReply" - | "onBlockReplyFlush" - | "runAbortController" ->; - -export async function prepareEmbeddedAttemptStreamRuntime(input: { - attempt: EmbeddedRunAttemptParams; - activeSession: StreamInput["activeSession"]; - sessionManager: HistoryInput["sessionManager"] & - NonNullable; - ownedTranscriptWriteContext: Parameters[0]; - runAbortController: AbortController; - externalAbortController: ExternalAbortController; - abortActiveSession: Parameters[0]["abortActiveSession"]; - abortState: Parameters[0]["state"]; - trackPromptSettlePromise: (promise: Promise) => Promise; - compactionTimeoutMs: number; - guards: StreamGuardPhaseInput; - history: HistoryPhaseInput; - stream: StreamPhaseInput; - lifecycle: { - isYieldDetected: StreamGuardInput["isYieldDetected"]; - markRejectedThinkingReplayRepaired: () => void; - markStreamReady: () => void; - markIdleTimedOut: () => void; - markExternalAbort: () => void; - markTimedOutDuringCompaction: () => void; - markTimedOutByRunBudget: () => void; - readRunState: StreamInput["getRunState"]; - setToolSearchCatalogExecutor: ( - executor: ReturnType["toolSearchCatalogExecutor"], - ) => void; - }; -}) { - const { activeSession, attempt, sessionManager } = input; - const idleTimeoutTriggerRef: { current?: (error: Error) => void } = {}; - const { cacheObservabilityEnabled, promptCacheTools } = installEmbeddedAttemptStreamGuards({ - ...input.guards, - attempt, - session: activeSession, - sessionManager, - isYieldDetected: input.lifecycle.isYieldDetected, - onRejectedThinkingReplayRepaired: input.lifecycle.markRejectedThinkingReplayRepaired, - onIdleTimeout: (error) => idleTimeoutTriggerRef.current?.(error), - abortSignal: input.runAbortController.signal, - }); - input.lifecycle.markStreamReady(); - - let preparedHistory: Awaited>; - try { - preparedHistory = await prepareEmbeddedAttemptHistory({ - ...input.history, - attempt, - activeSession, - sessionManager, - }); - } catch (error) { - await flushPendingToolResultsAfterIdle({ - agent: activeSession.agent, - sessionManager, - // An already-aborted setup must dispose immediately without orphaning tool calls. - ...(attempt.abortSignal?.aborted ? { timeoutMs: 0 } : {}), - }); - activeSession.dispose(); - throw error; - } - - const isProbeSession = attempt.sessionId?.startsWith("probe-") ?? false; - const queueHandleRef: { current?: EmbeddedAgentQueueHandle } = {}; - const abortRun = createEmbeddedAttemptRunAbort({ - abortActiveSession: input.abortActiveSession, - activeSession, - attempt, - getQueueHandle: () => queueHandleRef.current, - isProbeSession, - log, - runAbortController: input.runAbortController, - state: input.abortState, - }); - input.externalAbortController.setRunAbort(abortRun); - idleTimeoutTriggerRef.current = (error) => { - input.lifecycle.markIdleTimedOut(); - abortRun(true, error); - }; - const abortable = (promise: Promise): Promise => - abortableWithSignal(input.runAbortController.signal, promise); - const promptActiveSession = ( - prompt: string, - options?: Parameters[1], - ): Promise => - withOwnedSessionTranscriptWrites(input.ownedTranscriptWriteContext, async () => { - // Prompting starts its own agent loop; reject before creating a loop that - // an already-aborted attempt can no longer cancel. - if (input.runAbortController.signal.aborted) { - return abortable(Promise.resolve()); - } - return abortable(input.trackPromptSettlePromise(activeSession.prompt(prompt, options))); - }); - const onBlockReply = attempt.onBlockReply - ? bindOwnedSessionTranscriptWrites(input.ownedTranscriptWriteContext, attempt.onBlockReply) - : undefined; - const onBlockReplyFlush = attempt.onBlockReplyFlush - ? bindOwnedSessionTranscriptWrites(input.ownedTranscriptWriteContext, attempt.onBlockReplyFlush) - : undefined; - const preparedStream = prepareEmbeddedAttemptStream({ - ...input.stream, - attempt, - activeSession, - runAbortController: input.runAbortController, - abortRun, - markExternalAbort: input.lifecycle.markExternalAbort, - getRunState: input.lifecycle.readRunState, - onBlockReply, - onBlockReplyFlush, - }); - input.lifecycle.setToolSearchCatalogExecutor(preparedStream.toolSearchCatalogExecutor); - input.externalAbortController.setCompactionState({ - isPendingOrRetrying: preparedStream.subscription.isCompacting, - isInFlight: () => activeSession.isCompacting, - }); - queueHandleRef.current = preparedStream.queueHandle; - - const attemptTimeout = prepareEmbeddedAttemptTimeout({ - attempt, - activeSession, - compactionState: preparedStream.subscription, - compactionTimeoutMs: input.compactionTimeoutMs, - isProbeSession, - abortRun, - markTimedOutDuringCompaction: input.lifecycle.markTimedOutDuringCompaction, - markTimedOutByRunBudget: input.lifecycle.markTimedOutByRunBudget, - }); - - return { - abortable, - cache: { - observabilityEnabled: cacheObservabilityEnabled, - promptTools: promptCacheTools, - }, - history: preparedHistory, - isProbeSession, - onBlockReplyFlush, - promptActiveSession, - stream: preparedStream, - timeout: attemptTimeout, - }; -}