diff --git a/src/agents/embedded-agent-runner/runs.steering.test.ts b/src/agents/embedded-agent-runner/runs.steering.test.ts index b1018f2aa4ef..e697f90c78d7 100644 --- a/src/agents/embedded-agent-runner/runs.steering.test.ts +++ b/src/agents/embedded-agent-runner/runs.steering.test.ts @@ -5,30 +5,17 @@ import { setDiagnosticsEnabledForProcess } from "../../infra/diagnostic-events.j import { resetDiagnosticRunActivityForTest } from "../../logging/diagnostic-run-activity.js"; import { markDiagnosticToolStartedForTest } from "../../logging/diagnostic-run-activity.test-support.js"; import { resetDiagnosticSessionStateForTest } from "../../logging/diagnostic-session-state.js"; -import { queueEmbeddedAgentMessageWithOutcome, setActiveEmbeddedRun } from "./runs.js"; -import { testing } from "./runs.test-support.js"; +import { createUserTurnTranscriptRecorder } from "../../sessions/user-turn-transcript.js"; +import { createTestUserTurnTranscriptTarget } from "../../sessions/user-turn-transcript.test-support.js"; +import { + formatEmbeddedAgentQueueFailureSummary, + queueEmbeddedAgentMessageWithOutcome, + queueEmbeddedAgentMessageWithOutcomeAsync, + setActiveEmbeddedRun, +} from "./runs.js"; +import { createEmbeddedRunHandle, testing } from "./runs.test-support.js"; -type RunHandle = Parameters[1]; - -function createSteeringRunHandle( - overrides: { - isStreaming?: boolean; - isStopped?: () => boolean; - queueMessage?: RunHandle["queueMessage"]; - supportsQueueMessageImages?: boolean; - } = {}, -): RunHandle { - return { - queueMessage: overrides.queueMessage ?? (async () => {}), - isStreaming: () => overrides.isStreaming ?? true, - ...(overrides.isStopped ? { isStopped: overrides.isStopped } : {}), - isCompacting: () => false, - supportsQueueMessageImages: overrides.supportsQueueMessageImages, - abort: () => {}, - }; -} - -describe("embedded-agent runner steering admission", () => { +describe("embedded-agent active-run steering", () => { afterEach(() => { testing.resetActiveEmbeddedRuns(); replyRunTesting.resetReplyRunRegistry(); @@ -40,7 +27,7 @@ describe("embedded-agent runner steering admission", () => { it("passes steering options to active embedded runs", () => { const queueMessage = vi.fn(async () => {}); setActiveEmbeddedRun("session-steer", { - ...createSteeringRunHandle(), + ...createEmbeddedRunHandle(), sourceReplyDeliveryMode: "message_tool_only", queueMessage, }); @@ -61,7 +48,7 @@ describe("embedded-agent runner steering admission", () => { it("rejects images when the active run cannot preserve them", () => { const queueMessage = vi.fn(async () => {}); setActiveEmbeddedRun("session-images", { - ...createSteeringRunHandle(), + ...createEmbeddedRunHandle(), queueMessage, }); @@ -79,7 +66,7 @@ describe("embedded-agent runner steering admission", () => { setActiveEmbeddedRun( "session-images", - createSteeringRunHandle({ queueMessage, supportsQueueMessageImages: true }), + createEmbeddedRunHandle({ queueMessage, supportsQueueMessageImages: true }), ); expect( @@ -95,7 +82,7 @@ describe("embedded-agent runner steering admission", () => { it("rejects message-tool-only steering for active runs created without that mode", () => { const queueMessage = vi.fn(async () => {}); setActiveEmbeddedRun("session-automatic-source-reply", { - ...createSteeringRunHandle(), + ...createEmbeddedRunHandle(), queueMessage, }); @@ -131,7 +118,7 @@ describe("embedded-agent runner steering admission", () => { ])("rejects $label", ({ handleMode, requestMode }) => { const queueMessage = vi.fn(async () => {}); setActiveEmbeddedRun("session-task-suggestions", { - ...createSteeringRunHandle(), + ...createEmbeddedRunHandle(), taskSuggestionDeliveryMode: handleMode, queueMessage, }); @@ -153,7 +140,7 @@ describe("embedded-agent runner steering admission", () => { it("defaults active embedded steering to all pending messages", () => { const queueMessage = vi.fn(async () => {}); setActiveEmbeddedRun("session-default-steer", { - ...createSteeringRunHandle(), + ...createEmbeddedRunHandle(), queueMessage, }); @@ -168,7 +155,7 @@ describe("embedded-agent runner steering admission", () => { const queueMessage = vi.fn(async () => {}); setActiveEmbeddedRun( "session-active-non-streaming", - createSteeringRunHandle({ + createEmbeddedRunHandle({ isStreaming: false, isStopped: () => false, queueMessage, @@ -185,7 +172,7 @@ describe("embedded-agent runner steering admission", () => { vi.useFakeTimers(); try { const queueMessage = vi.fn(async () => {}); - setActiveEmbeddedRun("session-stale-steer", createSteeringRunHandle({ queueMessage })); + setActiveEmbeddedRun("session-stale-steer", createEmbeddedRunHandle({ queueMessage })); vi.advanceTimersByTime(10 * 60_000 + 1); @@ -207,7 +194,7 @@ describe("embedded-agent runner steering admission", () => { vi.useFakeTimers(); try { const queueMessage = vi.fn(async () => {}); - setActiveEmbeddedRun("session-quiet-tool-steer", createSteeringRunHandle({ queueMessage })); + setActiveEmbeddedRun("session-quiet-tool-steer", createEmbeddedRunHandle({ queueMessage })); markDiagnosticToolStartedForTest({ sessionId: "session-quiet-tool-steer", toolName: "exec", @@ -261,7 +248,7 @@ describe("embedded-agent runner steering admission", () => { const freshQueueMessage = vi.fn(async () => {}); setActiveEmbeddedRun( "session-fresh-steer", - createSteeringRunHandle({ queueMessage: freshQueueMessage }), + createEmbeddedRunHandle({ queueMessage: freshQueueMessage }), ); expect(queueEmbeddedAgentMessageWithOutcome("session-fresh-steer", "continue").queued).toBe( @@ -272,7 +259,7 @@ describe("embedded-agent runner steering admission", () => { const missingSnapshotQueueMessage = vi.fn(async () => {}); setActiveEmbeddedRun( "session-no-diagnostic-snapshot", - createSteeringRunHandle({ queueMessage: missingSnapshotQueueMessage }), + createEmbeddedRunHandle({ queueMessage: missingSnapshotQueueMessage }), ); resetDiagnosticRunActivityForTest(); @@ -286,7 +273,7 @@ describe("embedded-agent runner steering admission", () => { const queueMessage = vi.fn(async () => {}); setActiveEmbeddedRun( "session-stopped", - createSteeringRunHandle({ + createEmbeddedRunHandle({ isStreaming: true, isStopped: () => true, queueMessage, @@ -308,7 +295,7 @@ describe("embedded-agent runner steering admission", () => { const queueMessage = vi.fn(async () => {}); setActiveEmbeddedRun( "session-bad-state", - createSteeringRunHandle({ + createEmbeddedRunHandle({ isStopped: () => { throw new Error("bad stopped state"); }, @@ -326,4 +313,151 @@ describe("embedded-agent runner steering admission", () => { }); expect(queueMessage).not.toHaveBeenCalled(); }); + + it("returns a structured no-active-run queue failure", () => { + const outcome = queueEmbeddedAgentMessageWithOutcome("session-missing", "continue"); + + expect(outcome).toEqual({ + queued: false, + sessionId: "session-missing", + reason: "no_active_run", + gatewayHealth: "live", + }); + expect(formatEmbeddedAgentQueueFailureSummary(outcome)).toBe( + "queue_message_failed reason=no_active_run sessionId=session-missing gatewayHealth=live", + ); + }); + + it("returns structured queue failures for legacy, unavailable, or compacting runs", () => { + const legacyQueue = vi.fn(async () => {}); + const unavailableQueue = vi.fn(async () => {}); + setActiveEmbeddedRun( + "session-not-streaming", + createEmbeddedRunHandle({ isStreaming: false, queueMessage: legacyQueue }), + ); + setActiveEmbeddedRun( + "session-unavailable", + createEmbeddedRunHandle({ + messageInjection: { isAvailable: () => false, queueMessage: unavailableQueue }, + }), + ); + setActiveEmbeddedRun("session-compacting", createEmbeddedRunHandle({ isCompacting: true })); + + expect(queueEmbeddedAgentMessageWithOutcome("session-not-streaming", "continue")).toMatchObject( + { queued: false, reason: "not_streaming" }, + ); + expect(legacyQueue).not.toHaveBeenCalled(); + expect(queueEmbeddedAgentMessageWithOutcome("session-unavailable", "continue")).toMatchObject({ + queued: false, + reason: "not_streaming", + }); + expect(unavailableQueue).not.toHaveBeenCalled(); + expect(queueEmbeddedAgentMessageWithOutcome("session-compacting", "continue")).toMatchObject({ + queued: false, + reason: "compacting", + }); + }); + + it("returns runtime rejection details when async queue delivery fails", async () => { + setActiveEmbeddedRun("session-rejected", { + ...createEmbeddedRunHandle(), + queueMessage: async () => { + throw new Error("cannot steer a compact turn"); + }, + }); + + const outcome = await queueEmbeddedAgentMessageWithOutcomeAsync("session-rejected", "continue"); + + expect(outcome).toEqual({ + queued: false, + sessionId: "session-rejected", + reason: "runtime_rejected", + gatewayHealth: "live", + errorMessage: "cannot steer a compact turn", + }); + expect(formatEmbeddedAgentQueueFailureSummary(outcome)).toBe( + "queue_message_failed reason=runtime_rejected sessionId=session-rejected gatewayHealth=live error=cannot steer a compact turn", + ); + }); + + it("reports accepted steering without transcript confirmation as non-replayable", async () => { + setActiveEmbeddedRun("session-unconfirmed", { + ...createEmbeddedRunHandle(), + queueMessage: async () => ({ + transcriptCommit: "unconfirmed", + errorMessage: "receipt unavailable", + }), + }); + + const outcome = await queueEmbeddedAgentMessageWithOutcomeAsync( + "session-unconfirmed", + "continue", + ); + + expect(outcome).toEqual({ + queued: true, + sessionId: "session-unconfirmed", + target: "embedded_run", + gatewayHealth: "live", + transcriptCommit: "unconfirmed", + errorMessage: "receipt unavailable", + enqueuedAtMs: expect.any(Number), + }); + }); + + it("rejects transcript-commit waits for active handles without support", async () => { + const queueMessage = vi.fn(async () => {}); + setActiveEmbeddedRun("session-no-transcript-wait", { + ...createEmbeddedRunHandle(), + queueMessage, + }); + + const outcome = await queueEmbeddedAgentMessageWithOutcomeAsync( + "session-no-transcript-wait", + "continue", + { waitForTranscriptCommit: true }, + ); + + expect(outcome).toEqual({ + queued: false, + sessionId: "session-no-transcript-wait", + reason: "transcript_commit_wait_unsupported", + gatewayHealth: "live", + }); + expect(queueMessage).not.toHaveBeenCalled(); + }); + + it("rejects transcript-commit waits before reply-run fallback without an active handle", async () => { + const queueMessage = vi.fn(async () => {}); + const operation = createReplyOperation({ + sessionKey: "agent:main:main", + sessionId: "session-reply-run", + resetTriggered: false, + }); + operation.attachBackend({ + kind: "embedded", + cancel: vi.fn(), + isStreaming: () => true, + queueMessage, + }); + operation.setPhase("running"); + const recorder = createUserTurnTranscriptRecorder({ + input: { text: "visible group prompt", sender: { id: "user-42" } }, + target: createTestUserTurnTranscriptTarget(), + }); + + const outcome = await queueEmbeddedAgentMessageWithOutcomeAsync( + "session-reply-run", + "completion from child", + { waitForTranscriptCommit: true, userTurnTranscriptRecorder: recorder }, + ); + + expect(outcome).toEqual({ + queued: false, + sessionId: "session-reply-run", + reason: "transcript_commit_wait_unsupported", + gatewayHealth: "live", + }); + expect(queueMessage).not.toHaveBeenCalled(); + }); }); diff --git a/src/agents/embedded-agent-runner/runs.test-support.ts b/src/agents/embedded-agent-runner/runs.test-support.ts index e14e6ce14bac..bee10b38d951 100644 --- a/src/agents/embedded-agent-runner/runs.test-support.ts +++ b/src/agents/embedded-agent-runner/runs.test-support.ts @@ -1,4 +1,40 @@ import "./runs.js"; +import type { EmbeddedAgentQueueHandle } from "./run-state.js"; + +type RunHandle = EmbeddedAgentQueueHandle; + +export function createEmbeddedRunHandle( + overrides: { + abort?: () => void; + isAbortable?: boolean; + isCompacting?: boolean; + isStreaming?: boolean; + isStopped?: () => boolean; + messageInjection?: RunHandle["messageInjection"]; + runId?: string; + queueMessage?: RunHandle["queueMessage"]; + supportsQueueMessageImages?: boolean; + supportsTranscriptCommitWait?: boolean; + } = {}, +): RunHandle { + // Minimal handle fixture with overrideable lifecycle probes for registry + // behavior; individual tests supply queue/abort behavior when needed. + const abort = overrides.abort ?? (() => {}); + return { + runId: overrides.runId, + queueMessage: overrides.queueMessage ?? (async () => {}), + ...(overrides.messageInjection ? { messageInjection: overrides.messageInjection } : {}), + isStreaming: () => overrides.isStreaming ?? true, + ...(overrides.isStopped ? { isStopped: overrides.isStopped } : {}), + ...(overrides.isAbortable !== undefined + ? { isAbortable: () => overrides.isAbortable !== false } + : {}), + isCompacting: () => overrides.isCompacting ?? false, + supportsQueueMessageImages: overrides.supportsQueueMessageImages, + supportsTranscriptCommitWait: overrides.supportsTranscriptCommitWait, + abort, + }; +} type EmbeddedRunsTestApi = { persistForceClearedEmbeddedRunTerminalState(params: { diff --git a/src/agents/embedded-agent-runner/runs.test.ts b/src/agents/embedded-agent-runner/runs.test.ts index f89cb18e9f65..24166e35341d 100644 --- a/src/agents/embedded-agent-runner/runs.test.ts +++ b/src/agents/embedded-agent-runner/runs.test.ts @@ -11,8 +11,6 @@ import { getDiagnosticSessionState, resetDiagnosticSessionStateForTest, } from "../../logging/diagnostic-session-state.js"; -import { createUserTurnTranscriptRecorder } from "../../sessions/user-turn-transcript.js"; -import { createTestUserTurnTranscriptTarget } from "../../sessions/user-turn-transcript.test-support.js"; import { abortAndDrainEmbeddedAgentRun, abortEmbeddedAgentRun, @@ -21,48 +19,10 @@ import { isEmbeddedAgentRunAbortableForRunId, isEmbeddedAgentRunAbortableForCompaction, isEmbeddedAgentRunHandleActive, - formatEmbeddedAgentQueueFailureSummary, - queueEmbeddedAgentMessageWithOutcome, - queueEmbeddedAgentMessageWithOutcomeAsync, retainEmbeddedAgentRunAbortabilityForRunId, setActiveEmbeddedRun, } from "./runs.js"; -import { testing } from "./runs.test-support.js"; - -type RunHandle = Parameters[1]; - -function createRunHandle( - overrides: { - abort?: () => void; - isAbortable?: boolean; - isCompacting?: boolean; - isStreaming?: boolean; - isStopped?: () => boolean; - messageInjection?: RunHandle["messageInjection"]; - runId?: string; - queueMessage?: RunHandle["queueMessage"]; - supportsQueueMessageImages?: boolean; - supportsTranscriptCommitWait?: boolean; - } = {}, -): RunHandle { - // Minimal handle fixture with overrideable lifecycle probes for registry - // behavior; individual tests supply queue/abort behavior when needed. - const abort = overrides.abort ?? (() => {}); - return { - runId: overrides.runId, - queueMessage: overrides.queueMessage ?? (async () => {}), - ...(overrides.messageInjection ? { messageInjection: overrides.messageInjection } : {}), - isStreaming: () => overrides.isStreaming ?? true, - ...(overrides.isStopped ? { isStopped: overrides.isStopped } : {}), - ...(overrides.isAbortable !== undefined - ? { isAbortable: () => overrides.isAbortable !== false } - : {}), - isCompacting: () => overrides.isCompacting ?? false, - supportsQueueMessageImages: overrides.supportsQueueMessageImages, - supportsTranscriptCommitWait: overrides.supportsTranscriptCommitWait, - abort, - }; -} +import { createEmbeddedRunHandle, testing } from "./runs.test-support.js"; describe("embedded-agent runner run registry", () => { afterEach(() => { @@ -81,10 +41,10 @@ describe("embedded-agent runner run registry", () => { setActiveEmbeddedRun( "session-compacting", - createRunHandle({ isCompacting: true, abort: abortCompacting }), + createEmbeddedRunHandle({ isCompacting: true, abort: abortCompacting }), ); - setActiveEmbeddedRun("session-normal", createRunHandle({ abort: abortNormal })); + setActiveEmbeddedRun("session-normal", createEmbeddedRunHandle({ abort: abortNormal })); const aborted = abortEmbeddedAgentRun(undefined, { mode: "compacting" }); expect(aborted).toBe(true); @@ -110,9 +70,12 @@ describe("embedded-agent runner run registry", () => { const abortA = vi.fn(); const abortB = vi.fn(); - setActiveEmbeddedRun("session-a", createRunHandle({ isCompacting: true, abort: abortA })); + setActiveEmbeddedRun( + "session-a", + createEmbeddedRunHandle({ isCompacting: true, abort: abortA }), + ); - setActiveEmbeddedRun("session-b", createRunHandle({ abort: abortB })); + setActiveEmbeddedRun("session-b", createEmbeddedRunHandle({ abort: abortB })); const aborted = abortEmbeddedAgentRun(undefined, { mode: "all" }); expect(aborted).toBe(true); @@ -122,7 +85,7 @@ describe("embedded-agent runner run registry", () => { it("keeps finalizing runs active while rejecting abort requests", () => { const abort = vi.fn(); - const handle = createRunHandle({ abort, isAbortable: false }); + const handle = createEmbeddedRunHandle({ abort, isAbortable: false }); const operation = createReplyOperation({ sessionKey: "agent:main:finalizing", sessionId: "session-finalizing", @@ -155,7 +118,7 @@ describe("embedded-agent runner run registry", () => { it("keeps frozen run ownership through forced in-process restart", () => { const abort = vi.fn(); - const handle = createRunHandle({ abort, isAbortable: false }); + const handle = createEmbeddedRunHandle({ abort, isAbortable: false }); const operation = createReplyOperation({ sessionKey: "agent:main:restart-finalizing", sessionId: "session-restart-finalizing", @@ -185,7 +148,7 @@ describe("embedded-agent runner run registry", () => { }); it("binds abortability to the owning run id", () => { - const finalizing = createRunHandle({ + const finalizing = createEmbeddedRunHandle({ abort: vi.fn(), isAbortable: false, runId: "run-finalizing", @@ -203,7 +166,7 @@ describe("embedded-agent runner run registry", () => { clearActiveEmbeddedRun("session-shared", finalizing); expect(isEmbeddedAgentRunAbortableForRunId("run-finalizing")).toBe(false); - const queued = createRunHandle({ runId: "run-queued" }); + const queued = createEmbeddedRunHandle({ runId: "run-queued" }); setActiveEmbeddedRun("session-shared", queued); expect(isEmbeddedAgentRunAbortableForRunId("run-finalizing")).toBe(false); @@ -215,7 +178,7 @@ describe("embedded-agent runner run registry", () => { it("passes restart ownership to every aborted run", () => { const abort = vi.fn(); - setActiveEmbeddedRun("session-restart", createRunHandle({ abort })); + setActiveEmbeddedRun("session-restart", createEmbeddedRunHandle({ abort })); expect(abortEmbeddedAgentRun(undefined, { mode: "all", reason: "restart" })).toBe(true); expect(abort).toHaveBeenCalledWith("restart"); @@ -257,7 +220,7 @@ describe("embedded-agent runner run registry", () => { sessionId: "session-reply-stuck-live", resetTriggered: false, }); - const handle = createRunHandle({ + const handle = createEmbeddedRunHandle({ abort: () => { operation.abortByUser(); }, @@ -284,7 +247,7 @@ describe("embedded-agent runner run registry", () => { it("claims shared restart ownership before invoking an attached handle", () => { const abort = vi.fn(); - const handle = createRunHandle({ abort }); + const handle = createEmbeddedRunHandle({ abort }); const operation = createReplyOperation({ sessionKey: "agent:main:restart-owned", sessionId: "session-restart-owned", @@ -309,7 +272,7 @@ describe("embedded-agent runner run registry", () => { "does not bypass frozen shared ownership through %s handle aborts", (mode) => { const abort = vi.fn(); - const handle = createRunHandle({ abort, isCompacting: true }); + const handle = createEmbeddedRunHandle({ abort, isCompacting: true }); const sessionId = `session-restart-frozen-${mode}`; const operation = createReplyOperation({ sessionKey: `agent:main:restart-frozen-${mode}`, @@ -337,7 +300,7 @@ describe("embedded-agent runner run registry", () => { const abort = vi.fn(() => { throw new Error("cancel failed"); }); - const handle = createRunHandle({ abort }); + const handle = createEmbeddedRunHandle({ abort }); const operation = createReplyOperation({ sessionKey: "agent:main:restart-throwing", sessionId: "session-restart-throwing", @@ -359,7 +322,7 @@ describe("embedded-agent runner run registry", () => { it("does not bypass retained terminal ownership through compacting handle aborts", () => { const abort = vi.fn(); - const handle = createRunHandle({ abort, isCompacting: true }); + const handle = createEmbeddedRunHandle({ abort, isCompacting: true }); const operation = createReplyOperation({ sessionKey: "agent:main:restart-failed-compacting", sessionId: "session-restart-failed-compacting", @@ -385,7 +348,7 @@ describe("embedded-agent runner run registry", () => { it("records active run session files in diagnostic state for heartbeat recovery", () => { setDiagnosticsEnabledForProcess(true); const sessionFile = "/tmp/openclaw-run-registry-session.jsonl"; - const handle = createRunHandle(); + const handle = createEmbeddedRunHandle(); setActiveEmbeddedRun("session-file-diagnostics", handle, "agent:main:visible", sessionFile); @@ -393,151 +356,4 @@ describe("embedded-agent runner run registry", () => { sessionFile, ); }); - - it("returns a structured no-active-run queue failure", () => { - const outcome = queueEmbeddedAgentMessageWithOutcome("session-missing", "continue"); - - expect(outcome).toEqual({ - queued: false, - sessionId: "session-missing", - reason: "no_active_run", - gatewayHealth: "live", - }); - expect(formatEmbeddedAgentQueueFailureSummary(outcome)).toBe( - "queue_message_failed reason=no_active_run sessionId=session-missing gatewayHealth=live", - ); - }); - - it("returns structured queue failures for legacy, unavailable, or compacting runs", () => { - const legacyQueue = vi.fn(async () => {}); - const unavailableQueue = vi.fn(async () => {}); - setActiveEmbeddedRun( - "session-not-streaming", - createRunHandle({ isStreaming: false, queueMessage: legacyQueue }), - ); - setActiveEmbeddedRun( - "session-unavailable", - createRunHandle({ - messageInjection: { isAvailable: () => false, queueMessage: unavailableQueue }, - }), - ); - setActiveEmbeddedRun("session-compacting", createRunHandle({ isCompacting: true })); - - expect(queueEmbeddedAgentMessageWithOutcome("session-not-streaming", "continue")).toMatchObject( - { queued: false, reason: "not_streaming" }, - ); - expect(legacyQueue).not.toHaveBeenCalled(); - expect(queueEmbeddedAgentMessageWithOutcome("session-unavailable", "continue")).toMatchObject({ - queued: false, - reason: "not_streaming", - }); - expect(unavailableQueue).not.toHaveBeenCalled(); - expect(queueEmbeddedAgentMessageWithOutcome("session-compacting", "continue")).toMatchObject({ - queued: false, - reason: "compacting", - }); - }); - - it("returns runtime rejection details when async queue delivery fails", async () => { - setActiveEmbeddedRun("session-rejected", { - ...createRunHandle(), - queueMessage: async () => { - throw new Error("cannot steer a compact turn"); - }, - }); - - const outcome = await queueEmbeddedAgentMessageWithOutcomeAsync("session-rejected", "continue"); - - expect(outcome).toEqual({ - queued: false, - sessionId: "session-rejected", - reason: "runtime_rejected", - gatewayHealth: "live", - errorMessage: "cannot steer a compact turn", - }); - expect(formatEmbeddedAgentQueueFailureSummary(outcome)).toBe( - "queue_message_failed reason=runtime_rejected sessionId=session-rejected gatewayHealth=live error=cannot steer a compact turn", - ); - }); - - it("reports accepted steering without transcript confirmation as non-replayable", async () => { - setActiveEmbeddedRun("session-unconfirmed", { - ...createRunHandle(), - queueMessage: async () => ({ - transcriptCommit: "unconfirmed", - errorMessage: "receipt unavailable", - }), - }); - - const outcome = await queueEmbeddedAgentMessageWithOutcomeAsync( - "session-unconfirmed", - "continue", - ); - - expect(outcome).toEqual({ - queued: true, - sessionId: "session-unconfirmed", - target: "embedded_run", - gatewayHealth: "live", - transcriptCommit: "unconfirmed", - errorMessage: "receipt unavailable", - enqueuedAtMs: expect.any(Number), - }); - }); - - it("rejects transcript-commit waits for active handles without support", async () => { - const queueMessage = vi.fn(async () => {}); - setActiveEmbeddedRun("session-no-transcript-wait", { - ...createRunHandle(), - queueMessage, - }); - - const outcome = await queueEmbeddedAgentMessageWithOutcomeAsync( - "session-no-transcript-wait", - "continue", - { waitForTranscriptCommit: true }, - ); - - expect(outcome).toEqual({ - queued: false, - sessionId: "session-no-transcript-wait", - reason: "transcript_commit_wait_unsupported", - gatewayHealth: "live", - }); - expect(queueMessage).not.toHaveBeenCalled(); - }); - - it("rejects transcript-commit waits before reply-run fallback without an active handle", async () => { - const queueMessage = vi.fn(async () => {}); - const operation = createReplyOperation({ - sessionKey: "agent:main:main", - sessionId: "session-reply-run", - resetTriggered: false, - }); - operation.attachBackend({ - kind: "embedded", - cancel: vi.fn(), - isStreaming: () => true, - queueMessage, - }); - operation.setPhase("running"); - const recorder = createUserTurnTranscriptRecorder({ - input: { text: "visible group prompt", sender: { id: "user-42" } }, - target: createTestUserTurnTranscriptTarget(), - }); - - const outcome = await queueEmbeddedAgentMessageWithOutcomeAsync( - "session-reply-run", - "completion from child", - { waitForTranscriptCommit: true, userTurnTranscriptRecorder: recorder }, - ); - - expect(outcome).toEqual({ - queued: false, - sessionId: "session-reply-run", - reason: "transcript_commit_wait_unsupported", - gatewayHealth: "live", - }); - expect(queueMessage).not.toHaveBeenCalled(); - }); }); diff --git a/src/agents/sessions/agent-session-loop-correctness.test-support.ts b/src/agents/sessions/agent-session-loop-correctness.test-support.ts new file mode 100644 index 000000000000..2ea4b6efe9d0 --- /dev/null +++ b/src/agents/sessions/agent-session-loop-correctness.test-support.ts @@ -0,0 +1,169 @@ +import { + createAssistantMessageEventStream, + type AssistantMessage, + type Model, +} from "openclaw/plugin-sdk/llm"; +import { afterEach, beforeEach, vi } from "vitest"; +import { createResourceLoader } from "./agent-session-loop-resource-loader.test-support.js"; +import { AgentSession } from "./agent-session.js"; +import { AuthStorage } from "./auth-storage.js"; +import type { ToolDefinition } from "./extensions/types.js"; +import { ModelRegistry } from "./model-registry.js"; +import type { ResourceLoader } from "./resource-loader.js"; +import { createAgentSession, createAgentSessionForEmbeddedRunner } from "./sdk.js"; +import { SessionManager } from "./session-manager.js"; +import { SettingsManager } from "./settings-manager.js"; + +const hoistedStreamMocks = vi.hoisted(() => ({ + streamSimple: vi.fn(), +})); + +export const streamMocks = hoistedStreamMocks; + +export const testModel: Model = { + id: "test-model", + name: "Test Model", + api: "openai-responses", + provider: "test-provider", + baseUrl: "https://example.test", + reasoning: false, + input: ["text"], + cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 }, + contextWindow: 100, + maxTokens: 100, +}; + +const sessions: AgentSession[] = []; + +function createUsage(contextTokens: number) { + return { + input: contextTokens, + output: 1, + cacheRead: 0, + cacheWrite: 0, + totalTokens: contextTokens + 1, + contextUsage: { + state: "available" as const, + promptTokens: contextTokens, + totalTokens: contextTokens + 1, + }, + cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 }, + }; +} + +export function createAssistant( + activeModel: Model, + content: AssistantMessage["content"], + stopReason: AssistantMessage["stopReason"] = "stop", + contextTokens = 1, +): AssistantMessage { + return { + role: "assistant", + content, + api: activeModel.api, + provider: activeModel.provider, + model: activeModel.id, + usage: createUsage(contextTokens), + stopReason, + timestamp: Date.now(), + }; +} + +export function createAssistantResultStream(message: AssistantMessage) { + const stream = createAssistantMessageEventStream(); + queueMicrotask(() => { + if (message.stopReason === "error" || message.stopReason === "aborted") { + stream.push({ type: "error", reason: message.stopReason, error: message }); + } else { + stream.push({ type: "done", reason: message.stopReason, message }); + } + stream.end(); + }); + return stream; +} + +export const createOverflowAssistant = (activeModel: Model) => ({ + ...createAssistant(activeModel, [{ type: "text", text: "truncated answer" }], "length", 100), + usage: { ...createUsage(100), output: 0 }, +}); + +export const createAutoCompactionSettings = () => + SettingsManager.inMemory({ + compaction: { enabled: true, reserveTokens: 0, keepRecentTokens: 1 }, + retry: { enabled: false }, + }); + +export function mockInvalidThenTextSummary(recoveredText: string) { + let requests = 0; + streamMocks.streamSimple.mockImplementation((activeModel: Model) => { + return createAssistantResultStream( + createAssistant( + activeModel, + ++requests === 1 + ? [{ type: "thinking", thinking: "internal summary reasoning" }] + : [{ type: "text", text: recoveredText }], + ), + ); + }); + return () => requests; +} + +export async function createTestSession( + options: { + model?: Model; + settingsManager?: SettingsManager; + sessionManager?: SessionManager; + resourceLoader?: ResourceLoader; + customTools?: ToolDefinition[]; + contextOverflowRecoveryOwner?: "session" | "caller"; + } = {}, +) { + const model = options.model ?? testModel; + const authStorage = AuthStorage.inMemory(); + authStorage.setRuntimeApiKey(model.provider, "test-api-key"); + const settingsManager = + options.settingsManager ?? + SettingsManager.inMemory({ + compaction: { enabled: false }, + retry: { enabled: false }, + }); + const sessionManager = options.sessionManager ?? SessionManager.inMemory(); + const modelRegistry = ModelRegistry.inMemory(authStorage); + modelRegistry.registerProvider(model.provider, { + api: model.api, + streamSimple: streamMocks.streamSimple, + }); + const sessionOptions = { + model, + noTools: "builtin" as const, + customTools: options.customTools, + resourceLoader: options.resourceLoader ?? createResourceLoader(), + sessionManager, + settingsManager, + modelRegistry, + }; + const result = options.contextOverflowRecoveryOwner + ? await createAgentSessionForEmbeddedRunner(sessionOptions, { + contextOverflowRecoveryOwner: options.contextOverflowRecoveryOwner, + }) + : await createAgentSession(sessionOptions); + sessions.push(result.session); + return { ...result, settingsManager, sessionManager }; +} + +export function appendHistory(sessionManager: SessionManager, assistant: AssistantMessage): void { + sessionManager.appendMessage({ role: "user", content: "old prompt", timestamp: Date.now() - 2 }); + sessionManager.appendMessage({ ...assistant, timestamp: Date.now() - 1 }); +} + +export function registerAgentSessionLoopTestLifecycle(): void { + beforeEach(() => { + streamMocks.streamSimple.mockReset(); + }); + + afterEach(() => { + for (const session of sessions.splice(0)) { + session.dispose(); + } + }); +} diff --git a/src/agents/sessions/agent-session-loop-correctness.test.ts b/src/agents/sessions/agent-session-loop-correctness.test.ts index e1ffa2dee8a7..7b47de81052d 100644 --- a/src/agents/sessions/agent-session-loop-correctness.test.ts +++ b/src/agents/sessions/agent-session-loop-correctness.test.ts @@ -1,178 +1,33 @@ import { createAssistantMessageEventStream, - type AssistantMessage, type Context, type Model, - type SimpleStreamOptions, } from "openclaw/plugin-sdk/llm"; import { Type } from "typebox"; -import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; - -const streamMocks = vi.hoisted(() => ({ - streamSimple: vi.fn(), -})); - -import type { AgentTool } from "../runtime/index.js"; +import { describe, expect, it, vi } from "vitest"; import { agentSessionAutomaticCompaction } from "./agent-session-compaction.js"; +import { + appendHistory, + createAssistant, + createAssistantResultStream, + createAutoCompactionSettings, + createOverflowAssistant, + createTestSession, + mockInvalidThenTextSummary, + registerAgentSessionLoopTestLifecycle, + streamMocks, + testModel, +} from "./agent-session-loop-correctness.test-support.js"; import { createCompactionHandlers, createResourceLoader, } from "./agent-session-loop-resource-loader.test-support.js"; import type { AgentSessionEvent } from "./agent-session-types.js"; -import { AgentSession } from "./agent-session.js"; -import { AuthStorage } from "./auth-storage.js"; import type { ToolDefinition } from "./extensions/types.js"; -import { ModelRegistry } from "./model-registry.js"; -import type { ResourceLoader } from "./resource-loader.js"; -import { createAgentSession, createAgentSessionForEmbeddedRunner } from "./sdk.js"; import { SessionManager } from "./session-manager.js"; import { SettingsManager } from "./settings-manager.js"; -const testModel: Model = { - id: "test-model", - name: "Test Model", - api: "openai-responses", - provider: "test-provider", - baseUrl: "https://example.test", - reasoning: false, - input: ["text"], - cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 }, - contextWindow: 100, - maxTokens: 100, -}; - -const sessions: AgentSession[] = []; - -function createUsage(contextTokens: number) { - return { - input: contextTokens, - output: 1, - cacheRead: 0, - cacheWrite: 0, - totalTokens: contextTokens + 1, - contextUsage: { - state: "available" as const, - promptTokens: contextTokens, - totalTokens: contextTokens + 1, - }, - cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 }, - }; -} - -function createAssistant( - activeModel: Model, - content: AssistantMessage["content"], - stopReason: AssistantMessage["stopReason"] = "stop", - contextTokens = 1, -): AssistantMessage { - return { - role: "assistant", - content, - api: activeModel.api, - provider: activeModel.provider, - model: activeModel.id, - usage: createUsage(contextTokens), - stopReason, - timestamp: Date.now(), - }; -} - -function createAssistantResultStream(message: AssistantMessage) { - const stream = createAssistantMessageEventStream(); - queueMicrotask(() => { - if (message.stopReason === "error" || message.stopReason === "aborted") { - stream.push({ type: "error", reason: message.stopReason, error: message }); - } else { - stream.push({ type: "done", reason: message.stopReason, message }); - } - stream.end(); - }); - return stream; -} - -const createOverflowAssistant = (activeModel: Model) => ({ - ...createAssistant(activeModel, [{ type: "text", text: "truncated answer" }], "length", 100), - usage: { ...createUsage(100), output: 0 }, -}); - -const createAutoCompactionSettings = () => - SettingsManager.inMemory({ - compaction: { enabled: true, reserveTokens: 0, keepRecentTokens: 1 }, - retry: { enabled: false }, - }); - -function mockInvalidThenTextSummary(recoveredText: string) { - let requests = 0; - streamMocks.streamSimple.mockImplementation((activeModel: Model) => { - return createAssistantResultStream( - createAssistant( - activeModel, - ++requests === 1 - ? [{ type: "thinking", thinking: "internal summary reasoning" }] - : [{ type: "text", text: recoveredText }], - ), - ); - }); - return () => requests; -} - -async function createTestSession( - options: { - model?: Model; - settingsManager?: SettingsManager; - sessionManager?: SessionManager; - resourceLoader?: ResourceLoader; - customTools?: ToolDefinition[]; - contextOverflowRecoveryOwner?: "session" | "caller"; - } = {}, -) { - const model = options.model ?? testModel; - const authStorage = AuthStorage.inMemory(); - authStorage.setRuntimeApiKey(model.provider, "test-api-key"); - const settingsManager = - options.settingsManager ?? - SettingsManager.inMemory({ - compaction: { enabled: false }, - retry: { enabled: false }, - }); - const sessionManager = options.sessionManager ?? SessionManager.inMemory(); - const modelRegistry = ModelRegistry.inMemory(authStorage); - modelRegistry.registerProvider(model.provider, { - api: model.api, - streamSimple: streamMocks.streamSimple, - }); - const sessionOptions = { - model, - noTools: "builtin" as const, - customTools: options.customTools, - resourceLoader: options.resourceLoader ?? createResourceLoader(), - sessionManager, - settingsManager, - modelRegistry, - }; - const result = options.contextOverflowRecoveryOwner - ? await createAgentSessionForEmbeddedRunner(sessionOptions, { - contextOverflowRecoveryOwner: options.contextOverflowRecoveryOwner, - }) - : await createAgentSession(sessionOptions); - sessions.push(result.session); - return { ...result, settingsManager, sessionManager }; -} - -function appendHistory(sessionManager: SessionManager, assistant: AssistantMessage): void { - sessionManager.appendMessage({ role: "user", content: "old prompt", timestamp: Date.now() - 2 }); - sessionManager.appendMessage({ ...assistant, timestamp: Date.now() - 1 }); -} - -beforeEach(() => { - streamMocks.streamSimple.mockReset(); -}); - -afterEach(() => { - for (const session of sessions.splice(0)) { - session.dispose(); - } -}); +registerAgentSessionLoopTestLifecycle(); describe("AgentSession loop correctness", () => { it("carries the canonical assistant entry id through ordered terminal listeners", async () => { @@ -780,303 +635,4 @@ describe("AgentSession loop correctness", () => { expect(requests).toHaveLength(1); expect(JSON.stringify(requests[0]?.messages)).toContain("pending prompt"); }); - - it("drains a follow-up queued by an agent-end handler", async () => { - const sessionRef: { current?: AgentSession } = {}; - let queued = false; - const lifecycleEvents: string[] = []; - const handlers = new Map Promise>>([ - [ - "agent_end", - [ - async () => { - lifecycleEvents.push("agent_end"); - if (!queued) { - queued = true; - await sessionRef.current?.followUp("queued after end"); - } - return undefined; - }, - ], - ], - ["agent_settled", [async () => lifecycleEvents.push("agent_settled")]], - ]); - const requests: Context[] = []; - streamMocks.streamSimple.mockImplementation((activeModel: Model, context: Context) => { - requests.push(context); - return createAssistantResultStream( - createAssistant(activeModel, [{ type: "text", text: `answer ${requests.length}` }]), - ); - }); - const { session } = await createTestSession({ resourceLoader: createResourceLoader(handlers) }); - sessionRef.current = session; - - await session.prompt("initial prompt"); - - expect(requests).toHaveLength(2); - expect(JSON.stringify(requests[1]?.messages)).toContain("queued after end"); - expect(session.agent.hasQueuedMessages()).toBe(false); - expect(lifecycleEvents).toEqual(["agent_end", "agent_end", "agent_settled"]); - }); - - it("leaves queued messages dormant after a turn handoff", async () => { - const sessionRef: { current?: AgentSession } = {}; - const settled = vi.fn(); - const handlers = new Map Promise>>([ - ["agent_settled", [async () => settled()]], - ]); - const yieldTool: ToolDefinition = { - name: "yield_turn", - label: "Yield turn", - description: "ends the current turn for an external handoff", - parameters: Type.Object({}), - execute: async () => { - const activeSession = sessionRef.current; - if (!activeSession) { - throw new Error("session not ready"); - } - activeSession.agent.steer({ - role: "custom", - customType: "test.turn-handoff", - content: "resume only for external delivery", - display: false, - timestamp: Date.now(), - }); - activeSession.agent.abort({ code: "turn_handoff", turnHandoff: true }); - return { content: [{ type: "text", text: "yielded" }], details: { yielded: true } }; - }, - }; - streamMocks.streamSimple.mockImplementation((activeModel: Model) => - createAssistantResultStream( - createAssistant( - activeModel, - [{ type: "toolCall", id: "call-yield", name: "yield_turn", arguments: {} }], - "toolUse", - ), - ), - ); - const { session } = await createTestSession({ - customTools: [yieldTool], - resourceLoader: createResourceLoader(handlers), - }); - sessionRef.current = session; - - await session.prompt("yield now"); - - expect(streamMocks.streamSimple).toHaveBeenCalledOnce(); - expect(session.agent.hasQueuedMessages()).toBe(true); - expect(settled).not.toHaveBeenCalled(); - session.agent.clearAllQueues(); - }); - - it("applies session model, tool, and prompt changes on the following tool turn", async () => { - const nextModel = { ...testModel, id: "next-model" }; - const sessionRef: { current?: AgentSession } = {}; - const switchTool: ToolDefinition = { - name: "switch_state", - label: "Switch state", - description: "changes the next turn state", - parameters: Type.Object({}), - execute: async () => { - const activeSession = sessionRef.current; - if (!activeSession) { - throw new Error("session not ready"); - } - activeSession.setActiveToolsByName(["second_tool"]); - activeSession.agent.state.model = nextModel; - return { content: [{ type: "text", text: "switched" }], details: {} }; - }, - }; - const secondTool: ToolDefinition = { - name: "second_tool", - label: "Second tool", - description: "available after the switch", - parameters: Type.Object({}), - execute: async () => ({ content: [{ type: "text", text: "done" }], details: {} }), - }; - const handlers = new Map Promise>>([ - ["before_agent_start", [async () => ({ systemPrompt: "prompt override" })]], - ]); - const requests: Array<{ model: string; prompt: string; tools: string[] }> = []; - streamMocks.streamSimple.mockImplementation((activeModel: Model, context: Context) => { - requests.push({ - model: activeModel.id, - prompt: context.systemPrompt ?? "", - tools: context.tools?.map((tool) => tool.name) ?? [], - }); - const content: AssistantMessage["content"] = - requests.length === 1 - ? [{ type: "toolCall", id: "call-switch", name: "switch_state", arguments: {} }] - : [{ type: "text", text: "finished" }]; - return createAssistantResultStream( - createAssistant(activeModel, content, requests.length === 1 ? "toolUse" : "stop"), - ); - }); - const { session } = await createTestSession({ - resourceLoader: createResourceLoader(handlers), - customTools: [switchTool, secondTool], - }); - sessionRef.current = session; - session.setActiveToolsByName(["switch_state"]); - - await session.prompt("switch now"); - - expect(requests).toEqual([ - { model: testModel.id, prompt: "prompt override", tools: ["switch_state"] }, - { model: nextModel.id, prompt: "prompt override", tools: ["second_tool"] }, - ]); - }); - - it("preserves explicit updates from an existing next-turn hook", async () => { - const hookModel = { ...testModel, id: "hook-model" }; - const hookTool: AgentTool = { - name: "hook_tool", - label: "Hook tool", - description: "provided by the existing turn hook", - parameters: Type.Object({}), - execute: async () => ({ content: [{ type: "text", text: "done" }], details: {} }), - }; - const hookContext = { - systemPrompt: "hook prompt", - messages: [], - tools: [hookTool], - }; - let returnedUpdate = false; - const { session } = await createTestSession(); - session.agent.prepareNextTurn = () => { - if (returnedUpdate) { - return undefined; - } - returnedUpdate = true; - return { context: hookContext, model: hookModel, thinkingLevel: "high" }; - }; - const contextualHook = session.agent.prepareNextTurnWithContext; - if (!contextualHook) { - throw new Error("context-aware next-turn hook was not installed"); - } - const message = createAssistant(testModel, [{ type: "text", text: "turn complete" }]); - const newMessages = [message]; - - const firstUpdate = await contextualHook({ - message, - toolResults: [], - context: { systemPrompt: "loop prompt", messages: [], tools: [] }, - newMessages, - }); - const secondUpdate = await contextualHook({ - message, - toolResults: [], - context: firstUpdate?.context ?? hookContext, - newMessages, - }); - - for (const update of [firstUpdate, secondUpdate]) { - expect(update).toMatchObject({ - context: { - systemPrompt: "hook prompt", - tools: [expect.objectContaining({ name: "hook_tool" })], - }, - model: hookModel, - thinkingLevel: "high", - }); - } - }); - - it("preserves fields omitted by an existing next-turn context replacement", async () => { - const sessionTool: AgentTool = { - name: "session_tool", - label: "Session tool", - description: "available in session state", - parameters: Type.Object({}), - execute: async () => ({ content: [{ type: "text", text: "done" }], details: {} }), - }; - const initialHook = vi.fn(() => ({ - context: { systemPrompt: "stale prompt", messages: [], tools: [sessionTool] }, - })); - const replacementHook = vi.fn(() => ({ - context: { systemPrompt: "replacement prompt", messages: [] }, - })); - const { session } = await createTestSession({ customTools: [sessionTool] }); - session.setActiveToolsByName([sessionTool.name]); - session.agent.prepareNextTurn = initialHook; - session.agent.prepareNextTurn = replacementHook; - const message = createAssistant(testModel, [{ type: "text", text: "turn complete" }]); - const contextualHook = session.agent.prepareNextTurnWithContext; - if (!contextualHook) { - throw new Error("context-aware next-turn hook was not installed"); - } - - const update = await contextualHook({ - message, - toolResults: [], - context: { systemPrompt: "loop prompt", messages: [], tools: [sessionTool] }, - newMessages: [message], - }); - - expect(update?.context).toEqual({ systemPrompt: "replacement prompt", messages: [] }); - expect(replacementHook).toHaveBeenCalledOnce(); - expect(initialHook).not.toHaveBeenCalled(); - }); - - it("aborts in-flight work when disposed", async () => { - let providerSignal: AbortSignal | undefined; - streamMocks.streamSimple.mockImplementation( - (activeModel: Model, _context: Context, options?: SimpleStreamOptions) => { - providerSignal = options?.signal; - const stream = createAssistantMessageEventStream(); - options?.signal?.addEventListener( - "abort", - () => { - const message = createAssistant(activeModel, [], "aborted"); - stream.push({ type: "error", reason: "aborted", error: message }); - stream.end(); - }, - { once: true }, - ); - return stream; - }, - ); - const { session } = await createTestSession(); - const abortRetry = vi.spyOn(session, "abortRetry"); - const abortCompaction = vi.spyOn(session, "abortCompaction"); - const abortBranchSummary = vi.spyOn(session, "abortBranchSummary"); - const abortBash = vi.spyOn(session, "abortBash"); - const abortAgent = vi.spyOn(session.agent, "abort"); - abortRetry.mockImplementationOnce(() => { - throw new Error("retry abort failed"); - }); - const prompt = session.prompt("wait"); - await vi.waitFor(() => expect(providerSignal).toBeDefined()); - - session.dispose(); - await prompt; - - expect(providerSignal?.aborted).toBe(true); - expect(abortRetry).toHaveBeenCalledOnce(); - expect(abortCompaction).toHaveBeenCalledOnce(); - expect(abortBranchSummary).toHaveBeenCalledOnce(); - expect(abortBash).toHaveBeenCalledOnce(); - expect(abortAgent).toHaveBeenCalledOnce(); - }); - - it("resynchronizes queue modes when settings reload", async () => { - const settingsManager = SettingsManager.inMemory({ - steeringMode: "one-at-a-time", - followUpMode: "one-at-a-time", - compaction: { enabled: false }, - retry: { enabled: false }, - }); - const { session } = await createTestSession({ settingsManager }); - settingsManager.setSteeringMode("all"); - settingsManager.setFollowUpMode("all"); - await settingsManager.flush(); - - expect(session.agent.steeringMode).toBe("one-at-a-time"); - expect(session.agent.followUpMode).toBe("one-at-a-time"); - - await session.reload(); - - expect(session.agent.steeringMode).toBe("all"); - expect(session.agent.followUpMode).toBe("all"); - }); }); diff --git a/src/agents/sessions/agent-session-loop-next-turn.test.ts b/src/agents/sessions/agent-session-loop-next-turn.test.ts new file mode 100644 index 000000000000..6528b741bde2 --- /dev/null +++ b/src/agents/sessions/agent-session-loop-next-turn.test.ts @@ -0,0 +1,325 @@ +import { + createAssistantMessageEventStream, + type AssistantMessage, + type Context, + type Model, + type SimpleStreamOptions, +} from "openclaw/plugin-sdk/llm"; +import { Type } from "typebox"; +import { describe, expect, it, vi } from "vitest"; +import type { AgentTool } from "../runtime/index.js"; +import { + createAssistant, + createAssistantResultStream, + createTestSession, + registerAgentSessionLoopTestLifecycle, + streamMocks, + testModel, +} from "./agent-session-loop-correctness.test-support.js"; +import { createResourceLoader } from "./agent-session-loop-resource-loader.test-support.js"; +import type { AgentSession } from "./agent-session.js"; +import type { ToolDefinition } from "./extensions/types.js"; +import { SettingsManager } from "./settings-manager.js"; + +registerAgentSessionLoopTestLifecycle(); + +describe("AgentSession queue and next-turn lifecycle correctness", () => { + it("drains a follow-up queued by an agent-end handler", async () => { + const sessionRef: { current?: AgentSession } = {}; + let queued = false; + const lifecycleEvents: string[] = []; + const handlers = new Map Promise>>([ + [ + "agent_end", + [ + async () => { + lifecycleEvents.push("agent_end"); + if (!queued) { + queued = true; + await sessionRef.current?.followUp("queued after end"); + } + return undefined; + }, + ], + ], + ["agent_settled", [async () => lifecycleEvents.push("agent_settled")]], + ]); + const requests: Context[] = []; + streamMocks.streamSimple.mockImplementation((activeModel: Model, context: Context) => { + requests.push(context); + return createAssistantResultStream( + createAssistant(activeModel, [{ type: "text", text: `answer ${requests.length}` }]), + ); + }); + const { session } = await createTestSession({ resourceLoader: createResourceLoader(handlers) }); + sessionRef.current = session; + + await session.prompt("initial prompt"); + + expect(requests).toHaveLength(2); + expect(JSON.stringify(requests[1]?.messages)).toContain("queued after end"); + expect(session.agent.hasQueuedMessages()).toBe(false); + expect(lifecycleEvents).toEqual(["agent_end", "agent_end", "agent_settled"]); + }); + + it("leaves queued messages dormant after a turn handoff", async () => { + const sessionRef: { current?: AgentSession } = {}; + const settled = vi.fn(); + const handlers = new Map Promise>>([ + ["agent_settled", [async () => settled()]], + ]); + const yieldTool: ToolDefinition = { + name: "yield_turn", + label: "Yield turn", + description: "ends the current turn for an external handoff", + parameters: Type.Object({}), + execute: async () => { + const activeSession = sessionRef.current; + if (!activeSession) { + throw new Error("session not ready"); + } + activeSession.agent.steer({ + role: "custom", + customType: "test.turn-handoff", + content: "resume only for external delivery", + display: false, + timestamp: Date.now(), + }); + activeSession.agent.abort({ code: "turn_handoff", turnHandoff: true }); + return { content: [{ type: "text", text: "yielded" }], details: { yielded: true } }; + }, + }; + streamMocks.streamSimple.mockImplementation((activeModel: Model) => + createAssistantResultStream( + createAssistant( + activeModel, + [{ type: "toolCall", id: "call-yield", name: "yield_turn", arguments: {} }], + "toolUse", + ), + ), + ); + const { session } = await createTestSession({ + customTools: [yieldTool], + resourceLoader: createResourceLoader(handlers), + }); + sessionRef.current = session; + + await session.prompt("yield now"); + + expect(streamMocks.streamSimple).toHaveBeenCalledOnce(); + expect(session.agent.hasQueuedMessages()).toBe(true); + expect(settled).not.toHaveBeenCalled(); + session.agent.clearAllQueues(); + }); + + it("applies session model, tool, and prompt changes on the following tool turn", async () => { + const nextModel = { ...testModel, id: "next-model" }; + const sessionRef: { current?: AgentSession } = {}; + const switchTool: ToolDefinition = { + name: "switch_state", + label: "Switch state", + description: "changes the next turn state", + parameters: Type.Object({}), + execute: async () => { + const activeSession = sessionRef.current; + if (!activeSession) { + throw new Error("session not ready"); + } + activeSession.setActiveToolsByName(["second_tool"]); + activeSession.agent.state.model = nextModel; + return { content: [{ type: "text", text: "switched" }], details: {} }; + }, + }; + const secondTool: ToolDefinition = { + name: "second_tool", + label: "Second tool", + description: "available after the switch", + parameters: Type.Object({}), + execute: async () => ({ content: [{ type: "text", text: "done" }], details: {} }), + }; + const handlers = new Map Promise>>([ + ["before_agent_start", [async () => ({ systemPrompt: "prompt override" })]], + ]); + const requests: Array<{ model: string; prompt: string; tools: string[] }> = []; + streamMocks.streamSimple.mockImplementation((activeModel: Model, context: Context) => { + requests.push({ + model: activeModel.id, + prompt: context.systemPrompt ?? "", + tools: context.tools?.map((tool) => tool.name) ?? [], + }); + const content: AssistantMessage["content"] = + requests.length === 1 + ? [{ type: "toolCall", id: "call-switch", name: "switch_state", arguments: {} }] + : [{ type: "text", text: "finished" }]; + return createAssistantResultStream( + createAssistant(activeModel, content, requests.length === 1 ? "toolUse" : "stop"), + ); + }); + const { session } = await createTestSession({ + resourceLoader: createResourceLoader(handlers), + customTools: [switchTool, secondTool], + }); + sessionRef.current = session; + session.setActiveToolsByName(["switch_state"]); + + await session.prompt("switch now"); + + expect(requests).toEqual([ + { model: testModel.id, prompt: "prompt override", tools: ["switch_state"] }, + { model: nextModel.id, prompt: "prompt override", tools: ["second_tool"] }, + ]); + }); + + it("preserves explicit updates from an existing next-turn hook", async () => { + const hookModel = { ...testModel, id: "hook-model" }; + const hookTool: AgentTool = { + name: "hook_tool", + label: "Hook tool", + description: "provided by the existing turn hook", + parameters: Type.Object({}), + execute: async () => ({ content: [{ type: "text", text: "done" }], details: {} }), + }; + const hookContext = { + systemPrompt: "hook prompt", + messages: [], + tools: [hookTool], + }; + let returnedUpdate = false; + const { session } = await createTestSession(); + session.agent.prepareNextTurn = () => { + if (returnedUpdate) { + return undefined; + } + returnedUpdate = true; + return { context: hookContext, model: hookModel, thinkingLevel: "high" }; + }; + const contextualHook = session.agent.prepareNextTurnWithContext; + if (!contextualHook) { + throw new Error("context-aware next-turn hook was not installed"); + } + const message = createAssistant(testModel, [{ type: "text", text: "turn complete" }]); + const newMessages = [message]; + + const firstUpdate = await contextualHook({ + message, + toolResults: [], + context: { systemPrompt: "loop prompt", messages: [], tools: [] }, + newMessages, + }); + const secondUpdate = await contextualHook({ + message, + toolResults: [], + context: firstUpdate?.context ?? hookContext, + newMessages, + }); + + for (const update of [firstUpdate, secondUpdate]) { + expect(update).toMatchObject({ + context: { + systemPrompt: "hook prompt", + tools: [expect.objectContaining({ name: "hook_tool" })], + }, + model: hookModel, + thinkingLevel: "high", + }); + } + }); + + it("preserves fields omitted by an existing next-turn context replacement", async () => { + const sessionTool: AgentTool = { + name: "session_tool", + label: "Session tool", + description: "available in session state", + parameters: Type.Object({}), + execute: async () => ({ content: [{ type: "text", text: "done" }], details: {} }), + }; + const initialHook = vi.fn(() => ({ + context: { systemPrompt: "stale prompt", messages: [], tools: [sessionTool] }, + })); + const replacementHook = vi.fn(() => ({ + context: { systemPrompt: "replacement prompt", messages: [] }, + })); + const { session } = await createTestSession({ customTools: [sessionTool] }); + session.setActiveToolsByName([sessionTool.name]); + session.agent.prepareNextTurn = initialHook; + session.agent.prepareNextTurn = replacementHook; + const message = createAssistant(testModel, [{ type: "text", text: "turn complete" }]); + const contextualHook = session.agent.prepareNextTurnWithContext; + if (!contextualHook) { + throw new Error("context-aware next-turn hook was not installed"); + } + + const update = await contextualHook({ + message, + toolResults: [], + context: { systemPrompt: "loop prompt", messages: [], tools: [sessionTool] }, + newMessages: [message], + }); + + expect(update?.context).toEqual({ systemPrompt: "replacement prompt", messages: [] }); + expect(replacementHook).toHaveBeenCalledOnce(); + expect(initialHook).not.toHaveBeenCalled(); + }); + + it("aborts in-flight work when disposed", async () => { + let providerSignal: AbortSignal | undefined; + streamMocks.streamSimple.mockImplementation( + (activeModel: Model, _context: Context, options?: SimpleStreamOptions) => { + providerSignal = options?.signal; + const stream = createAssistantMessageEventStream(); + options?.signal?.addEventListener( + "abort", + () => { + const message = createAssistant(activeModel, [], "aborted"); + stream.push({ type: "error", reason: "aborted", error: message }); + stream.end(); + }, + { once: true }, + ); + return stream; + }, + ); + const { session } = await createTestSession(); + const abortRetry = vi.spyOn(session, "abortRetry"); + const abortCompaction = vi.spyOn(session, "abortCompaction"); + const abortBranchSummary = vi.spyOn(session, "abortBranchSummary"); + const abortBash = vi.spyOn(session, "abortBash"); + const abortAgent = vi.spyOn(session.agent, "abort"); + abortRetry.mockImplementationOnce(() => { + throw new Error("retry abort failed"); + }); + const prompt = session.prompt("wait"); + await vi.waitFor(() => expect(providerSignal).toBeDefined()); + + session.dispose(); + await prompt; + + expect(providerSignal?.aborted).toBe(true); + expect(abortRetry).toHaveBeenCalledOnce(); + expect(abortCompaction).toHaveBeenCalledOnce(); + expect(abortBranchSummary).toHaveBeenCalledOnce(); + expect(abortBash).toHaveBeenCalledOnce(); + expect(abortAgent).toHaveBeenCalledOnce(); + }); + + it("resynchronizes queue modes when settings reload", async () => { + const settingsManager = SettingsManager.inMemory({ + steeringMode: "one-at-a-time", + followUpMode: "one-at-a-time", + compaction: { enabled: false }, + retry: { enabled: false }, + }); + const { session } = await createTestSession({ settingsManager }); + settingsManager.setSteeringMode("all"); + settingsManager.setFollowUpMode("all"); + await settingsManager.flush(); + + expect(session.agent.steeringMode).toBe("one-at-a-time"); + expect(session.agent.followUpMode).toBe("one-at-a-time"); + + await session.reload(); + + expect(session.agent.steeringMode).toBe("all"); + expect(session.agent.followUpMode).toBe("all"); + }); +}); diff --git a/src/auto-reply/reply/commands-context-report.test.ts b/src/auto-reply/reply/commands-context-report.test.ts index 080dda3b7c40..63f2eb9092e1 100644 --- a/src/auto-reply/reply/commands-context-report.test.ts +++ b/src/auto-reply/reply/commands-context-report.test.ts @@ -32,6 +32,7 @@ function makeParams( currentTurn?: NonNullable["currentTurn"]; }, ): HandleCommandsParams { + const totalTokensFresh = options?.totalTokensFresh ?? true; return { command: { commandBodyNormalized, @@ -50,8 +51,8 @@ function makeParams( sessionEntry: { ...(options?.sessionId ? { sessionId: options.sessionId } : {}), totalTokens: options?.totalTokens ?? 123, - totalTokensFresh: options?.totalTokensFresh ?? true, - totalTokensVersion: 1 as const, + totalTokensFresh, + ...(totalTokensFresh ? { totalTokensVersion: 1 as const } : {}), inputTokens: 100, outputTokens: 23, systemPromptReport: { diff --git a/ui/src/components/app-sidebar-render.ts b/ui/src/components/app-sidebar-render.ts index 49d91edc8dee..c20f907f84c6 100644 --- a/ui/src/components/app-sidebar-render.ts +++ b/ui/src/components/app-sidebar-render.ts @@ -9,6 +9,7 @@ import { isSessionRouteId } from "../app-route-paths.ts"; import { sessionHasPendingApproval } from "../app/approval-presentation.ts"; import { isNativeWebChromeHost } from "../app/native-web-chrome.ts"; import { readPresenceEntries, resolveCurrentSelfUser } from "../app/user-profile.ts"; +import { CONTROL_UI_BUILD_INFO } from "../build-info.ts"; import { t } from "../i18n/index.ts"; import { normalizeAgentLabel, resolveAgentTextAvatar } from "../lib/agents/display.ts"; import { deriveAvatarInitial, resolveAgentAvatarUrl } from "../lib/avatar.ts"; @@ -31,6 +32,7 @@ import type { SidebarWorkboardBoard } from "./app-sidebar-workboard.ts"; import { icons } from "./icons.ts"; import { renderSessionGlyph, renderSessionUnreadBadge } from "./session-glyph.ts"; import { renderSessionRowBadges } from "./session-row-badges.ts"; +import { formatSidebarBuildSubtitle } from "./sidebar-build-chip-format.ts"; type AppSidebarRenderHost = AppSidebarSessionNavigationElement & { activePluginTabId: string; @@ -243,11 +245,17 @@ export function renderAppSidebarFooterBar(host: AppSidebarRenderHost) { watchedSessions: [], }; const gateway = host.offline ? null : readSidebarNativeGateway(); + const buildSubtitle = formatSidebarBuildSubtitle(CONTROL_UI_BUILD_INFO); // Health is visual-only here by budget decision; the header picker owns health accessibility. const gatewayPrimaryTag = gateway?.isPrimary ? t("chat.sessionHeader.gatewayPicker.primaryTag") : null; const identityMenuLabel = t("profilePage.identity.menuButtonLabel", { name: selfLabel }); + const identityDetail = host.offline + ? t("connection.reconnecting") + : gateway + ? `${gateway.name}${gatewayPrimaryTag ? `, ${gatewayPrimaryTag}` : ""}` + : buildSubtitle; return html`