diff --git a/CHANGELOG.md b/CHANGELOG.md index 1dc089c99c8a..7f6994542ba0 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -16,6 +16,7 @@ Docs: https://docs.openclaw.ai ### Fixes +- CLI/context engines: bootstrap and finalize non-legacy context engines for CLI turns while preserving transcript snapshots and deferred maintenance ownership. (#81869) Thanks @sahilsatralkar. - Telegram: persist polling updates through restart replay so queued same-topic messages resume in order instead of losing context after a gateway restart. (#82256) Thanks @VACInc. - Gateway/Gmail: abort in-flight Gmail watcher startup and hot-reload restarts before shutdown so reloads cannot spawn `gog serve` after the Gateway is closing. Thanks @frankekn. - MCP plugin tools: forward host MCP `tools/call` `AbortSignal` through `createPluginToolsMcpHandlers().callTool` into plugin `tool.execute`, so host cancellation actually cancels in-flight plugin tool calls instead of letting them run to completion. Fixes #82424. (#82443) Thanks @joshavant. diff --git a/src/agents/cli-runner.context-engine.test.ts b/src/agents/cli-runner.context-engine.test.ts new file mode 100644 index 000000000000..09e981f04e06 --- /dev/null +++ b/src/agents/cli-runner.context-engine.test.ts @@ -0,0 +1,443 @@ +import type { AgentMessage } from "@earendil-works/pi-agent-core"; +import { beforeEach, describe, expect, it, vi } from "vitest"; +import type { ContextEngine } from "../context-engine/types.js"; +import type { PreparedCliRunContext } from "./cli-runner/types.js"; + +const { + executePreparedCliRunMock, + loadCliSessionContextEngineMessagesMock, + loadCliSessionHistoryMessagesMock, + getGlobalHookRunnerMock, +} = vi.hoisted(() => ({ + executePreparedCliRunMock: vi.fn(), + loadCliSessionContextEngineMessagesMock: vi.fn(), + loadCliSessionHistoryMessagesMock: vi.fn(), + getGlobalHookRunnerMock: vi.fn(() => null), +})); + +vi.mock("./cli-runner/execute.runtime.js", () => ({ + executePreparedCliRun: executePreparedCliRunMock, +})); + +vi.mock("./cli-runner/session-history.js", () => ({ + loadCliSessionContextEngineMessages: loadCliSessionContextEngineMessagesMock, + loadCliSessionHistoryMessages: loadCliSessionHistoryMessagesMock, +})); + +vi.mock("../plugins/hook-runner-global.js", () => ({ + getGlobalHookRunner: getGlobalHookRunnerMock, +})); + +function textMessage(role: "user" | "assistant", text: string, timestamp: number): AgentMessage { + return { + role, + content: [{ type: "text", text }], + timestamp, + } as AgentMessage; +} + +function createContextEngine(overrides: Partial = {}): ContextEngine { + return { + info: { id: "test-context-engine", name: "Test context engine" }, + ingest: vi.fn(async () => ({ ingested: true })), + assemble: vi.fn(async (params) => ({ + messages: params.messages, + estimatedTokens: 0, + })), + compact: vi.fn(async () => ({ ok: true, compacted: false })), + ...overrides, + }; +} + +function createMaintenanceResult() { + return { + changed: false, + bytesFreed: 0, + rewrittenEntries: 0, + }; +} + +function buildPreparedContext(contextEngine: ContextEngine): PreparedCliRunContext { + const backend = { + command: "claude", + args: ["--print"], + output: "text" as const, + input: "arg" as const, + sessionMode: "existing" as const, + serialize: true, + }; + + return { + params: { + sessionId: "openclaw-session-1", + sessionKey: "agent:main:main", + agentId: "main", + sessionFile: "session.jsonl", + workspaceDir: "/tmp/openclaw-cli-context-engine-test", + prompt: "visible ask", + transcriptPrompt: "transcript visible ask", + provider: "claude-cli", + model: "sonnet-4.6", + thinkLevel: "low", + timeoutMs: 1_000, + runId: "run-1", + }, + started: Date.now(), + workspaceDir: "/tmp/openclaw-cli-context-engine-test", + backendResolved: { + id: "claude-cli", + config: backend, + bundleMcp: false, + pluginId: "anthropic", + }, + preparedBackend: { + backend, + env: {}, + }, + reusableCliSession: { + sessionId: "existing-external-cli-session", + }, + hadSessionFile: true, + contextEngineConfig: {}, + contextEngine, + contextEngineTurnPrompt: "transcript visible ask", + modelId: "sonnet-4.6", + normalizedModel: "sonnet-4.6", + systemPrompt: "You are a helpful assistant.", + systemPromptReport: {} as PreparedCliRunContext["systemPromptReport"], + bootstrapPromptWarningLines: [], + authEpochVersion: 2, + }; +} + +function expectMessageText(message: AgentMessage | undefined, expected: string): void { + expect(message).toBeDefined(); + const content = (message as { content?: unknown } | undefined)?.content; + if (typeof content === "string") { + expect(content).toBe(expected); + return; + } + expect(Array.isArray(content)).toBe(true); + expect((content as unknown[] | undefined)?.[0]).toMatchObject({ type: "text", text: expected }); +} + +describe("runPreparedCliAgent context engine lifecycle", () => { + beforeEach(() => { + executePreparedCliRunMock.mockReset(); + executePreparedCliRunMock.mockResolvedValue({ + text: " final answer ", + rawText: " final answer ", + sessionId: "external-cli-session-1", + usage: { input: 11, output: 7, total: 18 }, + finalPromptText: "prompt sent to cli", + }); + loadCliSessionContextEngineMessagesMock.mockReset(); + loadCliSessionContextEngineMessagesMock.mockResolvedValue([ + textMessage("user", "old ask", 1), + textMessage("assistant", "old answer", 2), + ]); + loadCliSessionHistoryMessagesMock.mockReset(); + loadCliSessionHistoryMessagesMock.mockResolvedValue([]); + getGlobalHookRunnerMock.mockReset(); + getGlobalHookRunnerMock.mockReturnValue(null); + }); + + it("finalizes successful CLI turns with the active context engine", async () => { + const bootstrap = vi.fn>(async () => ({ + bootstrapped: true, + })); + const afterTurn = vi.fn>(async () => {}); + const maintain = vi.fn>(async () => + createMaintenanceResult(), + ); + const dispose = vi.fn(async () => {}); + const contextEngine = createContextEngine({ bootstrap, afterTurn, maintain, dispose }); + const context = buildPreparedContext(contextEngine); + const { runPreparedCliAgent } = await import("./cli-runner.js"); + + const result = await runPreparedCliAgent(context); + + expect(result.meta.agentMeta?.sessionId).toBe("external-cli-session-1"); + expect(loadCliSessionContextEngineMessagesMock).toHaveBeenCalledWith({ + sessionId: "openclaw-session-1", + sessionFile: "session.jsonl", + sessionKey: "agent:main:main", + agentId: "main", + config: undefined, + }); + expect(loadCliSessionHistoryMessagesMock).not.toHaveBeenCalled(); + expect(bootstrap).toHaveBeenCalledWith({ + sessionId: "openclaw-session-1", + sessionKey: "agent:main:main", + sessionFile: "session.jsonl", + }); + expect(afterTurn).toHaveBeenCalledTimes(1); + const afterTurnParams = afterTurn.mock.calls[0]?.[0]; + expect(afterTurnParams).toMatchObject({ + sessionId: "openclaw-session-1", + sessionKey: "agent:main:main", + sessionFile: "session.jsonl", + prePromptMessageCount: 2, + tokenBudget: undefined, + runtimeContext: undefined, + }); + expect(afterTurnParams?.messages).toHaveLength(4); + expect(afterTurnParams?.messages.slice(0, 2)).toEqual([ + textMessage("user", "old ask", 1), + textMessage("assistant", "old answer", 2), + ]); + expectMessageText(afterTurnParams?.messages[2], "transcript visible ask"); + expectMessageText(afterTurnParams?.messages[3], "final answer"); + expect(afterTurnParams?.messages[3]).toMatchObject({ + role: "assistant", + provider: "claude-cli", + model: "sonnet-4.6", + usage: { input: 11, output: 7, total: 18 }, + }); + expect(maintain).toHaveBeenCalledTimes(2); + expect(maintain.mock.calls[1]?.[0]).toMatchObject({ + sessionId: "openclaw-session-1", + sessionKey: "agent:main:main", + sessionFile: "session.jsonl", + runtimeContext: { + rewriteTranscriptEntries: expect.any(Function), + llm: { complete: expect.any(Function) }, + }, + }); + expect(dispose).not.toHaveBeenCalled(); + }); + + it("does not synthesize a context-engine user turn for empty transcript prompts", async () => { + const afterTurn = vi.fn>(async () => {}); + const dispose = vi.fn(async () => {}); + const contextEngine = createContextEngine({ afterTurn, dispose }); + const context = buildPreparedContext(contextEngine); + context.params.transcriptPrompt = ""; + context.contextEngineTurnPrompt = ""; + const { runPreparedCliAgent } = await import("./cli-runner.js"); + + await runPreparedCliAgent(context); + + const afterTurnParams = afterTurn.mock.calls[0]?.[0]; + expect(afterTurnParams?.messages).toHaveLength(3); + expect(afterTurnParams?.prePromptMessageCount).toBe(2); + expect(afterTurnParams?.messages.slice(0, 2)).toEqual([ + textMessage("user", "old ask", 1), + textMessage("assistant", "old answer", 2), + ]); + const turnMessages = afterTurnParams?.messages.slice(afterTurnParams.prePromptMessageCount); + expect(turnMessages).toHaveLength(1); + expectMessageText(turnMessages?.[0], "final answer"); + expect(dispose).not.toHaveBeenCalled(); + }); + + it("does not finalize prepared model prompt as transcript turn text", async () => { + const afterTurn = vi.fn>(async () => {}); + const dispose = vi.fn(async () => {}); + const contextEngine = createContextEngine({ afterTurn, dispose }); + const context = buildPreparedContext(contextEngine); + context.params.prompt = "runtime context\n\noriginal user ask"; + delete context.params.transcriptPrompt; + context.contextEngineTurnPrompt = "original user ask"; + const { runPreparedCliAgent } = await import("./cli-runner.js"); + + await runPreparedCliAgent(context); + + const afterTurnParams = afterTurn.mock.calls[0]?.[0]; + expect(afterTurnParams?.messages).toHaveLength(4); + expect(afterTurnParams?.prePromptMessageCount).toBe(2); + const turnMessages = afterTurnParams?.messages.slice(afterTurnParams.prePromptMessageCount); + expect(turnMessages).toHaveLength(2); + expectMessageText(turnMessages?.[0], "original user ask"); + expectMessageText(turnMessages?.[1], "final answer"); + expect(dispose).not.toHaveBeenCalled(); + }); + + it("loads unbounded context-engine history separately from hook history", async () => { + const afterTurn = vi.fn>(async () => {}); + const dispose = vi.fn(async () => {}); + const contextEngine = createContextEngine({ afterTurn, dispose }); + const context = buildPreparedContext(contextEngine); + const fullHistory = Array.from({ length: 101 }, (_, index) => + textMessage("user", `old ask ${index}`, index), + ); + loadCliSessionContextEngineMessagesMock.mockResolvedValueOnce(fullHistory); + const { runPreparedCliAgent } = await import("./cli-runner.js"); + + await runPreparedCliAgent(context); + + const afterTurnParams = afterTurn.mock.calls[0]?.[0]; + expect(loadCliSessionContextEngineMessagesMock).toHaveBeenCalledTimes(1); + expect(loadCliSessionHistoryMessagesMock).not.toHaveBeenCalled(); + expect(afterTurnParams?.prePromptMessageCount).toBe(101); + expect(afterTurnParams?.messages.slice(0, 101)).toEqual(fullHistory); + }); + + it("loads context-engine history after bootstrap lifecycle runs", async () => { + const postBootstrapHistory = [textMessage("user", "post-bootstrap history", 9)]; + const bootstrap = vi.fn>(async () => { + loadCliSessionContextEngineMessagesMock.mockResolvedValueOnce(postBootstrapHistory); + return { bootstrapped: true }; + }); + const afterTurn = vi.fn>(async () => {}); + const dispose = vi.fn(async () => {}); + const contextEngine = createContextEngine({ bootstrap, afterTurn, dispose }); + const context = buildPreparedContext(contextEngine); + const { runPreparedCliAgent } = await import("./cli-runner.js"); + + await runPreparedCliAgent(context); + + expect(bootstrap).toHaveBeenCalledTimes(1); + expect(loadCliSessionContextEngineMessagesMock).toHaveBeenCalledTimes(1); + const bootstrapOrder = bootstrap.mock.invocationCallOrder[0]; + const loadHistoryOrder = loadCliSessionContextEngineMessagesMock.mock.invocationCallOrder[0]; + if (typeof bootstrapOrder !== "number" || typeof loadHistoryOrder !== "number") { + throw new Error("Expected bootstrap and history load invocation order"); + } + expect(bootstrapOrder).toBeLessThan(loadHistoryOrder); + const afterTurnParams = afterTurn.mock.calls[0]?.[0]; + expect(afterTurnParams?.prePromptMessageCount).toBe(1); + expect(afterTurnParams?.messages[0]).toEqual(postBootstrapHistory[0]); + }); + + it("falls back to ingestBatch and still runs turn maintenance", async () => { + const ingestBatch = vi.fn>(async () => ({ + ingestedCount: 2, + })); + const maintain = vi.fn>(async () => + createMaintenanceResult(), + ); + const dispose = vi.fn(async () => {}); + const contextEngine = createContextEngine({ ingestBatch, maintain, dispose }); + const { runPreparedCliAgent } = await import("./cli-runner.js"); + + await runPreparedCliAgent(buildPreparedContext(contextEngine)); + + expect(ingestBatch).toHaveBeenCalledTimes(1); + const ingestBatchParams = ingestBatch.mock.calls[0]?.[0]; + expect(ingestBatchParams).toMatchObject({ + sessionId: "openclaw-session-1", + sessionKey: "agent:main:main", + }); + expect(ingestBatchParams?.messages).toHaveLength(2); + expectMessageText(ingestBatchParams?.messages[0], "transcript visible ask"); + expectMessageText(ingestBatchParams?.messages[1], "final answer"); + expect(maintain).toHaveBeenCalledTimes(2); + expect(dispose).not.toHaveBeenCalled(); + }); + + it("preserves deferred maintenance ownership for background engines", async () => { + const maintain = vi.fn>(async () => + createMaintenanceResult(), + ); + const dispose = vi.fn(async () => {}); + const contextEngine = createContextEngine({ + info: { + id: "test-background-context-engine", + name: "Test background context engine", + turnMaintenanceMode: "background", + }, + maintain, + dispose, + }); + const { runPreparedCliAgent } = await import("./cli-runner.js"); + const context = buildPreparedContext(contextEngine); + + await runPreparedCliAgent(context); + + expect(dispose).not.toHaveBeenCalled(); + expect(context.contextEngineDeferredTurnMaintenance).toBeDefined(); + await context.contextEngineDeferredTurnMaintenance; + expect(dispose).not.toHaveBeenCalled(); + }); + + it("does not dispose background engines when no deferred turn maintenance is queued", async () => { + const dispose = vi.fn(async () => {}); + const contextEngine = createContextEngine({ + info: { + id: "test-background-context-engine", + name: "Test background context engine", + turnMaintenanceMode: "background", + }, + dispose, + }); + const { runPreparedCliAgent } = await import("./cli-runner.js"); + + await runPreparedCliAgent(buildPreparedContext(contextEngine)); + + expect(dispose).not.toHaveBeenCalled(); + }); + + it("does not dispose background engines after failed CLI attempts", async () => { + executePreparedCliRunMock.mockRejectedValue(new Error("cli boom")); + const maintain = vi.fn>(async () => + createMaintenanceResult(), + ); + const dispose = vi.fn(async () => {}); + const contextEngine = createContextEngine({ + info: { + id: "test-background-context-engine", + name: "Test background context engine", + turnMaintenanceMode: "background", + }, + maintain, + dispose, + }); + const { runPreparedCliAgent } = await import("./cli-runner.js"); + + await expect(runPreparedCliAgent(buildPreparedContext(contextEngine))).rejects.toThrow( + "cli boom", + ); + + expect(maintain).toHaveBeenCalledTimes(1); + expect(dispose).not.toHaveBeenCalled(); + }); + + it("does not finalize or run turn maintenance on failed CLI attempts", async () => { + executePreparedCliRunMock.mockRejectedValue(new Error("cli boom")); + const bootstrap = vi.fn>(async () => ({ + bootstrapped: true, + })); + const afterTurn = vi.fn>(async () => {}); + const ingestBatch = vi.fn>(async () => ({ + ingestedCount: 0, + })); + const maintain = vi.fn>(async () => + createMaintenanceResult(), + ); + const dispose = vi.fn(async () => {}); + const contextEngine = createContextEngine({ + bootstrap, + afterTurn, + ingestBatch, + maintain, + dispose, + }); + const { runPreparedCliAgent } = await import("./cli-runner.js"); + + await expect(runPreparedCliAgent(buildPreparedContext(contextEngine))).rejects.toThrow( + "cli boom", + ); + + expect(bootstrap).toHaveBeenCalledTimes(1); + expect(afterTurn).not.toHaveBeenCalled(); + expect(ingestBatch).not.toHaveBeenCalled(); + expect(maintain).toHaveBeenCalledTimes(1); + expect(dispose).not.toHaveBeenCalled(); + }); + + it("does not dispose context engines when CLI attempts fail", async () => { + executePreparedCliRunMock.mockRejectedValue(new Error("cli boom")); + const dispose = vi.fn(async () => { + throw new Error("dispose boom"); + }); + const contextEngine = createContextEngine({ dispose }); + const { runPreparedCliAgent } = await import("./cli-runner.js"); + + await expect(runPreparedCliAgent(buildPreparedContext(contextEngine))).rejects.toThrow( + "cli boom", + ); + + expect(dispose).not.toHaveBeenCalled(); + }); +}); diff --git a/src/agents/cli-runner.reliability.test.ts b/src/agents/cli-runner.reliability.test.ts index 3a104711e23f..582cb2d4d534 100644 --- a/src/agents/cli-runner.reliability.test.ts +++ b/src/agents/cli-runner.reliability.test.ts @@ -145,6 +145,8 @@ function buildPreparedContext(params?: { env: {}, }, reusableCliSession: params?.cliSessionId ? { sessionId: params.cliSessionId } : {}, + hadSessionFile: false, + contextEngineConfig: {}, modelId: "gpt-5.4", normalizedModel: "gpt-5.4", contextWindowInfo: { diff --git a/src/agents/cli-runner.spawn.test.ts b/src/agents/cli-runner.spawn.test.ts index 38c2068602ad..2c85d5c6d226 100644 --- a/src/agents/cli-runner.spawn.test.ts +++ b/src/agents/cli-runner.spawn.test.ts @@ -130,6 +130,8 @@ function buildPreparedCliRunContext(params: { ...(params.mcpConfigHash ? { mcpConfigHash: params.mcpConfigHash } : {}), }, reusableCliSession: {}, + hadSessionFile: false, + contextEngineConfig: {}, modelId: params.model, normalizedModel: params.model, systemPrompt: "You are a helpful assistant.", @@ -272,6 +274,8 @@ describe("runCliAgent spawn path", () => { env: {}, }, reusableCliSession: {}, + hadSessionFile: false, + contextEngineConfig: {}, modelId: "sonnet", normalizedModel: "sonnet", systemPrompt: "You are a helpful assistant.", diff --git a/src/agents/cli-runner.ts b/src/agents/cli-runner.ts index 5b4842a7430a..742a37b832f2 100644 --- a/src/agents/cli-runner.ts +++ b/src/agents/cli-runner.ts @@ -1,3 +1,4 @@ +import type { AgentMessage } from "@earendil-works/pi-agent-core"; import { SessionManager } from "@earendil-works/pi-coding-agent"; import type { ReplyPayload } from "../auto-reply/reply-payload.js"; import { SILENT_REPLY_TOKEN } from "../auto-reply/tokens.js"; @@ -6,9 +7,17 @@ import { createSubsystemLogger } from "../logging/subsystem.js"; import { buildAgentHookContextChannelFields } from "../plugins/hook-agent-context.js"; import { resolveBlockMessage } from "../plugins/hook-decision-types.js"; import { getGlobalHookRunner } from "../plugins/hook-runner-global.js"; -import { loadCliSessionHistoryMessages } from "./cli-runner/session-history.js"; +import { + loadCliSessionContextEngineMessages, + loadCliSessionHistoryMessages, +} from "./cli-runner/session-history.js"; import type { PreparedCliRunContext, RunCliAgentParams } from "./cli-runner/types.js"; import { FailoverError, isFailoverError, resolveFailoverStatus } from "./failover-error.js"; +import { + bootstrapHarnessContextEngine, + finalizeHarnessContextEngineTurn, + runHarnessContextEngineMaintenance, +} from "./harness/context-engine-lifecycle.js"; import { buildAgentHookContext } from "./harness/hook-context.js"; import { buildAgentHookConversationMessages } from "./harness/hook-history.js"; import { @@ -72,6 +81,89 @@ function buildCliHookAssistantMessage(params: { }; } +function isAgentMessage(value: unknown): value is AgentMessage { + return Boolean(value && typeof value === "object" && "role" in value); +} + +function buildCliContextEngineUserMessage(prompt: string): AgentMessage { + return { + role: "user", + content: prompt, + timestamp: Date.now(), + } as AgentMessage; +} + +function buildCliContextEngineAssistantMessage(params: { + text: string; + provider: string; + model: string; + usage?: { + input?: number; + output?: number; + cacheRead?: number; + cacheWrite?: number; + total?: number; + }; +}): AgentMessage { + return buildCliHookAssistantMessage(params) as AgentMessage; +} + +async function finalizeCliContextEngineTurn(params: { + context: PreparedCliRunContext; + historyMessages: unknown[]; + assistantText: string; + output: Awaited< + ReturnType + >; +}): Promise { + const { context } = params; + if (!context.contextEngine) { + return; + } + + const { params: runParams } = context; + const prePromptMessages = params.historyMessages.filter(isAgentMessage); + const turnMessages: AgentMessage[] = []; + if (context.contextEngineTurnPrompt) { + turnMessages.push(buildCliContextEngineUserMessage(context.contextEngineTurnPrompt)); + } + if (params.assistantText) { + turnMessages.push( + buildCliContextEngineAssistantMessage({ + text: params.assistantText, + provider: runParams.provider, + model: context.modelId, + usage: params.output.usage, + }), + ); + } + + let deferredTurnMaintenance: Promise | undefined; + const result = await finalizeHarnessContextEngineTurn({ + contextEngine: context.contextEngine, + promptError: false, + aborted: runParams.abortSignal?.aborted === true, + yieldAborted: false, + sessionIdUsed: runParams.sessionId, + sessionKey: runParams.sessionKey, + sessionFile: runParams.sessionFile, + messagesSnapshot: [...prePromptMessages, ...turnMessages], + prePromptMessageCount: prePromptMessages.length, + config: context.contextEngineConfig, + runMaintenance: async (maintenanceParams) => + await runHarnessContextEngineMaintenance({ + ...maintenanceParams, + onDeferredMaintenance: (promise) => { + deferredTurnMaintenance = promise; + }, + }), + warn: (message) => log.warn(message), + }); + if (result.postTurnFinalizationSucceeded && deferredTurnMaintenance) { + context.contextEngineDeferredTurnMaintenance = deferredTurnMaintenance; + } +} + export async function runCliAgent(params: RunCliAgentParams): Promise { // Cron gate must fire before prepareCliRunContext — that call allocates // backend resources released only by runPreparedCliAgent's try…finally. @@ -138,16 +230,16 @@ export async function runPreparedCliAgent( const hasLlmOutputHooks = hookRunner?.hasHooks("llm_output") === true; const hasAgentEndHooks = hookRunner?.hasHooks("agent_end") === true; const hasBeforeAgentRunHooks = hookRunner?.hasHooks("before_agent_run") === true; - const historyMessages = - hasLlmInputHooks || hasAgentEndHooks || hasBeforeAgentRunHooks - ? await loadCliSessionHistoryMessages({ - sessionId: params.sessionId, - sessionFile: params.sessionFile, - sessionKey: params.sessionKey, - agentId: params.agentId, - config: params.config, - }) - : []; + const needsHookHistory = hasLlmInputHooks || hasAgentEndHooks || hasBeforeAgentRunHooks; + const historyMessages = needsHookHistory + ? await loadCliSessionHistoryMessages({ + sessionId: params.sessionId, + sessionFile: params.sessionFile, + sessionKey: params.sessionKey, + agentId: params.agentId, + config: params.config, + }) + : []; const llmInputEvent = { runId: params.runId, sessionId: params.sessionId, @@ -416,6 +508,25 @@ export async function runPreparedCliAgent( // Try with the provided CLI session ID first try { + await bootstrapHarnessContextEngine({ + hadSessionFile: context.hadSessionFile, + contextEngine: context.contextEngine, + sessionId: params.sessionId, + sessionKey: params.sessionKey, + sessionFile: params.sessionFile, + config: context.contextEngineConfig, + warn: (message) => log.warn(message), + }); + const contextEngineHistoryMessages = context.contextEngine + ? await loadCliSessionContextEngineMessages({ + sessionId: params.sessionId, + sessionFile: params.sessionFile, + sessionKey: params.sessionKey, + agentId: params.agentId, + config: params.config, + }) + : []; + if (hasBeforeAgentRunHooks && hookRunner) { let beforeRunResult: | Awaited["runBeforeAgentRun"]>> @@ -479,7 +590,14 @@ export async function runPreparedCliAgent( const { output, lastAssistant } = await executeCliAttempt( context.reusableCliSession.sessionId, ); + const assistantText = output.text.trim(); const effectiveCliSessionId = output.sessionId ?? context.reusableCliSession.sessionId; + await finalizeCliContextEngineTurn({ + context, + historyMessages: context.contextEngine ? contextEngineHistoryMessages : historyMessages, + assistantText, + output, + }); runAgentHarnessAgentEndHook({ event: { messages: buildAgentEndMessages(lastAssistant), @@ -502,7 +620,16 @@ export async function runPreparedCliAgent( // For now, retry without the session ID to create a new session try { const { output, lastAssistant } = await executeCliAttempt(undefined); + const assistantText = output.text.trim(); const effectiveCliSessionId = output.sessionId; + await finalizeCliContextEngineTurn({ + context, + historyMessages: context.contextEngine + ? contextEngineHistoryMessages + : historyMessages, + assistantText, + output, + }); runAgentHarnessAgentEndHook({ event: { messages: buildAgentEndMessages(lastAssistant), diff --git a/src/agents/cli-runner/execute.supervisor-capture.test.ts b/src/agents/cli-runner/execute.supervisor-capture.test.ts index be5caae7de16..41b139a896e5 100644 --- a/src/agents/cli-runner/execute.supervisor-capture.test.ts +++ b/src/agents/cli-runner/execute.supervisor-capture.test.ts @@ -44,6 +44,8 @@ function buildPreparedCliRunContext(params: { env: {}, }, reusableCliSession: {}, + hadSessionFile: false, + contextEngineConfig: {}, modelId: "model", normalizedModel: "model", systemPrompt: "system", diff --git a/src/agents/cli-runner/prepare.test.ts b/src/agents/cli-runner/prepare.test.ts index 7807adb34c69..9f3449a51945 100644 --- a/src/agents/cli-runner/prepare.test.ts +++ b/src/agents/cli-runner/prepare.test.ts @@ -4,6 +4,12 @@ import path from "node:path"; import { CURRENT_SESSION_VERSION } from "@earendil-works/pi-coding-agent"; import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import type { OpenClawConfig } from "../../config/types.openclaw.js"; +import { registerLegacyContextEngine } from "../../context-engine/legacy.registration.js"; +import { + registerContextEngine, + registerContextEngineForOwner, +} from "../../context-engine/registry.js"; +import type { ContextEngine } from "../../context-engine/types.js"; import { getGlobalHookRunner } from "../../plugins/hook-runner-global.js"; import { __testing as cliBackendsTesting } from "../cli-backends.js"; import { hashCliSessionText } from "../cli-session.js"; @@ -15,6 +21,12 @@ import { shouldSkipLocalCliCredentialEpoch, } from "./prepare.js"; +const getRuntimeConfigMock = vi.hoisted(() => vi.fn(() => ({}))); + +vi.mock("../../config/config.js", () => ({ + getRuntimeConfig: getRuntimeConfigMock, +})); + vi.mock("../../plugins/hook-runner-global.js", () => ({ getGlobalHookRunner: vi.fn(() => null), })); @@ -174,12 +186,14 @@ describe("shouldSkipLocalCliCredentialEpoch", () => { resolveOpenClawReferencePaths: vi.fn(async () => ({ docsPath: null, sourcePath: null })), }); mockGetGlobalHookRunner.mockReturnValue(null); + getRuntimeConfigMock.mockReturnValue({}); mockBuildActiveVideoGenerationTaskPromptContextForSession.mockReturnValue(undefined); mockBuildActiveMusicGenerationTaskPromptContextForSession.mockReturnValue(undefined); }); afterEach(() => { cliBackendsTesting.resetDepsForTest(); + getRuntimeConfigMock.mockReset(); mockGetGlobalHookRunner.mockReset(); mockBuildActiveVideoGenerationTaskPromptContextForSession.mockReset(); mockBuildActiveMusicGenerationTaskPromptContextForSession.mockReset(); @@ -281,6 +295,7 @@ describe("shouldSkipLocalCliCredentialEpoch", () => { }); expect(context.params.prompt).toBe("history:2\n\nlatest ask"); + expect(context.contextEngineTurnPrompt).toBe("latest ask"); expect(context.systemPrompt).toBe( "prepend system\n\nhook system\n\nappend system\n\nCurrent model identity: test-cli/test-model. If asked what model you are, answer with this value for the current run.", ); @@ -377,6 +392,7 @@ describe("shouldSkipLocalCliCredentialEpoch", () => { "Sender (untrusted metadata):\nsender_id=U123 trusted hook context\n\nlatest ask\n\ntrusted hook tail", ); expect(context.params.transcriptPrompt).toBe("latest ask"); + expect(context.contextEngineTurnPrompt).toBe("latest ask"); expect(hookRunner.runBeforePromptBuild).toHaveBeenCalledTimes(1); const beforePromptBuildCalls = hookRunner.runBeforePromptBuild.mock.calls as unknown as Array< [unknown, unknown] @@ -560,6 +576,174 @@ describe("shouldSkipLocalCliCredentialEpoch", () => { } }); + it("does not allocate a non-legacy context engine before fallible CLI preparation finishes", async () => { + const { dir, sessionFile } = createSessionFile(); + const engineId = `cli-prepare-late-engine-${Date.now().toString(36)}`; + const dispose = vi.fn(async () => {}); + const factory = vi.fn((): ContextEngine => { + return { + info: { id: engineId, name: "CLI prepare late engine" }, + ingest: vi.fn(async () => ({ ingested: true })), + assemble: vi.fn(async ({ messages }) => ({ messages, estimatedTokens: 0 })), + compact: vi.fn(async () => ({ ok: true, compacted: false })), + dispose, + }; + }); + registerContextEngine(engineId, factory); + setCliRunnerPrepareTestDeps({ + resolveOpenClawReferencePaths: vi.fn(async () => { + throw new Error("reference path lookup failed"); + }), + }); + + try { + await expect( + prepareCliRunContext({ + sessionId: "session-test", + sessionFile, + workspaceDir: dir, + prompt: "latest ask", + provider: "test-cli", + model: "test-model", + timeoutMs: 1_000, + runId: "run-test-prepare-failure", + config: { + ...createCliBackendConfig(), + plugins: { slots: { contextEngine: engineId } }, + }, + }), + ).rejects.toThrow("reference path lookup failed"); + + expect(factory).not.toHaveBeenCalled(); + expect(dispose).not.toHaveBeenCalled(); + } finally { + fs.rmSync(dir, { recursive: true, force: true }); + } + }); + + it("cleans up prepared CLI backend when context-engine resolution fails", async () => { + const { dir, sessionFile } = createSessionFile(); + const cleanup = vi.fn(async () => {}); + const prepareExecution = vi.fn(async () => ({ cleanup })); + registerContextEngineForOwner( + "legacy", + () => { + throw new Error("context engine failed"); + }, + "core", + { allowSameOwnerRefresh: true }, + ); + cliBackendsTesting.setDepsForTest({ + resolvePluginSetupCliBackend: () => undefined, + resolveRuntimeCliBackends: () => [ + { + id: "test-cli", + pluginId: "test-plugin", + bundleMcp: false, + prepareExecution, + config: { + command: "test-cli", + args: ["--print"], + systemPromptArg: "--system-prompt", + systemPromptWhen: "first", + sessionMode: "existing", + output: "text", + input: "arg", + }, + }, + ], + }); + + try { + await expect( + prepareCliRunContext({ + sessionId: "session-test", + sessionFile, + workspaceDir: dir, + prompt: "latest ask", + provider: "test-cli", + model: "test-model", + timeoutMs: 1_000, + runId: "run-test-context-engine-resolution-failure", + config: createCliBackendConfig(), + }), + ).rejects.toThrow("context engine failed"); + + expect(prepareExecution).toHaveBeenCalledOnce(); + expect(cleanup).toHaveBeenCalledOnce(); + } finally { + registerLegacyContextEngine(); + fs.rmSync(dir, { recursive: true, force: true }); + } + }); + + it("uses runtime config when resolving the CLI context engine", async () => { + const { dir, sessionFile } = createSessionFile(); + const engineId = `cli-runtime-config-engine-${Date.now().toString(36)}`; + const runtimeAgentDir = path.join(dir, "runtime-agent"); + const runtimeConfig = { + agents: { + list: [{ id: "main", default: true, agentDir: runtimeAgentDir }], + }, + plugins: { slots: { contextEngine: engineId } }, + } satisfies OpenClawConfig; + const factory = vi.fn((_ctx: unknown): ContextEngine => { + return { + info: { id: engineId, name: "CLI runtime config engine" }, + ingest: vi.fn(async () => ({ ingested: true })), + assemble: vi.fn(async ({ messages }) => ({ messages, estimatedTokens: 0 })), + compact: vi.fn(async () => ({ ok: true, compacted: false })), + }; + }); + registerContextEngine(engineId, factory); + getRuntimeConfigMock.mockReturnValue(runtimeConfig); + cliBackendsTesting.setDepsForTest({ + resolvePluginSetupCliBackend: () => undefined, + resolveRuntimeCliBackends: () => [ + { + id: "test-cli", + pluginId: "test-plugin", + bundleMcp: false, + config: { + command: "test-cli", + args: ["--print"], + systemPromptArg: "--system-prompt", + systemPromptWhen: "first", + sessionMode: "existing", + output: "text", + input: "arg", + }, + }, + ], + }); + + try { + const context = await prepareCliRunContext({ + sessionId: "session-test", + sessionFile, + workspaceDir: dir, + prompt: "latest ask", + provider: "test-cli", + model: "test-model", + timeoutMs: 1_000, + runId: "run-test-runtime-config-context-engine", + }); + + expect(context.contextEngine?.info.id).toBe(engineId); + expect(context.contextEngineConfig).toBe(runtimeConfig); + expect(context.params.config).toBe(runtimeConfig); + expect(factory).toHaveBeenCalledWith( + expect.objectContaining({ + agentDir: runtimeAgentDir, + config: runtimeConfig, + workspaceDir: dir, + }), + ); + } finally { + fs.rmSync(dir, { recursive: true, force: true }); + } + }); + it("uses explicit static prompt text for CLI session reuse hashing", async () => { const { dir, sessionFile } = createSessionFile(); try { diff --git a/src/agents/cli-runner/prepare.ts b/src/agents/cli-runner/prepare.ts index 60d33fc67397..b33481256e6a 100644 --- a/src/agents/cli-runner/prepare.ts +++ b/src/agents/cli-runner/prepare.ts @@ -1,4 +1,6 @@ import { getRuntimeConfig } from "../../config/config.js"; +import { ensureContextEnginesInitialized } from "../../context-engine/init.js"; +import { resolveContextEngine } from "../../context-engine/registry.js"; import { ensureMcpLoopbackServer } from "../../gateway/mcp-http.js"; import { createMcpLoopbackServerConfig, @@ -53,6 +55,7 @@ import { buildCliAgentSystemPrompt, normalizeCliModel } from "./helpers.js"; import { cliBackendLog } from "./log.js"; import { buildCliSessionHistoryPrompt, + hasCliSessionTranscript, loadCliSessionHistoryMessages, loadCliSessionReseedMessages, } from "./session-history.js"; @@ -475,25 +478,65 @@ export async function prepareCliRunContext( runtimeContextChars: 0, }, }); + const contextEngineConfig = params.config ?? getRuntimeConfig(); + try { + ensureContextEnginesInitialized(); + const { sessionAgentId: contextEngineSessionAgentId } = resolveSessionAgentIds({ + sessionKey: params.sessionKey, + config: contextEngineConfig, + agentId: params.agentId, + }); + const contextEngineAgentDir = resolveAgentDir(contextEngineConfig, contextEngineSessionAgentId); + const resolvedContextEngine = await resolveContextEngine(contextEngineConfig, { + agentDir: contextEngineAgentDir, + workspaceDir, + }); + const contextEngine = + resolvedContextEngine.info.id !== "legacy" ? resolvedContextEngine : undefined; + const hadSessionFile = await hasCliSessionTranscript({ + sessionId: params.sessionId, + sessionFile: params.sessionFile, + sessionKey: params.sessionKey, + agentId: params.agentId, + config: contextEngineConfig, + }); + const contextEngineTurnPrompt = params.transcriptPrompt ?? params.prompt; + const preparedParams: RunCliAgentParams = { + ...params, + config: contextEngineConfig, + prompt: preparedPrompt, + }; - return { - params: preparedPrompt === params.prompt ? params : { ...params, prompt: preparedPrompt }, - effectiveAuthProfileId, - started, - workspaceDir, - backendResolved, - preparedBackend: preparedBackendFinal, - reusableCliSession, - modelId, - normalizedModel, - contextWindowInfo, - systemPrompt, - systemPromptReport, - bootstrapPromptWarningLines: bootstrapPromptWarning.lines, - ...(openClawHistoryPrompt ? { openClawHistoryPrompt } : {}), - heartbeatPrompt, - authEpoch, - authEpochVersion: CLI_AUTH_EPOCH_VERSION, - extraSystemPromptHash, - }; + return { + params: preparedParams, + effectiveAuthProfileId, + started, + workspaceDir, + backendResolved, + preparedBackend: preparedBackendFinal, + reusableCliSession, + hadSessionFile, + contextEngineConfig, + contextEngine, + contextEngineTurnPrompt, + modelId, + normalizedModel, + contextWindowInfo, + systemPrompt, + systemPromptReport, + bootstrapPromptWarningLines: bootstrapPromptWarning.lines, + ...(openClawHistoryPrompt ? { openClawHistoryPrompt } : {}), + heartbeatPrompt, + authEpoch, + authEpochVersion: CLI_AUTH_EPOCH_VERSION, + extraSystemPromptHash, + }; + } catch (err) { + try { + await preparedBackendFinal.cleanup?.(); + } catch (cleanupErr) { + cliBackendLog.warn(`cli backend cleanup after prepare failure failed: ${String(cleanupErr)}`); + } + throw err; + } } diff --git a/src/agents/cli-runner/session-history.test.ts b/src/agents/cli-runner/session-history.test.ts index 14f324915b80..35447c024654 100644 --- a/src/agents/cli-runner/session-history.test.ts +++ b/src/agents/cli-runner/session-history.test.ts @@ -5,6 +5,8 @@ import { CURRENT_SESSION_VERSION } from "@earendil-works/pi-coding-agent"; import { afterEach, describe, expect, it, vi } from "vitest"; import { buildCliSessionHistoryPrompt, + hasCliSessionTranscript, + loadCliSessionContextEngineMessages, loadCliSessionHistoryMessages, loadCliSessionReseedMessages, MAX_CLI_SESSION_HISTORY_FILE_BYTES, @@ -80,6 +82,19 @@ function expectCompactionSummary(value: unknown, summary: string) { expect(message.summary).toBe(summary); } +function expectCustomMessage(value: unknown, expected: { customType: string; content: string }) { + const message = requireRecord(value, "custom message"); + expect(message.role).toBe("custom"); + expect(message.customType).toBe(expected.customType); + expect(message.content).toBe(expected.content); +} + +function expectBranchSummary(value: unknown, summary: string) { + const message = requireRecord(value, "branch summary"); + expect(message.role).toBe("branchSummary"); + expect(message.summary).toBe(summary); +} + describe("loadCliSessionHistoryMessages", () => { afterEach(() => { vi.unstubAllEnvs(); @@ -116,6 +131,37 @@ describe("loadCliSessionHistoryMessages", () => { } }); + it("detects canonical transcripts when callers pass stale external session paths", async () => { + const stateDir = fs.mkdtempSync(path.join(os.tmpdir(), "openclaw-cli-state-")); + const outsideDir = fs.mkdtempSync(path.join(os.tmpdir(), "openclaw-cli-outside-")); + vi.stubEnv("OPENCLAW_STATE_DIR", stateDir); + createSessionTranscript({ + rootDir: stateDir, + sessionId: "session-test", + messages: ["expected history"], + }); + const outsideFile = createSessionTranscript({ + rootDir: outsideDir, + sessionId: "session-test", + filePath: path.join(outsideDir, "stale.jsonl"), + messages: ["stale history"], + }); + + try { + await expect( + hasCliSessionTranscript({ + sessionId: "session-test", + sessionFile: outsideFile, + sessionKey: "agent:main:main", + agentId: "main", + }), + ).resolves.toBe(true); + } finally { + fs.rmSync(stateDir, { recursive: true, force: true }); + fs.rmSync(outsideDir, { recursive: true, force: true }); + } + }); + it("keeps only the newest bounded history window", async () => { const stateDir = fs.mkdtempSync(path.join(os.tmpdir(), "openclaw-cli-state-")); vi.stubEnv("OPENCLAW_STATE_DIR", stateDir); @@ -146,6 +192,115 @@ describe("loadCliSessionHistoryMessages", () => { } }); + it("keeps complete history for context-engine snapshots", async () => { + const stateDir = fs.mkdtempSync(path.join(os.tmpdir(), "openclaw-cli-state-")); + vi.stubEnv("OPENCLAW_STATE_DIR", stateDir); + const sessionFile = createSessionTranscript({ + rootDir: stateDir, + sessionId: "session-context-engine-history", + messages: Array.from( + { length: MAX_CLI_SESSION_HISTORY_MESSAGES + 25 }, + (_, index) => `msg-${index}`, + ), + }); + + try { + const history = await loadCliSessionContextEngineMessages({ + sessionId: "session-context-engine-history", + sessionFile, + sessionKey: "agent:main:main", + agentId: "main", + }); + expect(history).toHaveLength(MAX_CLI_SESSION_HISTORY_MESSAGES + 25); + expectMessageFields(history[0], { role: "user", content: "msg-0" }); + expectMessageFields(history.at(-1), { + role: "user", + content: `msg-${MAX_CLI_SESSION_HISTORY_MESSAGES + 24}`, + }); + } finally { + fs.rmSync(stateDir, { recursive: true, force: true }); + } + }); + + it("uses the latest compaction summary and complete tail for context-engine snapshots", async () => { + const stateDir = fs.mkdtempSync(path.join(os.tmpdir(), "openclaw-cli-state-")); + vi.stubEnv("OPENCLAW_STATE_DIR", stateDir); + const sessionFile = createSessionTranscript({ + rootDir: stateDir, + sessionId: "session-context-engine-compacted", + messages: ["old ask"], + }); + fs.appendFileSync( + sessionFile, + `${JSON.stringify({ + type: "compaction", + id: "compact-1", + timestamp: new Date(2).toISOString(), + summary: "Earlier compacted context", + })}\n`, + "utf-8", + ); + fs.appendFileSync( + sessionFile, + `${JSON.stringify({ + type: "custom_message", + id: "custom-tail", + parentId: "compaction-1", + timestamp: new Date(3).toISOString(), + customType: "runtime-note", + content: "tail custom context", + display: false, + })}\n`, + "utf-8", + ); + fs.appendFileSync( + sessionFile, + `${JSON.stringify({ + type: "branch_summary", + id: "branch-tail", + parentId: "custom-tail", + fromId: "custom-tail", + timestamp: new Date(4).toISOString(), + summary: "tail branch context", + })}\n`, + "utf-8", + ); + fs.appendFileSync( + sessionFile, + `${JSON.stringify({ + type: "message", + id: "msg-tail", + parentId: "branch-tail", + timestamp: new Date(5).toISOString(), + message: { + role: "assistant", + content: "tail answer", + timestamp: 5, + }, + })}\n`, + "utf-8", + ); + + try { + const history = await loadCliSessionContextEngineMessages({ + sessionId: "session-context-engine-compacted", + sessionFile, + sessionKey: "agent:main:main", + agentId: "main", + }); + expect(history).toHaveLength(4); + expectCompactionSummary(history[0], "Earlier compacted context"); + expectCustomMessage(history[1], { + customType: "runtime-note", + content: "tail custom context", + }); + expectBranchSummary(history[2], "tail branch context"); + expectMessageFields(history[3], { role: "assistant", content: "tail answer" }); + } finally { + fs.rmSync(stateDir, { recursive: true, force: true }); + } + }); + it("rejects symlinked transcripts instead of following them outside the sessions directory", async () => { const stateDir = fs.mkdtempSync(path.join(os.tmpdir(), "openclaw-cli-state-")); const outsideDir = fs.mkdtempSync(path.join(os.tmpdir(), "openclaw-cli-outside-")); diff --git a/src/agents/cli-runner/session-history.ts b/src/agents/cli-runner/session-history.ts index dddf4578ed58..baced8b447f2 100644 --- a/src/agents/cli-runner/session-history.ts +++ b/src/agents/cli-runner/session-history.ts @@ -1,5 +1,6 @@ import fsp from "node:fs/promises"; import path from "node:path"; +import type { AgentMessage } from "@earendil-works/pi-agent-core"; import { migrateSessionEntries, parseSessionEntries } from "@earendil-works/pi-coding-agent"; import { resolveSessionFilePath, @@ -26,6 +27,15 @@ type HistoryEntry = { type?: unknown; message?: unknown; summary?: unknown; + customType?: unknown; + content?: unknown; + display?: unknown; + details?: unknown; + timestamp?: unknown; + fromId?: unknown; + firstKeptEntryId?: unknown; + tokensBefore?: unknown; + tokensAfter?: unknown; }; type RawTranscriptReseedReason = @@ -62,6 +72,48 @@ function coerceHistoryText(content: unknown): string { .trim(); } +function coerceHistoryTimestamp(value: unknown): number | string { + if (typeof value === "number" && Number.isFinite(value)) { + return value; + } + if (typeof value === "string") { + return value; + } + return 0; +} + +function historyEntryToContextEngineMessage(entry: HistoryEntry): AgentMessage | undefined { + if (entry.type === "message") { + return entry.message as AgentMessage; + } + if (entry.type === "custom_message") { + return { + role: "custom", + customType: typeof entry.customType === "string" ? entry.customType : "custom", + content: entry.content, + display: entry.display !== false, + details: entry.details, + timestamp: coerceHistoryTimestamp(entry.timestamp), + } as AgentMessage; + } + if (entry.type === "branch_summary") { + return { + role: "branchSummary", + summary: typeof entry.summary === "string" ? entry.summary : "", + fromId: typeof entry.fromId === "string" ? entry.fromId : "root", + timestamp: coerceHistoryTimestamp(entry.timestamp), + } as AgentMessage; + } + return undefined; +} + +function loadContextEngineMessagesFromEntries(entries: unknown[]): AgentMessage[] { + return entries.flatMap((entry) => { + const message = historyEntryToContextEngineMessage(entry as HistoryEntry); + return message ? [message] : []; + }); +} + export function buildCliSessionHistoryPrompt(params: { messages: unknown[]; prompt: string; @@ -185,6 +237,35 @@ async function loadCliSessionEntries(params: { } } +export async function hasCliSessionTranscript(params: { + sessionId: string; + sessionFile: string; + sessionKey?: string; + agentId?: string; + config?: OpenClawConfig; +}): Promise { + try { + const { sessionFile, sessionsDir } = resolveSafeCliSessionFile(params); + const entryStat = await fsp.lstat(sessionFile); + if (!entryStat.isFile() || entryStat.isSymbolicLink()) { + return false; + } + const realSessionsDir = (await safeRealpath(sessionsDir)) ?? path.resolve(sessionsDir); + const realSessionFile = await safeRealpath(sessionFile); + if ( + !realSessionFile || + realSessionFile === realSessionsDir || + !isPathInside(realSessionsDir, realSessionFile) + ) { + return false; + } + const stat = await fsp.stat(realSessionFile); + return stat.isFile() && stat.size <= MAX_CLI_SESSION_HISTORY_FILE_BYTES; + } catch { + return false; + } +} + export async function loadCliSessionHistoryMessages(params: { sessionId: string; sessionFile: string; @@ -199,6 +280,49 @@ export async function loadCliSessionHistoryMessages(params: { return limitAgentHookHistoryMessages(history, MAX_CLI_SESSION_HISTORY_MESSAGES); } +export async function loadCliSessionContextEngineMessages(params: { + sessionId: string; + sessionFile: string; + sessionKey?: string; + agentId?: string; + config?: OpenClawConfig; +}): Promise { + const entries = await loadCliSessionEntries(params); + const latestCompactionIndex = entries.findLastIndex((entry) => { + const candidate = entry as HistoryEntry; + return candidate.type === "compaction" && typeof candidate.summary === "string"; + }); + if (latestCompactionIndex < 0) { + return loadContextEngineMessagesFromEntries(entries); + } + + const compaction = entries[latestCompactionIndex] as HistoryEntry; + const summary = typeof compaction.summary === "string" ? compaction.summary.trim() : ""; + if (!summary) { + return loadContextEngineMessagesFromEntries(entries); + } + + const tailMessages = loadContextEngineMessagesFromEntries( + entries.slice(latestCompactionIndex + 1), + ); + return [ + { + role: "compactionSummary", + summary, + timestamp: coerceHistoryTimestamp(compaction.timestamp), + tokensBefore: typeof compaction.tokensBefore === "number" ? compaction.tokensBefore : 0, + ...(typeof compaction.tokensAfter === "number" + ? { tokensAfter: compaction.tokensAfter } + : {}), + ...(typeof compaction.firstKeptEntryId === "string" + ? { firstKeptEntryId: compaction.firstKeptEntryId } + : {}), + ...(compaction.details !== undefined ? { details: compaction.details } : {}), + }, + ...tailMessages, + ]; +} + export async function loadCliSessionReseedMessages(params: { sessionId: string; sessionFile: string; diff --git a/src/agents/cli-runner/types.ts b/src/agents/cli-runner/types.ts index 5c42f234f360..97c294b6f30b 100644 --- a/src/agents/cli-runner/types.ts +++ b/src/agents/cli-runner/types.ts @@ -7,6 +7,7 @@ import type { CliSessionBinding } from "../../config/sessions.js"; import type { SessionSystemPromptReport } from "../../config/sessions/types.js"; import type { CliBackendConfig } from "../../config/types.js"; import type { OpenClawConfig } from "../../config/types.openclaw.js"; +import type { ContextEngine } from "../../context-engine/types.js"; import type { PromptImageOrderEntry } from "../../media/prompt-image-order.js"; import type { InputProvenance } from "../../sessions/input-provenance.js"; import type { BootstrapContextMode } from "../bootstrap-files.js"; @@ -114,6 +115,11 @@ export type PreparedCliRunContext = { backendResolved: ResolvedCliBackend; preparedBackend: CliPreparedBackend; reusableCliSession: CliReusableSession; + hadSessionFile: boolean; + contextEngineConfig: OpenClawConfig; + contextEngine?: ContextEngine; + contextEngineTurnPrompt?: string; + contextEngineDeferredTurnMaintenance?: Promise; modelId: string; normalizedModel: string; contextWindowInfo?: ContextWindowInfo; diff --git a/src/agents/harness/context-engine-lifecycle.ts b/src/agents/harness/context-engine-lifecycle.ts index f9c66e47d76b..53abe8bc3d48 100644 --- a/src/agents/harness/context-engine-lifecycle.ts +++ b/src/agents/harness/context-engine-lifecycle.ts @@ -230,6 +230,7 @@ export async function runHarnessContextEngineMaintenance(params: { sessionManager?: unknown; runtimeContext?: ContextEngineRuntimeContext; executionMode?: "foreground" | "background"; + onDeferredMaintenance?: (promise: Promise) => void; config?: SessionWriteLockAcquireTimeoutConfig; }) { return await runContextEngineMaintenance({ @@ -243,6 +244,7 @@ export async function runHarnessContextEngineMaintenance(params: { >[0]["sessionManager"], runtimeContext: params.runtimeContext, executionMode: params.executionMode, + onDeferredMaintenance: params.onDeferredMaintenance, config: params.config, }); } diff --git a/src/agents/pi-embedded-runner/context-engine-maintenance.test.ts b/src/agents/pi-embedded-runner/context-engine-maintenance.test.ts index d41c79e5b8fa..62d848ed0bb4 100644 --- a/src/agents/pi-embedded-runner/context-engine-maintenance.test.ts +++ b/src/agents/pi-embedded-runner/context-engine-maintenance.test.ts @@ -645,6 +645,7 @@ describe("runContextEngineMaintenance", () => { const sessionKey = "agent:main:session-rerun"; let releaseFirstMaintenance: (() => void) | undefined; + let releaseSecondMaintenance: (() => void) | undefined; let maintenanceCalls = 0; const maintain = vi.fn(async () => { maintenanceCalls += 1; @@ -653,6 +654,11 @@ describe("runContextEngineMaintenance", () => { releaseFirstMaintenance = resolve; }); } + if (maintenanceCalls === 2) { + await new Promise((resolve) => { + releaseSecondMaintenance = resolve; + }); + } return { changed: false, bytesFreed: 0, @@ -674,6 +680,7 @@ describe("runContextEngineMaintenance", () => { compact: async () => ({ ok: true, compacted: false }), maintain, } as NonNullable[0]["contextEngine"]>; + const deferredPromises: Promise[] = []; await runContextEngineMaintenance({ contextEngine: backgroundEngine, @@ -681,6 +688,9 @@ describe("runContextEngineMaintenance", () => { sessionKey, sessionFile: "/tmp/session-rerun.jsonl", reason: "turn", + onDeferredMaintenance: (promise) => { + deferredPromises.push(promise); + }, }); await waitForAssertion(() => expect(maintain).toHaveBeenCalledTimes(1)); @@ -691,6 +701,14 @@ describe("runContextEngineMaintenance", () => { sessionKey, sessionFile: "/tmp/session-rerun.jsonl", reason: "turn", + onDeferredMaintenance: (promise) => { + deferredPromises.push(promise); + }, + }); + expect(deferredPromises).toHaveLength(2); + let secondDeferredSettled = false; + const secondDeferred = deferredPromises[1].then(() => { + secondDeferredSettled = true; }); if (!releaseFirstMaintenance) { @@ -698,6 +716,15 @@ describe("runContextEngineMaintenance", () => { } releaseFirstMaintenance(); await waitForAssertion(() => expect(maintain).toHaveBeenCalledTimes(2)); + await Promise.resolve(); + expect(secondDeferredSettled).toBe(false); + + if (!releaseSecondMaintenance) { + throw new Error("Expected second maintenance release callback to be initialized"); + } + releaseSecondMaintenance(); + await secondDeferred; + expect(secondDeferredSettled).toBe(true); const tasks = listTasksForOwnerKey(sessionKey).filter( (task) => task.taskKind === TURN_MAINTENANCE_TASK_KIND, diff --git a/src/agents/pi-embedded-runner/context-engine-maintenance.ts b/src/agents/pi-embedded-runner/context-engine-maintenance.ts index 274c3ca2d130..72170f345db6 100644 --- a/src/agents/pi-embedded-runner/context-engine-maintenance.ts +++ b/src/agents/pi-embedded-runner/context-engine-maintenance.ts @@ -525,16 +525,18 @@ async function runDeferredTurnMaintenanceWorker(params: { } } -function scheduleDeferredTurnMaintenance(params: DeferredTurnMaintenanceScheduleParams): void { +function scheduleDeferredTurnMaintenance( + params: DeferredTurnMaintenanceScheduleParams, +): Promise | undefined { const sessionKey = normalizeSessionKey(params.sessionKey); if (!sessionKey) { - return; + return undefined; } const activeRun = activeDeferredTurnMaintenanceRuns.get(sessionKey); if (activeRun) { activeRun.rerunRequested = true; activeRun.latestParams = { ...params, sessionKey }; - return; + return activeRun.promise; } const existingTask = findActiveSessionTask({ @@ -589,7 +591,7 @@ function scheduleDeferredTurnMaintenance(params: DeferredTurnMaintenanceSchedule taskId: task.taskId, error: err, }); - return; + return undefined; } let state!: DeferredTurnMaintenanceRunState; const trackedPromise = runPromise @@ -600,7 +602,7 @@ function scheduleDeferredTurnMaintenance(params: DeferredTurnMaintenanceSchedule error: err, }); }) - .finally(() => { + .finally(async () => { schedulerAbort.dispose(); const current = activeDeferredTurnMaintenanceRuns.get(sessionKey); if (current !== state) { @@ -611,7 +613,7 @@ function scheduleDeferredTurnMaintenance(params: DeferredTurnMaintenanceSchedule current.rerunRequested && !shutdownTriggered ? current.latestParams : undefined; activeDeferredTurnMaintenanceRuns.delete(sessionKey); if (rerunParams) { - scheduleDeferredTurnMaintenance(rerunParams); + await scheduleDeferredTurnMaintenance(rerunParams); } }); state = { @@ -621,6 +623,7 @@ function scheduleDeferredTurnMaintenance(params: DeferredTurnMaintenanceSchedule }; activeDeferredTurnMaintenanceRuns.set(sessionKey, state); void trackedPromise; + return trackedPromise; } /** @@ -636,6 +639,7 @@ export async function runContextEngineMaintenance(params: { runtimeContext?: ContextEngineRuntimeContext; agentId?: string; executionMode?: "foreground" | "background"; + onDeferredMaintenance?: (promise: Promise) => void; config?: OpenClawConfig; }): Promise { if (typeof params.contextEngine?.maintain !== "function") { @@ -650,7 +654,7 @@ export async function runContextEngineMaintenance(params: { if (shouldDefer) { try { - scheduleDeferredTurnMaintenance({ + const deferred = scheduleDeferredTurnMaintenance({ contextEngine: params.contextEngine, sessionId: params.sessionId, sessionKey: params.sessionKey ?? params.sessionId, @@ -660,6 +664,9 @@ export async function runContextEngineMaintenance(params: { agentId: params.agentId, config: params.config, }); + if (deferred) { + params.onDeferredMaintenance?.(deferred); + } } catch (err) { log.warn(`failed to schedule deferred context engine maintenance: ${String(err)}`); }