diff --git a/config/assertion-safety-baseline.txt b/config/assertion-safety-baseline.txt index 56ca9ec62b14..30dc8c2036b6 100644 --- a/config/assertion-safety-baseline.txt +++ b/config/assertion-safety-baseline.txt @@ -1370,7 +1370,6 @@ src/agents/embedded-agent-runner/model.provider-hooks.ts 8 src/agents/embedded-agent-runner/model.registry-resolution.ts 3 src/agents/embedded-agent-runner/model.ts 5 src/agents/embedded-agent-runner/openrouter-model-capabilities.ts 2 -src/agents/embedded-agent-runner/prompt-cache-observability.ts 2 src/agents/embedded-agent-runner/replay-history.ts 20 src/agents/embedded-agent-runner/run-entry.ts 2 src/agents/embedded-agent-runner/run.overflow-compaction.harness.ts 4 diff --git a/docs/reference/prompt-caching.md b/docs/reference/prompt-caching.md index 0821d61f60f0..07bddffbe8b3 100644 --- a/docs/reference/prompt-caching.md +++ b/docs/reference/prompt-caching.md @@ -385,7 +385,9 @@ diagnostics: ### What to inspect -Prompt-cache observations record `input`, `cacheRead`, and `cacheWrite` per completed foreground model request alongside its stable system-prefix, volatile-suffix, and tools fingerprints, and flag cache-read drops from the previous request, including reported zero reads; billing totals remain separate. A flagged drop lists the tracked changes since the last request (`model`, `cacheRetention`, `transport`, `streamStrategy`, `systemPrompt`, `systemPromptSuffix`, `tools`, `aggregateToolResultTruncation`). Observations and warnings require cache tracing (`diagnostics.cacheTrace.enabled` or `OPENCLAW_CACHE_TRACE=1`) or debug logging, and trace results identify each request within its attempt. +Prompt-cache observations record `input`, `cacheRead`, and `cacheWrite` per completed foreground model request alongside its stable system-prefix, volatile-suffix, and tools fingerprints, and flag cache-read drops from the previous request, including reported zero reads; billing totals remain separate. A flagged drop lists the tracked changes since the last request (`model`, `cacheRetention`, `transport`, `streamStrategy`, `systemPrompt`, `systemPromptSuffix`, `tools`, `aggregateToolResultTruncation`). Trace results require cache tracing (`diagnostics.cacheTrace.enabled` or `OPENCLAW_CACHE_TRACE=1`) and identify each request within its attempt. + +OpenClaw also checks that each converted request history extends the previous request in the same session. An undeclared edit, removal, or reorder records `historyRewrite` and warns once for the session. Set `OPENCLAW_PROMPT_CACHE_ASSERT=1` to throw at the first differing message during development or tests. Compaction, pruning, transient runtime-context removal, and image cleanup declare their rewrites as `compaction`, `pruning`, `runtimeContextCarrier`, and `imageCleanup`; model, transport, or retention changes start a new history series. Content-block fingerprints reuse hashes only while every primitive property still matches; unchanged large text and image data are not hashed again. Nested block values and message envelopes are checked on every observation, so in-place edits are detected even when wrappers or content arrays are reused. String-content hashes live with the bounded history baseline. Provider-owned tool schema declarations are fingerprinted once per object identity. - Cache trace events are JSONL with staged snapshots like `session:loaded`, `prompt:before`, `stream:context`, and `session:after`. - Per-turn cache token impact is visible in normal usage surfaces: `cacheRead` and `cacheWrite` show up in `/usage tokens`, `/status`, session usage summaries, and custom `messages.usageTemplate` layouts. diff --git a/docs/reference/transcript-hygiene.md b/docs/reference/transcript-hygiene.md index 73757b535187..06348e1ddc28 100644 --- a/docs/reference/transcript-hygiene.md +++ b/docs/reference/transcript-hygiene.md @@ -124,6 +124,11 @@ turns, so a result adjacent to a repeated call stays with that occurrence. A dis result is moved only when exactly one unresolved occurrence can own it; ambiguous extras are dropped and missing occurrences receive synthetic error results. +Synthetic missing results tell the model that the outcome is unknown: retry only +read-only or idempotent operations, and verify current state before repeating an +operation that may have had side effects. Responses-family transports retain +their `aborted` placeholder. Neither placeholder proves that the tool did not run. + Implementation: `sanitizeToolUseResultPairing` in `src/agents/session-transcript-repair.ts` diff --git a/packages/agent-core/src/harness/session/tool-result-pairing.ts b/packages/agent-core/src/harness/session/tool-result-pairing.ts index 3064d23fc855..4bb35f36e16b 100644 --- a/packages/agent-core/src/harness/session/tool-result-pairing.ts +++ b/packages/agent-core/src/harness/session/tool-result-pairing.ts @@ -1,3 +1,4 @@ +import { DEFAULT_MISSING_TOOL_RESULT_TEXT } from "@openclaw/llm-core/types"; import { asOptionalObjectRecord } from "@openclaw/normalization-core/record-coerce"; import { normalizeOptionalString } from "@openclaw/normalization-core/string-coerce"; import { normalizeUniqueTrimmedStringList } from "@openclaw/normalization-core/string-normalization"; @@ -6,7 +7,7 @@ import type { SessionTreeEntry } from "../types.js"; const TOOL_CALL_TYPES = new Set(["toolCall", "toolUse", "functionCall"]); export const SYNTHETIC_MISSING_TOOL_RESULT_DETAIL_KEY = "openclawSyntheticMissingToolResult"; -export const DEFAULT_MISSING_TOOL_RESULT_TEXT = +export const LEGACY_MISSING_TOOL_RESULT_TEXT = "[openclaw] missing tool result in session history; inserted synthetic error result for transcript repair."; type ToolCallLike = { @@ -169,7 +170,7 @@ export function isSyntheticMissingToolResult(message: { Array.isArray(content) && content.some((block) => { const record = asOptionalObjectRecord(block); - return record?.type === "text" && record.text === DEFAULT_MISSING_TOOL_RESULT_TEXT; + return record?.type === "text" && record.text === LEGACY_MISSING_TOOL_RESULT_TEXT; }) ); } diff --git a/packages/ai/src/providers/google-interactions-shared.convert.test.ts b/packages/ai/src/providers/google-interactions-shared.convert.test.ts index c1022399d9f2..4cfaa5b95560 100644 --- a/packages/ai/src/providers/google-interactions-shared.convert.test.ts +++ b/packages/ai/src/providers/google-interactions-shared.convert.test.ts @@ -1,3 +1,4 @@ +import { DEFAULT_MISSING_TOOL_RESULT_TEXT } from "@openclaw/llm-core/types"; import { describe, expect, it } from "vitest"; import type { AssistantMessage, @@ -201,7 +202,7 @@ describe("buildGoogleInteractionsParams", () => { type: "function_result", call_id: "call_123", name: "getWeather", - result: [{ type: "text", text: "No result provided" }], + result: [{ type: "text", text: DEFAULT_MISSING_TOOL_RESULT_TEXT }], is_error: true, }, ]); @@ -301,7 +302,7 @@ describe("buildGoogleInteractionsParams", () => { type: "function_result", call_id: "call_legacy", name: "lookup", - result: [{ type: "text", text: "No result provided" }], + result: [{ type: "text", text: DEFAULT_MISSING_TOOL_RESULT_TEXT }], is_error: true, }, ]); diff --git a/packages/ai/src/transcript-transform.ts b/packages/ai/src/transcript-transform.ts index e633fd22e082..52c85af30a3a 100644 --- a/packages/ai/src/transcript-transform.ts +++ b/packages/ai/src/transcript-transform.ts @@ -1,9 +1,11 @@ +import { DEFAULT_MISSING_TOOL_RESULT_TEXT } from "@openclaw/llm-core/types"; import { isImageWithMediaPayload } from "./media-payload.js"; import { resolveModelBoundThinkingReplayMode } from "./providers/anthropic-model-contract.js"; import { FAILED_ASSISTANT_REPLAY_TEXT, resolveFailedAssistantReplay, } from "./replay-turn-classification.js"; +import { OPENAI_RESPONSES_APIS } from "./transports/openai-responses-contracts.js"; import type { Api, AssistantMessage, @@ -149,6 +151,9 @@ export function transformMessages( ) : messages; const supportsImages = model.input.includes("image"); + const missingResultText = OPENAI_RESPONSES_APIS.has(model.api) + ? "aborted" + : DEFAULT_MISSING_TOOL_RESULT_TEXT; const result: Message[] = []; const pendingAsyncCalls = new Map(); let pendingToolCalls: ToolCall[] = []; @@ -160,7 +165,8 @@ export function transformMessages( role: "toolResult", toolCallId: call.id, toolName: call.name, - content: [{ type: "text", text: "No result provided" }], + content: [{ type: "text", text: missingResultText }], + details: { openclawSyntheticMissingToolResult: true }, isError: true, timestamp: Date.now(), }); diff --git a/packages/ai/src/transports/openai-responses-compaction-replay.test.ts b/packages/ai/src/transports/openai-responses-compaction-replay.test.ts index 6f73d0f7224b..3dee1c3d6806 100644 --- a/packages/ai/src/transports/openai-responses-compaction-replay.test.ts +++ b/packages/ai/src/transports/openai-responses-compaction-replay.test.ts @@ -5,6 +5,7 @@ import type { Model, ProviderReplayState, } from "@openclaw/llm-core"; +import { DEFAULT_MISSING_TOOL_RESULT_TEXT } from "@openclaw/llm-core/types"; import { describe, expect, it } from "vitest"; import { makeTextToolResult } from "../../../../test/helpers/text-tool-result.js"; import { convertResponsesMessages as convertProviderResponsesMessages } from "../providers/openai-responses-shared.js"; @@ -776,7 +777,7 @@ describe("OpenAI Responses compaction replay", () => { expect(input.filter((item) => item.type === "function_call_output")).toMatchObject([ { call_id: "call_before", output: "before output" }, ]); - expect(JSON.stringify(input)).not.toContain("No result provided"); + expect(JSON.stringify(input)).not.toContain(DEFAULT_MISSING_TOOL_RESULT_TEXT); expect(JSON.stringify(input)).not.toContain("aborted"); expect(owner.content).toEqual([ expect.objectContaining({ type: "toolCall", id: "call_before|fc_before" }), @@ -812,7 +813,7 @@ describe("OpenAI Responses compaction replay", () => { "function_call_output", ]); expect(input.filter((item) => item.type === "function_call_output")).toMatchObject([ - { call_id: "call_after", output: "No result provided" }, + { call_id: "call_after", output: "aborted" }, ]); expect(input.filter((item) => item.type === "function_call_output")).toHaveLength(1); expect(JSON.stringify(input)).not.toContain("call_before"); diff --git a/packages/llm-core/src/types.ts b/packages/llm-core/src/types.ts index 58b977bb4913..b9344a384269 100644 --- a/packages/llm-core/src/types.ts +++ b/packages/llm-core/src/types.ts @@ -362,6 +362,9 @@ export const PROVIDER_FAILURE_WITH_OUTPUT_ERROR_CODE = "PROVIDER_FAILURE_WITH_OU /** Pre-dispatch argument rejection; callers still enforce output and effect guards. */ export const MALFORMED_TOOL_CALL_ARGUMENTS_ERROR_CODE = "malformed_tool_call_arguments"; +export const DEFAULT_MISSING_TOOL_RESULT_TEXT = + "Tool call interrupted before a result was recorded; its outcome is unknown. Retry only if the operation is read-only or idempotent. If it may have had side effects, verify the current state first instead of repeating it."; + /** User turn in a text-model conversation. */ export interface UserMessage { role: "user"; diff --git a/scripts/pr-lib/wrapper-components.txt b/scripts/pr-lib/wrapper-components.txt index 1db966590ed7..5c9bd3abeb98 100644 --- a/scripts/pr-lib/wrapper-components.txt +++ b/scripts/pr-lib/wrapper-components.txt @@ -41,6 +41,7 @@ packages/gateway-protocol/src/validation-errors.ts packages/gateway-protocol/src/worker-capacity.ts packages/llm-core/package.json packages/llm-core/src/model-data.ts +packages/llm-core/src/types.ts packages/media-core/src/constants.ts packages/media-core/src/file-name.ts packages/media-core/src/mime.ts diff --git a/src/agents/embedded-agent-runner.cache.live.test.ts b/src/agents/embedded-agent-runner.cache.live.test.ts index 0f5262fa78eb..3edd72c4b832 100644 --- a/src/agents/embedded-agent-runner.cache.live.test.ts +++ b/src/agents/embedded-agent-runner.cache.live.test.ts @@ -7,7 +7,7 @@ import { Type } from "typebox"; import { afterAll, beforeAll, describe, expect, it } from "vitest"; import type { OpenClawConfig } from "../config/config.js"; import { disposeOpenClawAgentDatabaseByPath } from "../state/openclaw-agent-db.js"; -import { deleteTestEnvValue, setTestEnvValue } from "../test-utils/env.js"; +import { captureEnv, setTestEnvValue } from "../test-utils/env.js"; import { prepareSystemAgentRunAdmission } from "./admitted-run-context.js"; import { runEmbeddedAgent } from "./embedded-agent-runner.js"; import { compactEmbeddedAgentSessionOnDemand } from "./embedded-agent-runner/compact.runtime.js"; @@ -73,13 +73,7 @@ const NOOP_TOOL: Tool = { let liveTestPngBase64 = ""; let liveRunnerPaths: { rootDir: string; agentDir: string; storePath: string } | undefined; let liveCacheTraceFile: string | undefined; -let previousCacheTraceEnv: { - enabled?: string; - file?: string; - messages?: string; - prompt?: string; - system?: string; -} | null = null; +let previousCacheTraceEnv: ReturnType | undefined; type UserContent = Extract["content"]; @@ -803,13 +797,15 @@ describeCacheLive("embedded agent runner prompt caching (live)", () => { }; liveCacheTraceFile = path.join(rootDir, "cache-trace.jsonl"); liveTestPngBase64 = (await fs.readFile(LIVE_TEST_PNG_URL)).toString("base64"); - previousCacheTraceEnv = { - enabled: process.env.OPENCLAW_CACHE_TRACE, - file: process.env.OPENCLAW_CACHE_TRACE_FILE, - messages: process.env.OPENCLAW_CACHE_TRACE_MESSAGES, - prompt: process.env.OPENCLAW_CACHE_TRACE_PROMPT, - system: process.env.OPENCLAW_CACHE_TRACE_SYSTEM, - }; + previousCacheTraceEnv = captureEnv([ + "OPENCLAW_CACHE_TRACE", + "OPENCLAW_PROMPT_CACHE_ASSERT", + "OPENCLAW_CACHE_TRACE_FILE", + "OPENCLAW_CACHE_TRACE_MESSAGES", + "OPENCLAW_CACHE_TRACE_PROMPT", + "OPENCLAW_CACHE_TRACE_SYSTEM", + ]); + setTestEnvValue("OPENCLAW_PROMPT_CACHE_ASSERT", "1"); setTestEnvValue("OPENCLAW_CACHE_TRACE", "1"); setTestEnvValue("OPENCLAW_CACHE_TRACE_FILE", liveCacheTraceFile); setTestEnvValue("OPENCLAW_CACHE_TRACE_MESSAGES", "0"); @@ -818,29 +814,8 @@ describeCacheLive("embedded agent runner prompt caching (live)", () => { }, 120_000); afterAll(async () => { - if (previousCacheTraceEnv) { - const restore = ( - key: - | "OPENCLAW_CACHE_TRACE" - | "OPENCLAW_CACHE_TRACE_FILE" - | "OPENCLAW_CACHE_TRACE_MESSAGES" - | "OPENCLAW_CACHE_TRACE_PROMPT" - | "OPENCLAW_CACHE_TRACE_SYSTEM", - value: string | undefined, - ) => { - if (value === undefined) { - deleteTestEnvValue(key); - } else { - setTestEnvValue(key, value); - } - }; - restore("OPENCLAW_CACHE_TRACE", previousCacheTraceEnv.enabled); - restore("OPENCLAW_CACHE_TRACE_FILE", previousCacheTraceEnv.file); - restore("OPENCLAW_CACHE_TRACE_MESSAGES", previousCacheTraceEnv.messages); - restore("OPENCLAW_CACHE_TRACE_PROMPT", previousCacheTraceEnv.prompt); - restore("OPENCLAW_CACHE_TRACE_SYSTEM", previousCacheTraceEnv.system); - } - previousCacheTraceEnv = null; + previousCacheTraceEnv?.restore(); + previousCacheTraceEnv = undefined; liveCacheTraceFile = undefined; if (liveRunnerPaths) { disposeOpenClawAgentDatabaseByPath(liveRunnerPaths.storePath); diff --git a/src/agents/embedded-agent-runner/compact.queued-execution.ts b/src/agents/embedded-agent-runner/compact.queued-execution.ts index 92fa428e9b7f..7f6ea59d4a00 100644 --- a/src/agents/embedded-agent-runner/compact.queued-execution.ts +++ b/src/agents/embedded-agent-runner/compact.queued-execution.ts @@ -46,6 +46,7 @@ import { import { runContextEngineMaintenance } from "./context-engine-maintenance.js"; import { resolveGlobalLane, resolveSessionLane } from "./lanes.js"; import { log } from "./logger.js"; +import { declarePromptHistoryRewrite } from "./prompt-cache-observability.js"; import { attachCompactionAccountingRecorder, type CompactionAccountingReceipt, @@ -318,6 +319,7 @@ export async function executeQueuedContextEngineCompaction(input: { requestBudget: host.requestBudget, pendingUserEntryId: host.pendingUserEntryId, recordCompaction: (receipt) => { + declarePromptHistoryRewrite({ ...runtimeTarget, reason: "compaction" }); committedCompaction = receipt; }, }); @@ -438,6 +440,9 @@ export async function executeQueuedContextEngineCompaction(input: { }, }); tokensAfter = result.result?.tokensAfter; + if (!committedCompaction) { + declarePromptHistoryRewrite({ ...runtimeTarget, reason: "compaction" }); + } } catch (error) { if (!params.abortSignal?.aborted) { throw error; diff --git a/src/agents/embedded-agent-runner/compaction-session-execution.ts b/src/agents/embedded-agent-runner/compaction-session-execution.ts index 2a5dcb620aec..bc042662aa33 100644 --- a/src/agents/embedded-agent-runner/compaction-session-execution.ts +++ b/src/agents/embedded-agent-runner/compaction-session-execution.ts @@ -64,6 +64,7 @@ import { buildEmbeddedExtensionFactories } from "./extensions.js"; import { getHistoryLimitFromSessionKey, limitHistoryTurns } from "./history.js"; import { log } from "./logger.js"; import type { PreparedCompactionRuntime } from "./prepared-compaction-runtime.js"; +import { declarePromptHistoryRewrite } from "./prompt-cache-observability.js"; import { sanitizeSessionHistory, validateReplayTurns } from "./replay-history.js"; import { createEmbeddedAgentResourceLoader } from "./resource-loader.js"; import { wrapStreamFnWithDiagnosticModelCallEvents } from "./run/attempt.model-diagnostic-events.js"; @@ -116,7 +117,6 @@ export async function executePreparedCompactionSession(runtime: PreparedCompacti try { const compactionTimeoutMs = resolveCompactionTimeoutMs(params.config); const accountingRecorder = readCompactionAccountingRecorder(params.contextEngineRuntimeContext); - const recordCompaction = accountingRecorder?.recordCompaction; const memoryTranscript = accountingRecorder?.memoryTranscript; const sessionTarget = memoryTranscript?.sessionTarget ?? @@ -129,6 +129,9 @@ export async function executePreparedCompactionSession(runtime: PreparedCompacti sessionKey: params.sessionKey, sessionTarget: params.sessionTarget, })); + const recordCompaction = + accountingRecorder?.recordCompaction ?? + (() => declarePromptHistoryRewrite({ ...sessionTarget, reason: "compaction" })); const assertActive = memoryTranscript?.assertActive ?? captureOwnedTranscriptWriteAssertion(sessionTarget); assertActive(); @@ -269,14 +272,8 @@ export async function executePreparedCompactionSession(runtime: PreparedCompacti ); session = createdSession.session; session[agentSessionSetContextReplacementHook]( - recordCompaction - ? (tokensAfter, tokensBefore) => - recordCompaction({ - tokensBefore, - tokensAfter, - compactionKind: "context-engine", - }) - : undefined, + (tokensAfter, tokensBefore) => + recordCompaction({ tokensBefore, tokensAfter, compactionKind: "context-engine" }), assertActive, ); session.setActiveToolsByName(sessionToolAllowlist); @@ -474,7 +471,7 @@ export async function executePreparedCompactionSession(runtime: PreparedCompacti enabled: compactionReplayEnabled, }, }); - recordCompaction?.({ + recordCompaction({ tokensBefore, tokensAfter: serverTokensAfter, compactionKind: "server-endpoint", diff --git a/src/agents/embedded-agent-runner/prompt-cache-observability.test.ts b/src/agents/embedded-agent-runner/prompt-cache-observability.test.ts index 73dc6a1d63fc..5df063e167fd 100644 --- a/src/agents/embedded-agent-runner/prompt-cache-observability.test.ts +++ b/src/agents/embedded-agent-runner/prompt-cache-observability.test.ts @@ -1,11 +1,19 @@ // Coverage for prompt-cache diagnostic tracking across turns. import { SYSTEM_PROMPT_CACHE_BOUNDARY } from "@openclaw/ai/internal/shared"; +import * as cryptoDigest from "@openclaw/normalization-core/node-crypto"; import { Type } from "typebox"; -import { beforeEach, describe, expect, it } from "vitest"; +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; +import { convertToLlm } from "../../../packages/agent-core/src/harness/messages.js"; +import type { Message, TextContent } from "../../llm/types.js"; +import { withEnv } from "../../test-utils/env.js"; +import type { AgentMessage } from "../runtime/index.js"; +import { makeAgentAssistantMessage } from "../test-helpers/agent-message-fixtures.js"; +import { log } from "./logger.js"; import { beginPromptCacheObservation, collectPromptCacheTools, completePromptCacheObservation, + declarePromptHistoryRewrite, recordAggregateTruncation, } from "./prompt-cache-observability.js"; import { createPromptCacheRequestObserver } from "./prompt-cache-request-observer.js"; @@ -23,6 +31,7 @@ function beginOpenAIObservation( params: Pick & Partial, ) { return beginPromptCacheObservation({ + messages: [], provider: "openai", modelId: "gpt-5.4", modelApi: "openai-responses", @@ -34,6 +43,292 @@ function beginOpenAIObservation( } describe("prompt cache observability", () => { + afterEach(() => vi.restoreAllMocks()); + + it("keeps a two-turn tool loop append-only with bounded block hashing", () => { + withEnv({ OPENCLAW_PROMPT_CACHE_ASSERT: "1" }, () => { + const sessionId = scopedKey("two-turn-loop"); + const source: AgentMessage[] = [ + { role: "user", content: [{ type: "text", text: "Read the fixture" }], timestamp: 1 }, + { + role: "custom", + customType: "fixture-context", + content: [{ type: "text", text: "Fixture context" }], + display: false, + timestamp: 2, + }, + ]; + const hashes = vi.spyOn(cryptoDigest, "sha256Hex"); + const observed = vi.fn(); + const observer = createPromptCacheRequestObserver( + { sessionId, streamStrategy: "test" }, + observed, + ); + const request = () => { + const messages = convertToLlm(source); + observer.onModelRequest( + { provider: "openai", id: "test-model", api: "openai-responses" }, + { messages }, + ); + observer.onModelUsage({ cacheRead: 8_000 }); + return messages; + }; + const first = request(); + source.push( + makeAgentAssistantMessage({ + content: [{ type: "toolCall", id: "read-1", name: "read", arguments: {} }], + stopReason: "toolUse", + timestamp: 3, + }), + { + role: "toolResult", + toolCallId: "read-1", + toolName: "read", + content: [{ type: "text", text: "fixture result" }], + isError: false, + timestamp: 4, + }, + ); + const loop = request(); + source.push( + makeAgentAssistantMessage({ + content: [{ type: "text", text: "Read complete" }], + timestamp: 5, + }), + { role: "user", content: "Summarize it", timestamp: 6 }, + ); + request(); + expect(loop[0]).toBe(first[0]); + expect(loop[1]).not.toBe(first[1]); + expect(loop[1]?.content).toBe(first[1]?.content); + expect( + hashes.mock.calls.filter( + ([value]) => typeof value === "string" && value.includes('"role":'), + ), + ).toHaveLength(12); + expect(hashes).toHaveBeenCalledTimes(26); + expect(observed).toHaveBeenCalledTimes(3); + for (const [observation] of observed.mock.calls) { + expect(observation.changes).toBeNull(); + } + }); + }); + + it.each([false, true])( + "detects mutated block text with a reused content array (rebuilt wrapper=%s)", + (rebuildWrapper) => { + const sessionId = scopedKey("mutated-text"); + const block: TextContent = { type: "text", text: "original" }; + const first: Message = { role: "user", content: "question", timestamp: 1 }; + const message = makeAgentAssistantMessage({ content: [block], timestamp: 2 }); + beginOpenAIObservation({ sessionId, messages: [first, message] }); + completePromptCacheObservation({ sessionId, usage: { cacheRead: 8_000 } }); + block.text = "rewritten"; + const detail = + "message 1 (assistant) differs from the previous request; history must be append-only"; + withEnv({ OPENCLAW_PROMPT_CACHE_ASSERT: "1" }, () => { + expect(() => + beginOpenAIObservation({ + sessionId, + messages: [first, rebuildWrapper ? { ...message } : message], + }), + ).toThrow(detail); + }); + expect( + completePromptCacheObservation({ sessionId, usage: { input: 8_000, cacheRead: 0 } }) + ?.changes, + ).toEqual([{ code: "historyRewrite", detail }]); + }, + ); + + it("detects nested tool arguments mutated in place", () => { + withEnv({ OPENCLAW_PROMPT_CACHE_ASSERT: "1" }, () => { + const sessionId = scopedKey("mutated-arguments"); + const args = { options: { path: "before" } }; + const message = makeAgentAssistantMessage({ + content: [{ type: "toolCall", id: "call-1", name: "read", arguments: args }], + }); + beginOpenAIObservation({ sessionId, messages: [message] }); + args.options.path = "after"; + expect(() => beginOpenAIObservation({ sessionId, messages: [message] })).toThrow( + "message 0 (assistant)", + ); + }); + }); + + it("detects primitive property additions, changes, and removals", () => { + withEnv({ OPENCLAW_PROMPT_CACHE_ASSERT: "1" }, () => { + const sessionId = scopedKey("mutated-properties"); + const block: TextContent = { type: "text", text: "stable" }; + const message = makeAgentAssistantMessage({ content: [block] }); + beginOpenAIObservation({ sessionId, messages: [message] }); + for (const mutate of [ + () => { + block.textSignature = undefined; + }, + () => { + block.textSignature = "signature"; + }, + () => { + delete block.textSignature; + }, + ]) { + mutate(); + expect(() => beginOpenAIObservation({ sessionId, messages: [message] })).toThrow( + "message 0 (assistant)", + ); + } + }); + }); + + it.each(["block", "string"] as const)( + "hashes unchanged large %s content once across three observations", + (kind) => { + const sessionId = scopedKey(`large-${kind}`); + const text = "large-content-fixture ".repeat(50_000); + const message: Message = { + role: "user", + content: kind === "string" ? text : [{ type: "text", text }], + timestamp: 1, + }; + const hashes = vi.spyOn(cryptoDigest, "sha256Hex"); + for (let index = 0; index < 3; index++) { + expect( + beginOpenAIObservation({ sessionId, messages: [{ ...message }] }).changes, + ).toBeNull(); + } + expect( + hashes.mock.calls.filter( + ([value]) => typeof value === "string" && value.includes("large-content-fixture"), + ), + ).toHaveLength(1); + expect(hashes).toHaveBeenCalledTimes(10); + }, + ); + + it.each(["edit", "remove", "reorder"] as const)( + "reports the first history divergence after %s", + (kind) => { + const first: Message = { role: "user", content: "first", timestamp: 1 }; + const second = makeAgentAssistantMessage({ + content: [{ type: "text", text: "second" }], + timestamp: 2, + }); + const messages = [first, second]; + const changed = + kind === "edit" + ? [first, { ...second, content: [{ type: "text" as const, text: "rewritten" }] }] + : kind === "remove" + ? [first] + : [second, first]; + const index = kind === "reorder" ? 0 : 1; + const detail = `message ${index} (assistant) differs from the previous request; history must be append-only`; + const warn = vi.spyOn(log, "warn").mockImplementation(() => {}); + const sessionId = scopedKey(kind); + beginOpenAIObservation({ sessionId, messages }); + completePromptCacheObservation({ sessionId, usage: { cacheRead: 8_000 } }); + withEnv({ OPENCLAW_PROMPT_CACHE_ASSERT: undefined }, () => { + expect(beginOpenAIObservation({ sessionId, messages: changed }).changes).toEqual([ + { code: "historyRewrite", detail }, + ]); + expect( + completePromptCacheObservation({ sessionId, usage: { input: 8_000, cacheRead: 0 } }) + ?.changes, + ).toEqual([{ code: "historyRewrite", detail }]); + beginOpenAIObservation({ sessionId, messages }); + expect(warn).toHaveBeenCalledTimes(1); + }); + withEnv({ OPENCLAW_PROMPT_CACHE_ASSERT: "1" }, () => { + expect(() => beginOpenAIObservation({ sessionId, messages: changed })).toThrow(detail); + }); + }, + ); + + it.each(["compaction", "pruning", "runtimeContextCarrier", "imageCleanup"] as const)( + "consumes a declared %s rewrite for exactly one request", + (reason) => { + withEnv({ OPENCLAW_PROMPT_CACHE_ASSERT: "1" }, () => { + const sessionId = scopedKey(reason); + const messages: Message[] = [{ role: "user", content: "before", timestamp: 1 }]; + beginOpenAIObservation({ sessionId, messages }); + declarePromptHistoryRewrite({ sessionId, reason }); + const rewritten: Message[] = [{ role: "user", content: "after", timestamp: 1 }]; + expect(beginOpenAIObservation({ sessionId, messages: rewritten }).changes).toEqual([ + { code: reason, detail: `${reason} changed provider history` }, + ]); + expect(beginOpenAIObservation({ sessionId, messages: rewritten }).changes).toBeNull(); + expect(() => beginOpenAIObservation({ sessionId, messages })).toThrow("message 0 (user)"); + }); + }, + ); + + it.each(["compaction", "pruning", "runtimeContextCarrier", "imageCleanup"] as const)( + "shares a %s declaration across cache affinities only within the same session key", + (reason) => { + withEnv({ OPENCLAW_PROMPT_CACHE_ASSERT: "1" }, () => { + const identity = { sessionId: scopedKey("affinities"), sessionKey: scopedKey("primary") }; + const before: Message[] = [{ role: "user", content: "before", timestamp: 1 }]; + const after: Message[] = [{ role: "user", content: "after", timestamp: 1 }]; + const keys = [scopedKey("affinity-a"), scopedKey("affinity-b")]; + for (const promptCacheKey of keys) { + beginOpenAIObservation({ ...identity, promptCacheKey, messages: before }); + } + const unrelated = { + ...identity, + sessionKey: scopedKey("other"), + promptCacheKey: scopedKey("unrelated-affinity"), + }; + beginOpenAIObservation({ ...unrelated, messages: before }); + declarePromptHistoryRewrite({ ...identity, promptCacheKey: keys[0], reason }); + for (const promptCacheKey of keys) { + expect( + beginOpenAIObservation({ ...identity, promptCacheKey, messages: after }).changes, + ).toEqual([{ code: reason, detail: `${reason} changed provider history` }]); + expect( + beginOpenAIObservation({ ...identity, promptCacheKey, messages: after }).changes, + ).toBeNull(); + } + expect(() => beginOpenAIObservation({ ...unrelated, messages: after })).toThrow( + "message 0 (user)", + ); + }); + }, + ); + + it.each([ + { modelId: "other" }, + { transport: "websocket" }, + { cacheRetention: "long" as const }, + { sessionId: "new-session" }, + ])("restarts the history series for %j", (change) => { + withEnv({ OPENCLAW_PROMPT_CACHE_ASSERT: "1" }, () => { + const identity = { sessionId: scopedKey("series"), promptCacheKey: scopedKey("affinity") }; + beginOpenAIObservation({ + ...identity, + messages: [{ role: "user", content: "previous", timestamp: 1 }], + }); + expect(() => beginOpenAIObservation({ ...identity, ...change, messages: [] })).not.toThrow(); + }); + }); + + it("does not reuse a content digest for a different tool-call identity", () => { + withEnv({ OPENCLAW_PROMPT_CACHE_ASSERT: "1" }, () => { + const sessionId = scopedKey("shared-content"); + const result: Message = { + role: "toolResult", + toolName: "read", + toolCallId: "first", + content: [], + isError: false, + timestamp: 1, + }; + beginOpenAIObservation({ sessionId, messages: [result] }); + expect(() => + beginOpenAIObservation({ sessionId, messages: [{ ...result, toolCallId: "second" }] }), + ).toThrow("message 0 (toolResult)"); + }); + }); + beforeEach(() => { currentTestScope = String(++testScope); }); @@ -67,6 +362,7 @@ describe("prompt cache observability", () => { observer.onModelRequest( { provider: "anthropic", id: "claude-sonnet-4-6", api: "anthropic-messages" }, { + messages: [], systemPrompt: "stable prefix", tools: [{ name: "read", description: "Read text", parameters: Type.Object({}) }], }, @@ -84,6 +380,7 @@ describe("prompt cache observability", () => { const sessionId = scopedKey("small-cache-miss"); const begin = (systemPrompt: string) => beginPromptCacheObservation({ + messages: [], sessionId, provider: "anthropic", modelId: "claude-sonnet-4-6", @@ -187,65 +484,29 @@ describe("prompt cache observability", () => { expect(numberPrototype[0]?.schemaDigest).not.toBe(noPrototype[0]?.schemaDigest); }); - it("bounds hostile, circular, and unreadable schema fingerprints", () => { - const circular: Record = { type: "object" }; - circular.self = circular; + it("memoizes cycle-safe schemas and skips unreadable tools", () => { + const parameters: Record = { type: "object" }; + parameters.self = parameters; + const digest = vi.spyOn(cryptoDigest, "sha256Hex"); const unreadable = { name: "unreadable", - get parameters(): unknown { - throw new Error("schema getter exploded"); + get parameters(): object { + throw new Error("unreadable schema"); }, }; - const oversized = { - name: "oversized", - parameters: { - type: "object", - properties: Object.fromEntries( - Array.from({ length: 1_000 }, (_, index) => [ - `property_${String(index).padStart(4, "0")}`, - { type: "string", description: "x".repeat(10_000) }, - ]), - ), - }, - }; - - expect( - collectPromptCacheTools([oversized, unreadable, { name: "circular", parameters: circular }]), - ).toEqual([ - { name: "circular", schemaDigest: expect.stringMatching(/^[a-f0-9]{64}$/) }, - { name: "oversized", schemaDigest: expect.stringMatching(/^[a-f0-9]{64}$/) }, - { name: "unreadable", schemaDigest: expect.stringMatching(/^[a-f0-9]{64}$/) }, + const first = collectPromptCacheTools([{ name: "read", parameters }, unreadable]); + expect(first).toEqual([ + { name: "read", schemaDigest: expect.stringMatching(/^[a-f0-9]{64}$/) }, ]); + expect(collectPromptCacheTools([{ name: "read", parameters }])).toEqual(first); + expect(digest).toHaveBeenCalledTimes(1); }); - it("rejects wide schemas before reading values and ignores their insertion order", () => { - let propertyReads = 0; - const createWideSchema = (reversed: boolean) => { - const properties: Record = {}; - const names = Array.from( - { length: 256 }, - (_, index) => `property_${String(index).padStart(4, "0")}`, - ); - for (const name of reversed ? names.toReversed() : names) { - Object.defineProperty(properties, name, { - enumerable: true, - get: () => { - propertyReads += 1; - return { type: "string" }; - }, - }); - } - return { type: "object", properties }; - }; - - const first = collectPromptCacheTools([{ name: "wide", parameters: createWideSchema(false) }]); - const reversed = collectPromptCacheTools([ - { name: "wide", parameters: createWideSchema(true) }, - ]); - - expect(reversed).toEqual(first); - expect(first[0]?.schemaDigest).toMatch(/^[a-f0-9]{64}$/); - expect(propertyReads).toBe(0); + it("fingerprints complete schemas independently of property insertion order", () => { + const collect = (parameters: object) => collectPromptCacheTools([{ name: "read", parameters }]); + expect(collect({ type: "object", properties: { a: {}, b: {} } })).toEqual( + collect({ properties: { b: {}, a: {} }, type: "object" }), + ); }); it("tracks cache-relevant changes and reports a real cache-read drop", () => { @@ -302,6 +563,7 @@ describe("prompt cache observability", () => { it("suppresses cache-break events for small drops", () => { beginPromptCacheObservation({ + messages: [], sessionId: scopedKey("session-1"), provider: "anthropic", modelId: "claude-sonnet-4-6", @@ -316,6 +578,7 @@ describe("prompt cache observability", () => { }); beginPromptCacheObservation({ + messages: [], sessionId: scopedKey("session-1"), provider: "anthropic", modelId: "claude-sonnet-4-6", @@ -357,6 +620,7 @@ describe("prompt cache observability", () => { const sessionId = scopedKey("dynamic-system-suffix"); const stablePrefix = "stable instructions and tool capability directory"; beginPromptCacheObservation({ + messages: [], sessionId, provider: "anthropic", modelId: "claude-sonnet-4-6", @@ -368,6 +632,7 @@ describe("prompt cache observability", () => { completePromptCacheObservation({ sessionId, usage: { cacheRead: 8_000 } }); const next = beginPromptCacheObservation({ + messages: [], sessionId, provider: "anthropic", modelId: "claude-sonnet-4-6", diff --git a/src/agents/embedded-agent-runner/prompt-cache-observability.ts b/src/agents/embedded-agent-runner/prompt-cache-observability.ts index c6f4eceb77fb..6b957466c457 100644 --- a/src/agents/embedded-agent-runner/prompt-cache-observability.ts +++ b/src/agents/embedded-agent-runner/prompt-cache-observability.ts @@ -4,24 +4,21 @@ import { } from "@openclaw/ai/internal/shared"; import { stableStringify } from "@openclaw/normalization-core"; import { sha256Hex } from "@openclaw/normalization-core/node-crypto"; -import { truncateUtf16Safe } from "@openclaw/normalization-core/utf16-slice"; +import type { ContextEnginePromptCacheObservationChange as PromptCacheChange } from "../../context-engine/types.js"; +import { createDedupeCache } from "../../infra/dedupe.js"; import { pruneMapToMaxSize } from "../../infra/map-size.js"; +import type { Message } from "../../llm/types.js"; import type { NormalizedUsage } from "../usage.js"; +import { log } from "./logger.js"; -type PromptCacheChangeCode = - | "aggregateToolResultTruncation" - | "cacheRetention" - | "model" - | "streamStrategy" - | "systemPrompt" - | "systemPromptSuffix" - | "tools" - | "transport"; +type PromptHistoryRewriteReason = + | "compaction" + | "pruning" + | "runtimeContextCarrier" + | "imageCleanup"; +type PromptCacheIdentity = { sessionId: string; promptCacheKey?: string; sessionKey?: string }; -export type PromptCacheChange = { - code: PromptCacheChangeCode; - detail: string; -}; +export type { PromptCacheChange }; type PromptCacheToolSnapshot = { name: string; @@ -32,7 +29,7 @@ type PromptCacheToolSnapshot = { type PromptCacheToolDescriptor = { readonly name?: string; readonly description?: string; - readonly parameters?: unknown; + readonly parameters?: object; }; type PromptCacheSnapshot = { @@ -50,19 +47,11 @@ type PromptCacheSnapshot = { toolNames: string[]; }; -type PromptCacheObservationStart = { - snapshot: PromptCacheSnapshot; - changes: PromptCacheChange[] | null; - previousCacheRead: number | null; -}; - -type PromptCacheBreak = { - previousCacheRead: number; - cacheRead: number; - changes: PromptCacheChange[] | null; -}; - type PromptCacheTracker = { + sessionId: string; + sessionKey?: string; + history: PromptHistoryFingerprint[]; + declaredRewrites?: Set; snapshot: PromptCacheSnapshot; lastCacheRead: number | null; /** Missing usage must not bind an older hit to a new request fingerprint. */ @@ -70,91 +59,83 @@ type PromptCacheTracker = { pendingChanges: PromptCacheChange[] | null; }; +type PromptHistoryFingerprint = { + digest: string; + role: string; + stringBlock?: { value: string; digest: string }; +}; + const trackers = new Map(); +const blockFingerprints = new WeakMap< + object, + { digest: string; primitives: [string, unknown][] } +>(); +// Schemas are provider-owned declarations; unlike transcript blocks, they do not mutate in place. +const toolSchemaFingerprints = new WeakMap(); const MAX_TRACKERS = 512; -const MAX_TOOL_SCHEMA_FINGERPRINT_DEPTH = 24; -const MAX_TOOL_SCHEMA_FINGERPRINT_NODES = 2_048; -const MAX_TOOL_SCHEMA_FINGERPRINT_ENTRIES = 128; -const MAX_TOOL_SCHEMA_FINGERPRINT_STRING_CHARS = 4_096; +const historyRewriteWarnings = createDedupeCache({ ttlMs: 0, maxSize: MAX_TRACKERS }); + +function fingerprintBlock(block: object): string { + const primitives: [string, unknown][] = []; + const nested: [string, unknown][] = []; + for (const entry of Object.entries(block)) { + const value = entry[1]; + if (value === null || ["string", "number", "boolean", "undefined"].includes(typeof value)) { + primitives.push(entry); + } else { + nested.push(entry); + } + } + const previous = blockFingerprints.get(block); + const unchanged = + previous && + previous.primitives.length === primitives.length && + primitives.every(([key, value], index) => { + const cached = previous.primitives[index]; + return cached?.[0] === key && cached[1] === value; + }); + const memo = unchanged + ? previous + : { digest: sha256Hex(stableStringify(Object.fromEntries(primitives))), primitives }; + if (!unchanged) { + blockFingerprints.set(block, memo); + } + return nested.length + ? sha256Hex(`${memo.digest}:${stableStringify(Object.fromEntries(nested))}`) + : memo.digest; +} + +function fingerprintMessage( + message: Message, + previous?: PromptHistoryFingerprint, +): PromptHistoryFingerprint { + const { content, ...envelope } = message; + // String blocks share the bounded history lifetime instead of a process-wide string cache. + let stringBlock: PromptHistoryFingerprint["stringBlock"]; + let blocks: string[]; + if (typeof content === "string") { + stringBlock = + previous?.stringBlock?.value === content + ? previous.stringBlock + : { value: content, digest: sha256Hex(stableStringify(content)) }; + blocks = [stringBlock.digest]; + } else { + blocks = content.map(fingerprintBlock); + } + return { + digest: sha256Hex(stableStringify([envelope, blocks])), + role: message.role, + stringBlock, + }; +} const MIN_CACHE_BREAK_TOKEN_DROP = 1_000; const MAX_STABLE_CACHE_READ_RATIO = 0.95; -function buildTrackerKey(params: { - promptCacheKey?: string; - sessionKey?: string; - sessionId: string; -}): string { +function buildTrackerKey(params: PromptCacheIdentity): string { return params.promptCacheKey?.trim() || params.sessionKey?.trim() || params.sessionId; } -function normalizeToolSchemaFingerprint( - value: unknown, - state: { remainingNodes: number; stack: WeakSet }, - depth = 0, -): unknown { - if (depth >= MAX_TOOL_SCHEMA_FINGERPRINT_DEPTH || state.remainingNodes <= 0) { - return "[schema fingerprint limit]"; - } - state.remainingNodes -= 1; - if (value === null || typeof value === "boolean") { - return value; - } - if (typeof value === "number") { - return Number.isFinite(value) ? value : "[non-finite number]"; - } - if (typeof value === "string") { - return truncateUtf16Safe(value, MAX_TOOL_SCHEMA_FINGERPRINT_STRING_CHARS); - } - if (typeof value !== "object") { - return `[${typeof value}]`; - } - if (state.stack.has(value)) { - return "[circular schema]"; - } - state.stack.add(value); - try { - if (Array.isArray(value)) { - const entries = value - .slice(0, MAX_TOOL_SCHEMA_FINGERPRINT_ENTRIES) - .map((entry) => normalizeToolSchemaFingerprint(entry, state, depth + 1)); - if (value.length > MAX_TOOL_SCHEMA_FINGERPRINT_ENTRIES) { - entries.push({ omitted: value.length - MAX_TOOL_SCHEMA_FINGERPRINT_ENTRIES }); - } - return entries; - } - const record = value as Record; - const keys: string[] = []; - for (const key in record) { - if (!Object.hasOwn(record, key)) { - continue; - } - // A schema above this limit receives one order-independent marker. Do - // not materialize or sort an attacker-controlled complete key list. - if (keys.length >= MAX_TOOL_SCHEMA_FINGERPRINT_ENTRIES) { - return "[schema key limit]"; - } - keys.push(key); - } - keys.sort((left, right) => (left < right ? -1 : left > right ? 1 : 0)); - // Schema property names are untrusted; a null prototype keeps "__proto__" - // as an own fingerprinted key instead of invoking a prototype setter. - const entries = Object.create(null) as Record; - for (const key of keys) { - try { - entries[key] = normalizeToolSchemaFingerprint(record[key], state, depth + 1); - } catch { - entries[key] = "[unreadable schema value]"; - } - } - return entries; - } catch { - return "[unreadable schema]"; - } finally { - state.stack.delete(value); - } -} - function setTracker(key: string, tracker: PromptCacheTracker): void { trackers.delete(key); pruneMapToMaxSize(trackers, MAX_TRACKERS - 1); @@ -177,7 +158,7 @@ function diffSnapshots( detail: `${previous.modelApi ?? "unknown"} -> ${next.modelApi ?? "unknown"}`, }); } - for (const code of ["cacheRetention", "transport"] as const) { + for (const code of ["cacheRetention", "transport", "streamStrategy"] as const) { if (previous[code] !== next[code]) { changes.push({ code, @@ -185,26 +166,16 @@ function diffSnapshots( }); } } - if (previous.streamStrategy !== next.streamStrategy) { - changes.push({ - code: "streamStrategy", - detail: `${previous.streamStrategy} -> ${next.streamStrategy}`, - }); - } - if (previous.systemPromptDigest !== next.systemPromptDigest) { - changes.push({ - code: "systemPrompt", - detail: "system prompt digest changed", - }); - } // OpenAI Responses routes send the suffix inline in `instructions`, so a // suffix change re-caches from that point; Anthropic-style checkpoints lose // the later conversation checkpoint. Track it separately from the prefix. - if (previous.systemPromptSuffixDigest !== next.systemPromptSuffixDigest) { - changes.push({ - code: "systemPromptSuffix", - detail: "system prompt suffix digest changed", - }); + for (const [code, detail] of [ + ["systemPrompt", "system prompt digest changed"], + ["systemPromptSuffix", "system prompt suffix digest changed"], + ] as const) { + if (previous[`${code}Digest`] !== next[`${code}Digest`]) { + changes.push({ code, detail }); + } } if (previous.toolDigest !== next.toolDigest) { changes.push({ @@ -228,29 +199,20 @@ export function collectPromptCacheTools( if (!name) { continue; } - const snapshot: PromptCacheToolSnapshot = { name }; - try { - if (typeof tool.description === "string") { - snapshot.descriptionDigest = sha256Hex(tool.description); + const { description, parameters } = tool; + let schemaDigest: string | undefined; + if (parameters) { + schemaDigest = toolSchemaFingerprints.get(parameters); + if (!schemaDigest) { + schemaDigest = sha256Hex(stableStringify(parameters)); + toolSchemaFingerprints.set(parameters, schemaDigest); } - } catch { - snapshot.descriptionDigest = sha256Hex("[unreadable tool description]"); } - try { - if (tool.parameters !== undefined) { - snapshot.schemaDigest = sha256Hex( - stableStringify( - normalizeToolSchemaFingerprint(tool.parameters, { - remainingNodes: MAX_TOOL_SCHEMA_FINGERPRINT_NODES, - stack: new WeakSet(), - }), - ), - ); - } - } catch { - snapshot.schemaDigest = sha256Hex("[unreadable tool schema]"); - } - snapshots.push(snapshot); + snapshots.push({ + name, + descriptionDigest: description === undefined ? undefined : sha256Hex(description), + schemaDigest, + }); } catch { continue; } @@ -258,19 +220,19 @@ export function collectPromptCacheTools( return sortPromptCacheToolsByName(snapshots); } -export function beginPromptCacheObservation(params: { - sessionId: string; - promptCacheKey?: string; - sessionKey?: string; - provider: string; - modelId: string; - modelApi?: string | null; - cacheRetention?: "none" | "short" | "long"; - streamStrategy: string; - transport?: string; - systemPrompt: string; - tools: readonly PromptCacheToolSnapshot[]; -}): PromptCacheObservationStart { +export function beginPromptCacheObservation( + params: PromptCacheIdentity & { + provider: string; + modelId: string; + modelApi?: string | null; + cacheRetention?: "none" | "short" | "long"; + streamStrategy: string; + transport?: string; + systemPrompt: string; + tools: readonly PromptCacheToolSnapshot[]; + messages: readonly Message[]; + }, +) { const key = buildTrackerKey(params); const tools = sortPromptCacheToolsByName(params.tools); const splitSystemPrompt = splitSystemPromptCacheBoundary(params.systemPrompt); @@ -290,6 +252,9 @@ export function beginPromptCacheObservation(params: { toolNames: tools.map((tool) => tool.name), }; const previous = trackers.get(key); + const history = params.messages.map((message, index) => + fingerprintMessage(message, previous?.history[index]), + ); const changes = previous ? [ ...(previous.pendingChanges?.filter( @@ -298,12 +263,45 @@ export function beginPromptCacheObservation(params: { ...(diffSnapshots(previous.snapshot, snapshot) ?? []), ] : []; + for (const code of previous?.declaredRewrites ?? []) { + changes.push({ code, detail: `${code} changed provider history` }); + } + const restarted = + previous?.sessionId !== params.sessionId || + changes.some( + ({ code }) => code === "model" || code === "transport" || code === "cacheRetention", + ); + const divergence = + previous && !restarted && !previous.declaredRewrites?.size + ? previous.history.findIndex((message, index) => message.digest !== history[index]?.digest) + : -1; + const violation = + divergence < 0 + ? undefined + : { + code: "historyRewrite" as const, + detail: `message ${divergence} (${history[divergence]?.role ?? previous!.history[divergence]!.role}) differs from the previous request; history must be append-only`, + }; + if (violation) { + changes.push(violation); + } setTracker(key, { + sessionId: params.sessionId, + sessionKey: params.sessionKey?.trim(), + history, snapshot, lastCacheRead: previous?.lastCacheRead ?? null, lastCacheReadSnapshot: previous?.lastCacheReadSnapshot, pendingChanges: changes.length > 0 ? changes : null, }); + if (violation) { + if (process.env.OPENCLAW_PROMPT_CACHE_ASSERT === "1") { + throw new Error(violation.detail); + } + if (!historyRewriteWarnings.check(params.sessionKey?.trim() || params.sessionId)) { + log.warn(`[prompt-cache] ${violation.detail} sessionKey=${params.sessionKey ?? key}`); + } + } return { snapshot, changes: changes.length > 0 ? changes : null, @@ -311,40 +309,48 @@ export function beginPromptCacheObservation(params: { }; } -export function recordAggregateTruncation(params: { - sessionId: string; - promptCacheKey?: string; - sessionKey?: string; -}): void { +export function declarePromptHistoryRewrite( + params: PromptCacheIdentity & { reason: PromptHistoryRewriteReason }, +): void { + // Session projections are shared; each cache-affinity baseline consumes the rewrite once. + for (const tracker of trackers.values()) { + if ( + tracker.sessionId === params.sessionId && + tracker.sessionKey === params.sessionKey?.trim() + ) { + (tracker.declaredRewrites ??= new Set()).add(params.reason); + } + } +} + +export function recordAggregateTruncation(params: PromptCacheIdentity): void { const tracker = trackers.get(buildTrackerKey(params)); const changes = tracker?.pendingChanges ?? []; if (!tracker || changes.some((change) => change.code === "aggregateToolResultTruncation")) { return; } - tracker.pendingChanges = [ - ...changes, - { - code: "aggregateToolResultTruncation", - detail: "aggregate tool-result truncation changed provider prompt", - }, - ]; + changes.push({ + code: "aggregateToolResultTruncation", + detail: "aggregate tool-result truncation changed provider prompt", + }); + tracker.pendingChanges = changes; } -export function completePromptCacheObservation(params: { - sessionId: string; - promptCacheKey?: string; - sessionKey?: string; - usage?: NormalizedUsage; -}): PromptCacheBreak | null { +export function completePromptCacheObservation( + params: PromptCacheIdentity & { + usage?: NormalizedUsage; + }, +) { const key = buildTrackerKey(params); const tracker = trackers.get(key); if (!tracker) { return null; } + const changes = tracker.pendingChanges; + tracker.pendingChanges = null; const cacheRead = params.usage?.cacheRead; if (typeof cacheRead !== "number" || !Number.isFinite(cacheRead)) { - tracker.pendingChanges = null; return null; } const previousCacheRead = tracker.lastCacheRead; @@ -353,7 +359,6 @@ export function completePromptCacheObservation(params: { tracker.lastCacheReadSnapshot = tracker.snapshot; if (previousCacheRead == null || previousCacheRead <= 0) { - tracker.pendingChanges = null; return null; } @@ -366,14 +371,11 @@ export function completePromptCacheObservation(params: { (params.usage?.input ?? 0) > 0 && previousSnapshot !== undefined && diffSnapshots(previousSnapshot, tracker.snapshot) === null; - const result = - hasMeaningfulDrop || completeMiss - ? { - previousCacheRead, - cacheRead, - changes: tracker.pendingChanges, - } - : null; - tracker.pendingChanges = null; - return result; + return hasMeaningfulDrop || completeMiss + ? { + previousCacheRead, + cacheRead, + changes, + } + : null; } diff --git a/src/agents/embedded-agent-runner/prompt-cache-request-observer.ts b/src/agents/embedded-agent-runner/prompt-cache-request-observer.ts index 1922670db417..d821f5db111b 100644 --- a/src/agents/embedded-agent-runner/prompt-cache-request-observer.ts +++ b/src/agents/embedded-agent-runner/prompt-cache-request-observer.ts @@ -24,7 +24,7 @@ export type PromptCacheRequestObservation = { export function createPromptCacheRequestObserver( params: Omit< Parameters[0], - "provider" | "modelId" | "modelApi" | "systemPrompt" | "tools" + "provider" | "modelId" | "modelApi" | "systemPrompt" | "tools" | "messages" >, onObservation: ( observation: PromptCacheRequestObservation, @@ -38,7 +38,7 @@ export function createPromptCacheRequestObserver( return { onModelRequest: ( model: Pick[0], "provider" | "id" | "api">, - context: Pick[1], "systemPrompt" | "tools">, + context: Pick[1], "systemPrompt" | "tools" | "messages">, ) => { requestIndex += 1; request = beginPromptCacheObservation({ @@ -48,6 +48,7 @@ export function createPromptCacheRequestObserver( modelApi: model.api, systemPrompt: context.systemPrompt ?? "", tools: collectPromptCacheTools(context.tools ?? []), + messages: context.messages, }); onRequest?.({ ...request, requestIndex }); }, diff --git a/src/agents/embedded-agent-runner/run/attempt-execution-phase.ts b/src/agents/embedded-agent-runner/run/attempt-execution-phase.ts index 5ca550ffc80a..7c3f43e8c470 100644 --- a/src/agents/embedded-agent-runner/run/attempt-execution-phase.ts +++ b/src/agents/embedded-agent-runner/run/attempt-execution-phase.ts @@ -12,6 +12,7 @@ import { } from "../../agent-run-terminal-outcome.js"; import { agentSessionSetContextReplacementHook } from "../../sessions/agent-session-compaction.js"; import { log } from "../logger.js"; +import { declarePromptHistoryRewrite } from "../prompt-cache-observability.js"; import type { EmbeddedAgentQueueHandle } from "../runs.js"; import { flushPendingToolResultsAfterIdle } from "../wait-for-idle-before-flush.js"; import { abortable as abortableWithSignal } from "./abortable.js"; @@ -44,6 +45,7 @@ export async function runEmbeddedAttemptExecutionPhase( throw new Error("embedded attempt requires an active admitted run"); } activeSession[agentSessionSetContextReplacementHook]((tokensAfter) => { + declarePromptHistoryRewrite({ ...attempt, reason: "compaction" }); toolBase.skillInstructionDeliveryCache.clear(); attempt.onContextAccountingEvent?.({ kind: "compaction", tokensAfter }); }, assertActive); diff --git a/src/agents/embedded-agent-runner/run/attempt-llm-boundary.ts b/src/agents/embedded-agent-runner/run/attempt-llm-boundary.ts index 48efd71e418f..356f9c51e857 100644 --- a/src/agents/embedded-agent-runner/run/attempt-llm-boundary.ts +++ b/src/agents/embedded-agent-runner/run/attempt-llm-boundary.ts @@ -52,6 +52,7 @@ type LlmBoundaryOptions = { projectPersistedSenderContext?: boolean; userTranscriptContexts?: readonly UserTranscriptContext[]; currentUserTimestampOverride?: CurrentUserTimestampMatch; + onRuntimeContextCarrierRemoved?: (removed: AgentMessage[]) => void; }; /** A session keeps its model projection across replay and process restarts. */ @@ -126,6 +127,12 @@ export function normalizeMessagesForLlmBoundary( const retained = options?.appendOnlyRuntimeContext ? withPersistedSenderContext : stripHistoricalRuntimeContextCustomMessages(withPersistedSenderContext); + if (retained.length < withPersistedSenderContext.length) { + const retainedMessages = new Set(retained); + options?.onRuntimeContextCarrierRemoved?.( + withPersistedSenderContext.filter((message) => !retainedMessages.has(message)), + ); + } return usesEscapedRuntimeContext(options?.sessionVersion) ? projectRuntimeContextMessages(retained) : retained; diff --git a/src/agents/embedded-agent-runner/run/attempt-prompt-submit.test-support.ts b/src/agents/embedded-agent-runner/run/attempt-prompt-submit.test-support.ts index f0f5e9862787..636110f6ac42 100644 --- a/src/agents/embedded-agent-runner/run/attempt-prompt-submit.test-support.ts +++ b/src/agents/embedded-agent-runner/run/attempt-prompt-submit.test-support.ts @@ -3,6 +3,7 @@ import { vi } from "vitest"; import type { ImageContent } from "../../../llm/types.js"; import type { AgentMessage } from "../../runtime/index.js"; import { agentSessionQueuePromptContext } from "../../sessions/agent-session-prompting.js"; +import { convertToLlm } from "../../sessions/messages.js"; import { getEmbeddedSessionPromptState } from "../session-prompt-state.js"; export const sessionId = "attempt-prompt-submit-test"; @@ -18,6 +19,7 @@ export function createSession() { const agent = { state, streamFn: baseStreamFn, + convertToLlm, transformContext: originalTransformContext, reset: () => { state.messages = []; diff --git a/src/agents/embedded-agent-runner/run/attempt-prompt-submit.test.ts b/src/agents/embedded-agent-runner/run/attempt-prompt-submit.test.ts index 5d1ad2297e54..78aa64bc5988 100644 --- a/src/agents/embedded-agent-runner/run/attempt-prompt-submit.test.ts +++ b/src/agents/embedded-agent-runner/run/attempt-prompt-submit.test.ts @@ -7,6 +7,7 @@ import { } from "../../../config/sessions/session-accessor.js"; import type { Context, ImageContent, Model } from "../../../llm/types.js"; import { createUserTurnTranscriptRecorder } from "../../../sessions/user-turn-transcript.js"; +import { withEnvAsync } from "../../../test-utils/env.js"; import { withOpenClawTestState } from "../../../test-utils/openclaw-test-state.js"; import { prepareSystemAgentRunAdmission } from "../../admitted-run-context.js"; import { readBtwTranscriptMessages } from "../../btw-transcript.js"; @@ -355,7 +356,7 @@ describe("submitEmbeddedAttemptPrompt", () => { await prepareEmbeddedAttemptSessionBoundary({ activeSession: session, appendOnlyRuntimeContext: retryContext !== "transient", - attempt: { prompt: user.content, userTurnTranscriptRecorder: recorder }, + attempt: { sessionId, prompt: user.content, userTurnTranscriptRecorder: recorder }, getUserTranscriptContexts: () => undefined, isRawModelRun: false, preparedUserTurnMessage: user, @@ -874,7 +875,7 @@ describe("submitEmbeddedAttemptPrompt", () => { ).toEqual([{ type: "text", text: oversized }]); }); - it("records aggregate truncation on a provider-bound cache break", async () => { + it("declares new pruning once and records aggregate truncation on a provider-bound cache break", async () => { const { activeSession } = createSession(); const input = createBaseInput(); const promptCacheKey = `${sessionId}:aggregate-truncation`; @@ -888,13 +889,6 @@ describe("submitEmbeddedAttemptPrompt", () => { systemPrompt: input.systemPrompt, tools: [], } as const; - beginPromptCacheObservation(observation); - completePromptCacheObservation({ - sessionId, - promptCacheKey, - usage: { cacheRead: 8_000 }, - }); - beginPromptCacheObservation(observation); activeSession.agent.state.messages = [ { role: "user", content: "call tools", timestamp: 1 }, { @@ -917,21 +911,34 @@ describe("submitEmbeddedAttemptPrompt", () => { { role: "assistant", content: [{ type: "text", text: "results processed" }], timestamp: 4 }, { role: "user", content: "continue", timestamp: 5 }, ] as AgentMessage[]; - activeSession.agent.streamFn = (() => undefined as never) as StreamFn; - - await submitEmbeddedAttemptPrompt({ - ...input, - attempt: { sessionId, promptCacheKey }, - activeSession, - toolResultAggregateMaxChars: 6_000, - promptActiveSession: async () => { - await activeSession.agent.streamFn( - {} as never, - { messages: activeSession.messages } as never, - {} as never, - ); - }, + beginPromptCacheObservation({ + ...observation, + messages: activeSession.messages as Context["messages"], }); + completePromptCacheObservation({ sessionId, promptCacheKey, usage: { cacheRead: 8_000 } }); + const changes: ReturnType["changes"][] = []; + activeSession.agent.streamFn = (_model, context) => { + changes.push( + beginPromptCacheObservation({ ...observation, messages: context.messages }).changes, + ); + return undefined as never; + }; + + const submit = () => + submitEmbeddedAttemptPrompt({ + ...input, + attempt: { sessionId, promptCacheKey }, + activeSession, + toolResultAggregateMaxChars: 6_000, + promptActiveSession: async () => { + await activeSession.agent.streamFn( + {} as never, + { messages: activeSession.messages } as never, + {} as never, + ); + }, + }); + await withEnvAsync({ OPENCLAW_PROMPT_CACHE_ASSERT: "1" }, submit); expect( completePromptCacheObservation({ @@ -947,7 +954,16 @@ describe("submitEmbeddedAttemptPrompt", () => { code: "aggregateToolResultTruncation", detail: "aggregate tool-result truncation changed provider prompt", }, + { code: "pruning", detail: "pruning changed provider history" }, ], }); + await withEnvAsync({ OPENCLAW_PROMPT_CACHE_ASSERT: "1" }, submit); + expect(changes).toEqual([ + [ + expect.objectContaining({ code: "aggregateToolResultTruncation" }), + expect.objectContaining({ code: "pruning" }), + ], + null, + ]); }); }); diff --git a/src/agents/embedded-agent-runner/run/attempt-prompt-submit.ts b/src/agents/embedded-agent-runner/run/attempt-prompt-submit.ts index e97c0e9eda1b..c4030a30c135 100644 --- a/src/agents/embedded-agent-runner/run/attempt-prompt-submit.ts +++ b/src/agents/embedded-agent-runner/run/attempt-prompt-submit.ts @@ -11,7 +11,10 @@ import { import type { AgentSession } from "../../sessions/index.js"; import { withSessionManagerWrite } from "../../sessions/session-manager-write-admission.js"; import { ackPendingAgentSteeringItems } from "../../subagents/registry/subagent-registry.js"; -import { recordAggregateTruncation } from "../prompt-cache-observability.js"; +import { + declarePromptHistoryRewrite, + recordAggregateTruncation, +} from "../prompt-cache-observability.js"; import { updateActiveEmbeddedRunSnapshot } from "../runs.js"; import { type getEmbeddedSessionPromptState, @@ -165,6 +168,9 @@ export async function submitEmbeddedAttemptPrompt(input: { input.toolResultPromptProjectionState, ); const providerMessages = providerPromptHistoryTruncation.messages; + if (providerPromptHistoryTruncation.truncatedCount > 0) { + declarePromptHistoryRewrite({ ...attempt, reason: "pruning" }); + } if (providerPromptHistoryTruncation.aggregateTruncatedCount > 0) { recordAggregateTruncation(attempt); } diff --git a/src/agents/embedded-agent-runner/run/attempt-session-boundary.test.ts b/src/agents/embedded-agent-runner/run/attempt-session-boundary.test.ts index 1952652f683a..153f00cf13d6 100644 --- a/src/agents/embedded-agent-runner/run/attempt-session-boundary.test.ts +++ b/src/agents/embedded-agent-runner/run/attempt-session-boundary.test.ts @@ -22,6 +22,7 @@ import { } from "../../../config/sessions/transcript-write-context.js"; import { buildTimestampPrefix } from "../../../gateway/server-methods/agent-timestamp.js"; import { MAIN_SESSION_RESTART_RECOVERY_SOURCE_TOOL } from "../../../sessions/input-provenance.js"; +import { withEnvAsync } from "../../../test-utils/env.js"; import { withOpenClawTestState } from "../../../test-utils/openclaw-test-state.js"; import type { AgentMessage } from "../../runtime/index.js"; import { guardSessionManager } from "../../session-tool-result-guard-wrapper.js"; @@ -29,6 +30,7 @@ import type { AgentSession } from "../../sessions/index.js"; import { convertToLlm as convertHarnessMessages } from "../../sessions/messages.js"; import { SessionManager } from "../../sessions/session-manager.js"; import { makeAssistantMessageFixture } from "../../test-helpers/assistant-message-fixtures.js"; +import { beginPromptCacheObservation } from "../prompt-cache-observability.js"; import { prepareEmbeddedAttemptSessionBoundary } from "./attempt-session-prepare.js"; import { buildRuntimeContextCustomMessage } from "./runtime-context-prompt.js"; @@ -124,6 +126,7 @@ async function withPersistedOrphanBoundary( input: { activeSession, attempt: { + sessionId: target.sessionId, ...(options.restartRecovery ? { inputProvenance: { @@ -162,7 +165,7 @@ describe("prepareEmbeddedAttemptSessionBoundary", () => { await prepareEmbeddedAttemptSessionBoundary({ activeSession, appendOnlyRuntimeContext: false, - attempt: { prompt: "next question" }, + attempt: { sessionId: "session-boundary", prompt: "next question" }, getUserTranscriptContexts: () => undefined, isRawModelRun: false, preparedUserTurnMessage: undefined, @@ -179,78 +182,111 @@ describe("prepareEmbeddedAttemptSessionBoundary", () => { it.each([false, true])( "replays turn and tool-loop prefixes with append-only runtime context %s", - async (appendOnlyRuntimeContext) => { - const { activeSession } = createActiveSession(); - await prepareEmbeddedAttemptSessionBoundary({ - activeSession, - appendOnlyRuntimeContext, - attempt: { - config: { agents: { defaults: { userTimezone: "UTC" } } }, - prompt: "first question", - }, - getUserTranscriptContexts: () => undefined, - isRawModelRun: false, - preparedUserTurnMessage: undefined, - sessionManager: createSessionManager(), - setActiveSessionSystemPrompt: vi.fn(), - }); - const user = { - role: "user" as const, - content: [ - { - type: "text" as const, - text: `${markInboundContextLabel("Conversation info:")}\n\`\`\`json\n{"channel":"discord"}\n\`\`\`\n\nfirst question`, + async (appendOnlyRuntimeContext) => + withEnvAsync({ OPENCLAW_PROMPT_CACHE_ASSERT: "1" }, async () => { + const { activeSession } = createActiveSession(); + activeSession.agent.convertToLlm = convertHarnessMessages; + const sessionId = `boundary-cache-${appendOnlyRuntimeContext}`; + const observe = (messages: Parameters[0]["messages"]) => + beginPromptCacheObservation({ + sessionId, + provider: "test-provider", + modelId: "test-model", + streamStrategy: "test", + systemPrompt: "Stable system prompt", + tools: [], + messages, + }); + await prepareEmbeddedAttemptSessionBoundary({ + activeSession, + appendOnlyRuntimeContext, + attempt: { + sessionId, + config: { agents: { defaults: { userTimezone: "UTC" } } }, + prompt: "first question", }, - ], - timestamp: 1_717_570_800_000, - }; - const carrier = buildRuntimeContextCustomMessage("first turn context")!; - const messages: AgentMessage[] = appendOnlyRuntimeContext ? [user, carrier] : [carrier, user]; - const first = await activeSession.agent.convertToLlm(messages); - expect(first).toHaveLength(2); - expect(first[1]).toBe(carrier); - messages.push( - makeAssistantMessageFixture({ - content: [{ type: "toolCall", id: "call_read", name: "read", arguments: {} }], - stopReason: "toolUse", - }), - { - role: "toolResult", - toolCallId: "call_read", - toolName: "read", - content: [{ type: "text", text: "result" }], - isError: false, - timestamp: user.timestamp + 1, - }, - ); - const toolLoop = await activeSession.agent.convertToLlm(messages); - if (appendOnlyRuntimeContext) { - expect(JSON.stringify(toolLoop.slice(0, first.length))).toBe(JSON.stringify(first)); - } else { - expect(toolLoop.at(-1)).toBe(carrier); - } - const nextUser = { - role: "user" as const, - content: "next question", - timestamp: user.timestamp + 60_000, - }; - const nextCarrier = buildRuntimeContextCustomMessage("second turn context")!; - messages.push(makeAssistantMessageFixture({ content: [{ type: "text", text: "done" }] })); - messages.push( - ...(appendOnlyRuntimeContext ? [nextUser, nextCarrier] : [nextCarrier, nextUser]), - ); - const next = await activeSession.agent.convertToLlm(messages); - if (appendOnlyRuntimeContext) { - expect(JSON.stringify(next.slice(0, toolLoop.length))).toBe(JSON.stringify(toolLoop)); - expect(next[1]).toBe(carrier); - expect(next.at(-1)).toBe(nextCarrier); - expect(next[0]!.content).toContain("Conversation info:"); - } else { - expect(next).not.toContain(carrier); - expect(next.at(-1)).toBe(nextCarrier); - expect(next[0]!.content).not.toContain("Conversation info:"); - } - }, + getUserTranscriptContexts: () => undefined, + isRawModelRun: false, + preparedUserTurnMessage: undefined, + sessionManager: createSessionManager(), + setActiveSessionSystemPrompt: vi.fn(), + }); + const user = { + role: "user" as const, + content: [ + { + type: "text" as const, + text: `${markInboundContextLabel("Conversation info:")}\n\`\`\`json\n{"channel":"discord"}\n\`\`\`\n\nfirst question`, + }, + ], + timestamp: 1_717_570_800_000, + }; + const carrier = buildRuntimeContextCustomMessage("first turn context")!; + const messages: AgentMessage[] = appendOnlyRuntimeContext + ? [user, carrier] + : [carrier, user]; + const first = await activeSession.agent.convertToLlm(messages); + expect(observe(first).changes).toBeNull(); + expect(first).toHaveLength(2); + expect(first[1]).toMatchObject({ + role: "user", + content: [{ type: "text", text: carrier.content }], + }); + messages.push( + makeAssistantMessageFixture({ + content: [{ type: "toolCall", id: "call_read", name: "read", arguments: {} }], + stopReason: "toolUse", + }), + { + role: "toolResult", + toolCallId: "call_read", + toolName: "read", + content: [{ type: "text", text: "result" }], + isError: false, + timestamp: user.timestamp + 1, + }, + ); + const toolLoop = await activeSession.agent.convertToLlm(messages); + expect(observe(toolLoop).changes).toEqual( + appendOnlyRuntimeContext + ? null + : [expect.objectContaining({ code: "runtimeContextCarrier" })], + ); + if (appendOnlyRuntimeContext) { + expect(JSON.stringify(toolLoop.slice(0, first.length))).toBe(JSON.stringify(first)); + } else { + expect(toolLoop.at(-1)).toEqual(first[1]); + } + expect(observe(await activeSession.agent.convertToLlm(messages)).changes).toBeNull(); + const nextUser = { + role: "user" as const, + content: "next question", + timestamp: user.timestamp + 60_000, + }; + const nextCarrier = buildRuntimeContextCustomMessage("second turn context")!; + messages.push(makeAssistantMessageFixture({ content: [{ type: "text", text: "done" }] })); + messages.push( + ...(appendOnlyRuntimeContext ? [nextUser, nextCarrier] : [nextCarrier, nextUser]), + ); + const next = await activeSession.agent.convertToLlm(messages); + expect(observe(next).changes).toEqual( + appendOnlyRuntimeContext + ? null + : [expect.objectContaining({ code: "runtimeContextCarrier" })], + ); + if (appendOnlyRuntimeContext) { + expect(JSON.stringify(next.slice(0, toolLoop.length))).toBe(JSON.stringify(toolLoop)); + expect(next[1]).toEqual(first[1]); + expect(next[0]!.content).toContain("Conversation info:"); + } else { + expect(next).not.toContainEqual(first[1]); + expect(next[0]!.content).not.toContain("Conversation info:"); + } + expect(next.at(-1)).toMatchObject({ + role: "user", + content: [{ type: "text", text: nextCarrier.content }], + }); + }), ); it.each([false, true])( @@ -261,7 +297,7 @@ describe("prepareEmbeddedAttemptSessionBoundary", () => { await prepareEmbeddedAttemptSessionBoundary({ activeSession, appendOnlyRuntimeContext, - attempt: { prompt: "question" }, + attempt: { sessionId: "session-boundary", prompt: "question" }, getUserTranscriptContexts: () => undefined, isRawModelRun: false, preparedUserTurnMessage: undefined, @@ -444,7 +480,7 @@ describe("prepareEmbeddedAttemptSessionBoundary", () => { const boundary = await prepareEmbeddedAttemptSessionBoundary({ activeSession, - attempt: { prompt: "exact probe" }, + attempt: { sessionId: "session-boundary", prompt: "exact probe" }, getUserTranscriptContexts: () => undefined, isRawModelRun: true, preparedUserTurnMessage: undefined, @@ -485,6 +521,7 @@ describe("prepareEmbeddedAttemptSessionBoundary", () => { const boundary = await prepareEmbeddedAttemptSessionBoundary({ activeSession, attempt: { + sessionId: "session-boundary", operation: "settled-tool-finalization", prompt: "finalize exactly", }, @@ -518,6 +555,7 @@ describe("prepareEmbeddedAttemptSessionBoundary", () => { const boundary = await prepareEmbeddedAttemptSessionBoundary({ activeSession, attempt: { + sessionId: "session-boundary", config: { agents: { defaults: { userTimezone: "UTC" } } }, prompt: "Current ask", }, @@ -560,7 +598,7 @@ describe("prepareEmbeddedAttemptSessionBoundary", () => { const { activeSession } = createActiveSession(); await prepareEmbeddedAttemptSessionBoundary({ activeSession, - attempt: { prompt: "The launch is Friday" }, + attempt: { sessionId: "session-boundary", prompt: "The launch is Friday" }, getUserTranscriptContexts: () => [{ runtimeMessage, transcriptMessage }], isRawModelRun: false, preparedUserTurnMessage: undefined, @@ -587,7 +625,7 @@ describe("prepareEmbeddedAttemptSessionBoundary", () => { const { activeSession } = createActiveSession(); await prepareEmbeddedAttemptSessionBoundary({ activeSession, - attempt: { prompt: "The launch is Friday" }, + attempt: { sessionId: "session-boundary", prompt: "The launch is Friday" }, getUserTranscriptContexts: () => [ { runtimeMessage: initialRuntime, @@ -634,7 +672,7 @@ describe("prepareEmbeddedAttemptSessionBoundary", () => { const { activeSession } = createActiveSession(); await prepareEmbeddedAttemptSessionBoundary({ activeSession, - attempt: { prompt: "same" }, + attempt: { sessionId: "session-boundary", prompt: "same" }, getUserTranscriptContexts: () => [ { runtimeMessage: secondRuntime, @@ -704,6 +742,7 @@ describe("prepareEmbeddedAttemptSessionBoundary", () => { const boundary = await prepareEmbeddedAttemptSessionBoundary({ activeSession, attempt: { + sessionId: "session-boundary", onUserMessagePersistenceInvalidated, prompt: "current prompt", suppressNextUserMessagePersistence, @@ -776,6 +815,7 @@ describe("prepareEmbeddedAttemptSessionBoundary", () => { const boundary = await prepareEmbeddedAttemptSessionBoundary({ activeSession, attempt: { + sessionId: "session-boundary", onUserMessagePersistenceInvalidated, prompt: "current prompt", userTurnTranscriptRecorder: recorder, @@ -839,6 +879,7 @@ describe("prepareEmbeddedAttemptSessionBoundary", () => { const boundary = await prepareEmbeddedAttemptSessionBoundary({ activeSession, attempt: { + sessionId: "session-boundary", onUserMessagePersistenceInvalidated, prompt: "current prompt", userTurnTranscriptRecorder: recorder, @@ -899,6 +940,7 @@ describe("prepareEmbeddedAttemptSessionBoundary", () => { const boundary = await prepareEmbeddedAttemptSessionBoundary({ activeSession, attempt: { + sessionId: "session-boundary", inputProvenance: { kind: "internal_system", sourceTool: MAIN_SESSION_RESTART_RECOVERY_SOURCE_TOOL, diff --git a/src/agents/embedded-agent-runner/run/attempt-session-prepare.ts b/src/agents/embedded-agent-runner/run/attempt-session-prepare.ts index 7263320f5b7d..93841a0f59e7 100644 --- a/src/agents/embedded-agent-runner/run/attempt-session-prepare.ts +++ b/src/agents/embedded-agent-runner/run/attempt-session-prepare.ts @@ -37,7 +37,9 @@ import { resolveToolSearchCatalogTool } from "../../tool-search.js"; import { runContextEngineMaintenance } from "../context-engine-maintenance.js"; import { buildEmbeddedExtensionFactories } from "../extensions.js"; import { log } from "../logger.js"; +import { declarePromptHistoryRewrite } from "../prompt-cache-observability.js"; import { createEmbeddedAgentResourceLoader } from "../resource-loader.js"; +import { recordRuntimeContextProjection } from "../session-prompt-state.js"; import { applySystemPromptToSession } from "../system-prompt.js"; import { prepareEmbeddedAttemptClientTools } from "./attempt-client-tools.js"; import { createAttemptCompactionThinkingResolver } from "./attempt-compaction-thinking.js"; @@ -372,6 +374,9 @@ type SessionBoundaryAttempt = Pick< | "onUserMessagePersistenceInvalidated" | "operation" | "prompt" + | "promptCacheKey" + | "sessionId" + | "sessionKey" | "skipPreparedUserTurnMessage" | "suppressNextUserMessagePersistence" | "userTurnTranscriptRecorder" @@ -498,25 +503,35 @@ export async function prepareEmbeddedAttemptSessionBoundary(input: { }; }; - if (typeof activeSession.agent.convertToLlm === "function") { - const baseConvertToLlm = activeSession.agent.convertToLlm.bind(activeSession.agent); - activeSession.agent.convertToLlm = async (messages) => { - const normalized = normalizeMessagesForLlmBoundary(messages, buildBoundaryOptions()); - const converted = await baseConvertToLlm( - // Persisted carriers stay after their user turn, including during tool loops; - // moving one would change the prefix bound to later thinking signatures. - input.appendOnlyRuntimeContext - ? normalized - : relocateCurrentRuntimeContextCarrierToTail(normalized), - ); - for (const message of converted) { - if (message.role === "user" && message.runtimeContextCarrier) { - message.runtimeContextCarrierRetained = input.appendOnlyRuntimeContext; - } + const baseConvertToLlm = activeSession.agent.convertToLlm.bind(activeSession.agent); + activeSession.agent.convertToLlm = async (messages) => { + let removedRuntimeContext: AgentMessage[] | undefined; + const normalized = normalizeMessagesForLlmBoundary(messages, { + ...buildBoundaryOptions(), + onRuntimeContextCarrierRemoved: (removed) => { + removedRuntimeContext = removed; + }, + }); + const converted = await baseConvertToLlm( + // Persisted carriers stay after their user turn, including during tool loops; + // moving one would change the prefix bound to later thinking signatures. + input.appendOnlyRuntimeContext + ? normalized + : relocateCurrentRuntimeContextCarrierToTail(normalized), + ); + for (const message of converted) { + if (message.role === "user" && message.runtimeContextCarrier) { + message.runtimeContextCarrierRetained = input.appendOnlyRuntimeContext; } - return converted; - }; - } + } + if ( + !input.appendOnlyRuntimeContext && + recordRuntimeContextProjection(attempt.sessionId, removedRuntimeContext, converted) + ) { + declarePromptHistoryRewrite({ ...attempt, reason: "runtimeContextCarrier" }); + } + return converted; + }; return { boundaryTimezone, diff --git a/src/agents/embedded-agent-runner/run/attempt-setup.ts b/src/agents/embedded-agent-runner/run/attempt-setup.ts index b38f95874563..fe3e35c6f594 100644 --- a/src/agents/embedded-agent-runner/run/attempt-setup.ts +++ b/src/agents/embedded-agent-runner/run/attempt-setup.ts @@ -38,7 +38,11 @@ import { invalidateComputerFrameIfMissing } from "../../tools/computer-tool.js"; import { resolveAttemptWorkspaceSandbox } from "../../workspace-sandbox.js"; import { isCacheTtlEligibleProvider, readLastCacheTtlTimestamp } from "../cache-ttl.js"; import { log } from "../logger.js"; -import type { ToolResultPromptProjectionState } from "../session-prompt-state.js"; +import { declarePromptHistoryRewrite } from "../prompt-cache-observability.js"; +import { + getEmbeddedSessionPromptState, + type ToolResultPromptProjectionState, +} from "../session-prompt-state.js"; import { installContextEngineLoopHook, installToolResultContextGuard, @@ -267,6 +271,7 @@ export function installEmbeddedAttemptContextGuards(input: { // replay so the prefix already sent for this session does not change. pruneNewRounds: !input.getServerToolClearingEnabled(), onPruned: () => { + declarePromptHistoryRewrite({ ...attempt, reason: "pruning" }); lastCacheTouchAt = Date.now(); }, }); @@ -362,6 +367,16 @@ export function installEmbeddedAttemptContextGuards(input: { : undefined, onCurrentTurnImageFailure: input.onCurrentTurnImageFailure, }, + (pruned) => { + const promptState = getEmbeddedSessionPromptState(attempt.sessionId); + const keys = new Set( + [...pruned].map(([index, message]) => `${index}:${message.role}:${message.timestamp}`), + ); + if ([...keys].some((key) => !promptState.prunedImageMessages?.has(key))) { + declarePromptHistoryRewrite({ ...attempt, reason: "imageCleanup" }); + } + promptState.prunedImageMessages = keys; + }, ); const previousComputerFrameTransform = activeSession.agent.transformContext; activeSession.agent.transformContext = async (messages, signal) => { diff --git a/src/agents/embedded-agent-runner/run/attempt-spawn-workspace.test-support.ts b/src/agents/embedded-agent-runner/run/attempt-spawn-workspace.test-support.ts index 351a20a46aea..1cf33492ab31 100644 --- a/src/agents/embedded-agent-runner/run/attempt-spawn-workspace.test-support.ts +++ b/src/agents/embedded-agent-runner/run/attempt-spawn-workspace.test-support.ts @@ -28,6 +28,7 @@ import type { Agent, AgentMessage, StreamFn } from "../../runtime/index.js"; import { agentSessionSetContextReplacementHook } from "../../sessions/agent-session-compaction.js"; import { agentSessionSetPromptPreparation } from "../../sessions/agent-session-prompting.js"; import type { AgentSession, CreateAgentSessionOptions } from "../../sessions/index.js"; +import { convertToLlm } from "../../sessions/messages.js"; import { getModelRegistryRuntime, initializeModelRegistryRuntime, @@ -814,7 +815,7 @@ type MutableSession = { isStreaming: boolean; subscribe: AgentSession["subscribe"]; agent: { - convertToLlm?: (messages: AgentMessage[]) => AgentMessage[] | Promise; + convertToLlm: Agent["convertToLlm"]; prompt?: (...args: unknown[]) => Promise; streamFn?: (...args: Parameters) => Promise; transport?: string; @@ -862,12 +863,7 @@ type SessionPromptOverride = ( options?: { images?: unknown[]; preflightResult?: (submitted: boolean) => void }, ) => Promise; -type TestAgentStream = { - result: () => Promise; - [Symbol.asyncIterator]: () => AsyncIterator; -}; - -function createCompletedAssistantStream(): TestAgentStream { +function createCompletedAssistantStream() { return { async result() { return { role: "assistant", content: "done" }; @@ -1013,6 +1009,7 @@ export function createDefaultEmbeddedSession(params?: { isStreaming: false, subscribe: () => () => {}, agent: { + convertToLlm, prompt: async (prompt, options) => { pendingPrompt = { prompt: String(prompt), diff --git a/src/agents/embedded-agent-runner/run/attempt-stream.ts b/src/agents/embedded-agent-runner/run/attempt-stream.ts index bdbf895b4258..26e136324642 100644 --- a/src/agents/embedded-agent-runner/run/attempt-stream.ts +++ b/src/agents/embedded-agent-runner/run/attempt-stream.ts @@ -184,36 +184,33 @@ export function installEmbeddedAttemptStreamGuards( ); } }; - const cacheObservabilityEnabled = Boolean(cacheTrace) || log.isEnabled("debug"); - const cacheObserver = cacheObservabilityEnabled - ? createPromptCacheRequestObserver( - { - sessionId: attempt.sessionId, - sessionKey: attempt.sessionKey, - promptCacheKey: attempt.promptCacheKey, - cacheRetention: effectivePromptCacheRetention, - streamStrategy, - transport: effectiveAgentTransport, - }, - (observation, snapshot) => { - if (observation.broke) { - const changes = - observation.changes?.map((change) => `${change.code}(${change.detail})`).join(", ") ?? - "no tracked cache input change"; - log.warn( - `[prompt-cache] cache read dropped ${observation.previousCacheRead} -> ${observation.cacheRead} ` + - `runId=${attempt.runId} request=${observation.requestIndex} for ${snapshot.provider}/${snapshot.modelId} via ${streamStrategy}; ${changes}`, - ); - } - cacheTrace?.recordStage("cache:result", { options: { ...observation } }); - }, - (request) => { - cacheTrace?.recordStage("cache:state", { - options: { ...request, previousCacheRead: request.previousCacheRead ?? undefined }, - }); - }, - ) - : undefined; + const cacheObserver = createPromptCacheRequestObserver( + { + sessionId: attempt.sessionId, + sessionKey: attempt.sessionKey, + promptCacheKey: attempt.promptCacheKey, + cacheRetention: effectivePromptCacheRetention, + streamStrategy, + transport: effectiveAgentTransport, + }, + (observation, snapshot) => { + if (observation.broke) { + const changes = + observation.changes?.map((change) => `${change.code}(${change.detail})`).join(", ") ?? + "no tracked cache input change"; + log.warn( + `[prompt-cache] cache read dropped ${observation.previousCacheRead} -> ${observation.cacheRead} ` + + `runId=${attempt.runId} request=${observation.requestIndex} for ${snapshot.provider}/${snapshot.modelId} via ${streamStrategy}; ${changes}`, + ); + } + cacheTrace?.recordStage("cache:result", { options: { ...observation } }); + }, + (request) => { + cacheTrace?.recordStage("cache:state", { + options: { ...request, previousCacheRead: request.previousCacheRead ?? undefined }, + }); + }, + ); if (cacheTrace) { cacheTrace.recordStage("session:loaded", { messages: session.messages, @@ -446,17 +443,15 @@ export function installEmbeddedAttemptStreamGuards( ); } return { - onModelRequest: cacheObserver?.onModelRequest, - onModelUsage: cacheObserver - ? (usage: NormalizedUsage | undefined) => { - // Async-tool fragments also end messages. result() marks the terminal - // response before core commits its final fragment with normalized usage. - if (modelResponseTerminal) { - modelResponseTerminal = false; - cacheObserver.onModelUsage(usage); - } - } - : undefined, - getPromptCacheObservation: cacheObserver?.getObservation, + onModelRequest: cacheObserver.onModelRequest, + onModelUsage: (usage: NormalizedUsage | undefined) => { + // Async-tool fragments also end messages. result() marks the terminal + // response before core commits its final fragment with normalized usage. + if (modelResponseTerminal) { + modelResponseTerminal = false; + cacheObserver.onModelUsage(usage); + } + }, + getPromptCacheObservation: cacheObserver.getObservation, }; } diff --git a/src/agents/embedded-agent-runner/run/attempt.llm-boundary.cache-stability.test.ts b/src/agents/embedded-agent-runner/run/attempt.llm-boundary.cache-stability.test.ts index 156fe697bec8..f25fca8e88c9 100644 --- a/src/agents/embedded-agent-runner/run/attempt.llm-boundary.cache-stability.test.ts +++ b/src/agents/embedded-agent-runner/run/attempt.llm-boundary.cache-stability.test.ts @@ -1,7 +1,7 @@ import path from "node:path"; import { streamOpenAICompletions, streamOpenAIResponses } from "@openclaw/ai/internal/openai"; import { expectDefined } from "@openclaw/normalization-core"; -import { afterAll, describe, expect, it } from "vitest"; +import { afterAll, afterEach, beforeEach, describe, expect, it } from "vitest"; import { resolveResponsesContinuationRequest } from "../../../../packages/ai/src/transports/openai-responses-continuation.js"; import { loadTranscriptEvents } from "../../../config/sessions/session-accessor.js"; import { buildTimestampPrefix } from "../../../gateway/server-methods/agent-timestamp.js"; @@ -13,6 +13,7 @@ import { type UserTurnInput, } from "../../../sessions/user-turn-transcript.js"; import { persistUserTurnTranscript } from "../../../sessions/user-turn-transcript.test-support.js"; +import { captureEnv, setTestEnvValue } from "../../../test-utils/env.js"; import { useSessionStoreTempDirs } from "../../../test-utils/session-state-cleanup.js"; import { INTERNAL_RUNTIME_CONTEXT_BEGIN, @@ -106,6 +107,12 @@ async function capture(api: "openai-completions" | "openai-responses", messages: } describe("prompt-cache boundary regressions", () => { + let env: ReturnType; + beforeEach(() => { + env = captureEnv(["OPENCLAW_PROMPT_CACHE_ASSERT"]); + setTestEnvValue("OPENCLAW_PROMPT_CACHE_ASSERT", "1"); + }); + afterEach(() => env.restore()); it("rejects unknown session projection versions before submitting history", () => { expect(() => normalizeMessagesForLlmBoundary([], { sessionVersion: 99 })).toThrow( "Unsupported session prompt projection version", diff --git a/src/agents/embedded-agent-runner/run/attempt.skills-policy.test.ts b/src/agents/embedded-agent-runner/run/attempt.skills-policy.test.ts index a43cadb41269..aaaadf51692f 100644 --- a/src/agents/embedded-agent-runner/run/attempt.skills-policy.test.ts +++ b/src/agents/embedded-agent-runner/run/attempt.skills-policy.test.ts @@ -98,6 +98,7 @@ function toolDigest( session = { sessionId: "embedded-session", sessionKey: "agent:main:main" }, ) { return beginPromptCacheObservation({ + messages: [], ...session, provider: "openai", modelId: "gpt-test", diff --git a/src/agents/embedded-agent-runner/run/attempt.test.ts b/src/agents/embedded-agent-runner/run/attempt.test.ts index f9689f7bc4df..4f43f1019b23 100644 --- a/src/agents/embedded-agent-runner/run/attempt.test.ts +++ b/src/agents/embedded-agent-runner/run/attempt.test.ts @@ -899,7 +899,7 @@ describe("wrapStreamFnSanitizeMalformedToolCalls", () => { expect(repairedToolResult.content).toEqual([ { type: "text", - text: "[openclaw] missing tool result in session history; inserted synthetic error result for transcript repair.", + text: "Tool call interrupted before a result was recorded; its outcome is unknown. Retry only if the operation is read-only or idempotent. If it may have had side effects, verify the current state first instead of repeating it.", }, ]); expect(repairedToolResult.isError).toBe(true); diff --git a/src/agents/embedded-agent-runner/run/compaction-runtime.ts b/src/agents/embedded-agent-runner/run/compaction-runtime.ts index 24269954b1a1..090ea372854c 100644 --- a/src/agents/embedded-agent-runner/run/compaction-runtime.ts +++ b/src/agents/embedded-agent-runner/run/compaction-runtime.ts @@ -30,6 +30,7 @@ import { } from "../compaction-successor.js"; import { resolveContextEngineCapabilities } from "../context-engine-capabilities.js"; import { log } from "../logger.js"; +import { declarePromptHistoryRewrite } from "../prompt-cache-observability.js"; import { mergeUsageIntoAccumulator, type UsageAccumulator } from "../usage-accumulator.js"; import { attachCompactionAccountingRecorder } from "./compaction-accounting-bridge.js"; import type { resolveCompactionLiveModelSelection } from "./compaction-live-model-selection.js"; @@ -90,6 +91,11 @@ export async function compactEmbeddedRunForRecovery( const { runParams } = input; const owner = input.prepareRecoveryOwner(); const activeSession = owner.session; + const promptCacheIdentity = { + ...runParams, + sessionId: activeSession.id, + sessionKey: input.resolvedSessionKey, + }; const reason = recovery.trigger === "budget" ? "context budget recovery" @@ -224,6 +230,7 @@ export async function compactEmbeddedRunForRecovery( : undefined, recordUsage: (usage) => mergeUsageIntoAccumulator(input.usageAccumulator, usage), recordCompaction: ({ tokensAfter }) => { + declarePromptHistoryRewrite({ ...promptCacheIdentity, reason: "compaction" }); observedCompactions += 1; input.state.observeContextAccounting({ kind: "compaction", tokensAfter }); }, @@ -289,6 +296,9 @@ export async function compactEmbeddedRunForRecovery( ? await input.adoptCompactionTranscript(result, sameTarget ? undefined : recordTokensAfter) : undefined; input.assertRecoveryActive(); + if (result.compacted && observedCompactions === 0) { + declarePromptHistoryRewrite({ ...promptCacheIdentity, reason: "compaction" }); + } return { result, runtimeContext, runtimeSettings, previousSessionId }; } diff --git a/src/agents/embedded-agent-runner/run/history-image-prune.ts b/src/agents/embedded-agent-runner/run/history-image-prune.ts index 83b7b10a82ba..8dc86ddbd6be 100644 --- a/src/agents/embedded-agent-runner/run/history-image-prune.ts +++ b/src/agents/embedded-agent-runner/run/history-image-prune.ts @@ -299,18 +299,31 @@ export function pruneProcessedHistoryImages(messages: AgentMessage[]): AgentMess export function installHistoryImagePruneContextTransform( agent: PrunableContextAgent, mediaOptions?: Parameters[1], + onPruned?: (messages: ReadonlyMap) => void, ): () => void { const originalTransformContext = agent.transformContext; agent.transformContext = async (messages: AgentMessage[], signal?: AbortSignal) => { - const prunedInput = pruneProcessedHistoryImages(messages) ?? messages; + const pruned = new Map(); + const prune = (source: AgentMessage[]) => { + const projected = pruneProcessedHistoryImages(source); + projected?.forEach((message, index) => { + const original = source[index]; + if (original && message !== original) { + pruned.set(index, original); + } + }); + return projected ?? source; + }; + const prunedInput = prune(messages); const hydratedInput = mediaOptions ? await hydratePromptMediaMessages(prunedInput, mediaOptions) : prunedInput; const transformed = originalTransformContext ? await originalTransformContext.call(agent, hydratedInput, signal) : hydratedInput; - const sourceMessages = Array.isArray(transformed) ? transformed : hydratedInput; - return pruneProcessedHistoryImages(sourceMessages) ?? sourceMessages; + const result = prune(transformed); + onPruned?.(pruned); + return result; }; return () => { agent.transformContext = originalTransformContext; diff --git a/src/agents/embedded-agent-runner/session-prompt-state.ts b/src/agents/embedded-agent-runner/session-prompt-state.ts index 45cfcc82d68f..fdea9007c081 100644 --- a/src/agents/embedded-agent-runner/session-prompt-state.ts +++ b/src/agents/embedded-agent-runner/session-prompt-state.ts @@ -2,6 +2,7 @@ import { sha256Hex } from "@openclaw/normalization-core/node-crypto"; import { isRecord } from "@openclaw/normalization-core/record-coerce"; import { pruneMapToMaxSize } from "../../infra/map-size.js"; +import type { Message } from "../../llm/types.js"; import { resolveGlobalSingleton } from "../../shared/global-singleton.js"; import type { AgentMessage } from "../runtime/index.js"; @@ -23,6 +24,9 @@ type EmbeddedSessionPromptState = { activeProjectKeys: string[]; toolResults: ToolResultPromptProjectionState; sentUserTurnIds: Set; + prunedImageMessages?: Set; + removedRuntimeContextKeys?: Set; + runtimeContextCarrierPositions?: number[]; }; const MAX_SESSION_PROMPT_STATES = 64; @@ -131,6 +135,26 @@ export function getEmbeddedSessionPromptState(sessionId: string): EmbeddedSessio return created; } +export function recordRuntimeContextProjection( + sessionId: string, + removed: readonly AgentMessage[] | undefined, + converted: readonly Message[], +): boolean { + const state = getEmbeddedSessionPromptState(sessionId); + const keys = removed?.map((message, index) => `${index}:${message.timestamp}`); + const positions = converted.flatMap((message, index) => + message.role === "user" && message.runtimeContextCarrier ? [index] : [], + ); + const changed = + keys?.some((key) => !state.removedRuntimeContextKeys?.has(key)) || + state.runtimeContextCarrierPositions?.some((position, index) => positions[index] !== position); + if (keys) { + state.removedRuntimeContextKeys = new Set(keys); + } + state.runtimeContextCarrierPositions = positions; + return Boolean(changed); +} + export function hashToolResultProjectionSnapshot( snapshot: ReturnType, ): string { diff --git a/src/agents/session-transcript-repair.test.ts b/src/agents/session-transcript-repair.test.ts index 630bc738c4a2..11efefe0a9a3 100644 --- a/src/agents/session-transcript-repair.test.ts +++ b/src/agents/session-transcript-repair.test.ts @@ -1,4 +1,5 @@ // Verifies transcript repair pairs tool calls/results and sanitizes tool inputs. +import { DEFAULT_MISSING_TOOL_RESULT_TEXT } from "@openclaw/llm-core/types"; import type { AgentMessage } from "openclaw/plugin-sdk/agent-core"; import { describe, expect, it } from "vitest"; import { @@ -11,9 +12,6 @@ import { import { castAgentMessage, castAgentMessages } from "./test-helpers/agent-message-fixtures.js"; import { sparseAssistant, textToolResult } from "./test-helpers/sparse-transcript.test-support.js"; -const DEFAULT_MISSING_TOOL_RESULT_TEXT = - "[openclaw] missing tool result in session history; inserted synthetic error result for transcript repair."; - const TOOL_CALL_BLOCK_TYPES = new Set([ "toolCall", "toolUse", @@ -478,6 +476,7 @@ describe("repairToolUseResultPairing prefers real result over synthetic error", toolCallId, toolName: "read", content: [{ type: "text", text: DEFAULT_MISSING_TOOL_RESULT_TEXT }], + details: { openclawSyntheticMissingToolResult: true }, isError: true, }; } @@ -1336,7 +1335,7 @@ describe("sanitizeToolCallInputs allowed-name filtering", () => { }); describe("stripToolResultDetails", () => { - it("removes details only from toolResult messages", () => { + it("strips opaque details and keeps synthetic projections stable", () => { const input = castAgentMessages([ { role: "toolResult", @@ -1347,6 +1346,10 @@ describe("stripToolResultDetails", () => { }, { role: "assistant", content: [{ type: "text", text: "keep me" }], details: { no: "touch" } }, { role: "user", content: "hello" }, + { + ...makeMissingToolResult({ toolCallId: "missing", toolName: "read" }), + details: { internal: true, openclawSyntheticMissingToolResult: true }, + }, ]); const out = stripToolResultDetails(input) as unknown as Array>; @@ -1358,17 +1361,8 @@ describe("stripToolResultDetails", () => { expect(Object.hasOwn(out[1] ?? {}, "details")).toBe(true); expect((out[1] ?? {}).role).toBe("assistant"); expect((out[2] ?? {}).role).toBe("user"); - }); - - it("returns the same array reference when there are no toolResult details", () => { - const input = castAgentMessages([ - { role: "assistant", content: [{ type: "text", text: "a" }] }, - textToolResult("call_1", "read", "ok"), - { role: "user", content: "b" }, - ]); - - const out = stripToolResultDetails(input); - expect(out).toBe(input); + expect(out[3]?.details).toEqual({ openclawSyntheticMissingToolResult: true }); + expect(stripToolResultDetails(out)).toBe(out); }); }); /* oxlint-disable max-lines -- TODO: split this grandfathered oversized file. */ diff --git a/src/agents/sessions/session-manager-model-context.test.ts b/src/agents/sessions/session-manager-model-context.test.ts index b3024726f12f..de6044d6a4bd 100644 --- a/src/agents/sessions/session-manager-model-context.test.ts +++ b/src/agents/sessions/session-manager-model-context.test.ts @@ -3,7 +3,7 @@ import fs from "node:fs"; import path from "node:path"; import type { StatementSync } from "node:sqlite"; import { expect, it, vi } from "vitest"; -import { DEFAULT_MISSING_TOOL_RESULT_TEXT } from "../../../packages/agent-core/src/harness/session/tool-result-pairing.js"; +import { LEGACY_MISSING_TOOL_RESULT_TEXT } from "../../../packages/agent-core/src/harness/session/tool-result-pairing.js"; import { makeUserMessage } from "../../../test/helpers/user-message.js"; import { appendTranscriptEvent, @@ -840,7 +840,7 @@ it.each(["details", "text", "duplicate-object", "late-array-call"])( content: [ { type: "text", - text: marker === "details" ? "missing" : DEFAULT_MISSING_TOOL_RESULT_TEXT, + text: marker === "details" ? "missing" : LEGACY_MISSING_TOOL_RESULT_TEXT, }, ], ...(marker === "details" ? { details: { openclawSyntheticMissingToolResult: true } } : {}), @@ -867,7 +867,7 @@ it.each(["details", "text", "duplicate-object", "late-array-call"])( await waitForSessionTranscriptProjection(scope); const database = openOpenClawAgentDatabase({ agentId: "main", path: scope.storePath }); // Preserve duplicate members from imported JSON; JavaScript objects would collapse them. - const content = `{"part":{"type":"text","text":"ordinary"},"part":{"type":"text","text":${JSON.stringify(DEFAULT_MISSING_TOOL_RESULT_TEXT)}}}`; + const content = `{"part":{"type":"text","text":"ordinary"},"part":{"type":"text","text":${JSON.stringify(LEGACY_MISSING_TOOL_RESULT_TEXT)}}}`; database.db .prepare( "UPDATE transcript_events SET event_json = json_set(event_json, '$.message.content', json(?)) WHERE session_id = ? AND json_extract(event_json, '$.id') = ?", diff --git a/src/agents/transport-message-transform.test.ts b/src/agents/transport-message-transform.test.ts index 9c90ca08e399..379f7f7e6e7f 100644 --- a/src/agents/transport-message-transform.test.ts +++ b/src/agents/transport-message-transform.test.ts @@ -1,5 +1,6 @@ // Transport message transform tests cover replay cleanup for provider-specific // tool-call/result sequencing before messages are sent back to transports. +import { DEFAULT_MISSING_TOOL_RESULT_TEXT } from "@openclaw/llm-core/types"; import type { Api, Context, Model } from "openclaw/plugin-sdk/llm"; import { describe, expect, it } from "vitest"; import { makeAssistantMessageFixture } from "./test-helpers/assistant-message-fixtures.js"; @@ -803,7 +804,7 @@ describe("transformTransportMessages synthetic tool-result policy", () => { expect(requireToolResultMessage(result[1])).toMatchObject({ toolCallId: "call_repeated", isError: true, - content: [{ type: "text", text: "No result provided" }], + content: [{ type: "text", text: DEFAULT_MISSING_TOOL_RESULT_TEXT }], }); expect(JSON.stringify(result)).not.toContain("failed turn output"); }); @@ -843,7 +844,9 @@ describe("transformTransportMessages synthetic tool-result policy", () => { ); expect(googleAlias.map((msg) => msg.role)).toEqual(["assistant", "toolResult", "user"]); const googleToolResult = requireToolResultMessage(googleAlias[1]); - expect(googleToolResult.content).toEqual([{ type: "text", text: "No result provided" }]); + expect(googleToolResult.content).toEqual([ + { type: "text", text: DEFAULT_MISSING_TOOL_RESULT_TEXT }, + ]); const bedrockCanonical = transformTransportMessages( messages, diff --git a/src/agents/transport-message-transform.ts b/src/agents/transport-message-transform.ts index c5a15e29a3d1..543da0e37f0e 100644 --- a/src/agents/transport-message-transform.ts +++ b/src/agents/transport-message-transform.ts @@ -10,6 +10,7 @@ import { isReasoningOnlyLengthAssistantTurn, resolveFailedAssistantReplay, } from "@openclaw/ai/internal/shared"; +import { DEFAULT_MISSING_TOOL_RESULT_TEXT } from "@openclaw/llm-core/types"; import { isRecord } from "@openclaw/normalization-core/record-coerce"; import { isSyntheticMissingToolResult, @@ -27,10 +28,7 @@ const SYNTHETIC_TOOL_RESULT_APIS = new Set([ ...OPENAI_RESPONSES_APIS, ]); -// "aborted" is the OpenAI Responses-family synthetic result convention, -// inherited from upstream Codex history normalization. It applies to public, -// Codex, Azure, and their OpenClaw transport aliases; Gemini/Anthropic use their -// own text. tool-replay-repair.live.test.ts exercises both paths against real models. +// Responses-family APIs retain their inherited "aborted" synthetic-result convention. /** Transforms transcript messages into a provider-safe replay context. */ export function transformTransportMessages( messages: Context["messages"], @@ -46,16 +44,12 @@ export function transformTransportMessages( preserveUnframedToolResults?: boolean; }, ): Context["messages"] { - const allowSyntheticToolResults = SYNTHETIC_TOOL_RESULT_APIS.has(model.api); const syntheticToolResultText = OPENAI_RESPONSES_APIS.has(model.api) ? "aborted" - : "No result provided"; + : DEFAULT_MISSING_TOOL_RESULT_TEXT; const toolCallIdMap = new Map(); let hasCrossModelAsyncCalls = false; const transformed = messages.map((msg) => { - if (msg.role === "user") { - return msg; - } if (msg.role === "toolResult") { // Earlier history repair may already have paired this call. Apply the same // transport placeholder without rewriting persisted diagnostics or real output. @@ -167,29 +161,22 @@ export function transformTransportMessages( // Pairing-aware transports must let shared repair see errored tool-call frames and // their adjacent results together; pre-filtering the call can misattribute its result // to an older turn that reused the same provider id. - const requiresPairing = allowSyntheticToolResults || hasCrossModelAsyncCalls; + const requiresPairing = SYNTHETIC_TOOL_RESULT_APIS.has(model.api) || hasCrossModelAsyncCalls; let replayLength = 0; transformed.forEach((msg, index) => { const original = messages[index]; - let replayMessage = msg; - if (original) { - if (isReasoningOnlyLengthAssistantTurn(original)) { - return; - } - switch (resolveFailedAssistantReplay(original, { pairingAware: requiresPairing })) { - case "drop": - return; - case "marker": - replayMessage = { - ...msg, - content: [{ type: "text", text: FAILED_ASSISTANT_REPLAY_TEXT }], - }; - break; - case "keep": - break; - } + if (original && isReasoningOnlyLengthAssistantTurn(original)) { + return; + } + const replay = original + ? resolveFailedAssistantReplay(original, { pairingAware: requiresPairing }) + : "keep"; + if (replay !== "drop") { + transformed[replayLength++] = + replay === "marker" + ? { ...msg, content: [{ type: "text", text: FAILED_ASSISTANT_REPLAY_TEXT }] } + : msg; } - transformed[replayLength++] = replayMessage; }); transformed.length = replayLength; diff --git a/src/config/sessions/session-accessor.sqlite-model-context.ts b/src/config/sessions/session-accessor.sqlite-model-context.ts index c0e3b7758a77..85b681fbcda4 100644 --- a/src/config/sessions/session-accessor.sqlite-model-context.ts +++ b/src/config/sessions/session-accessor.sqlite-model-context.ts @@ -5,7 +5,11 @@ import { iterateSessionContextMessages, projectSessionEntryMessage, } from "../../../packages/agent-core/src/harness/session/session.js"; -import { classifyToolUseResultPairing } from "../../../packages/agent-core/src/harness/session/tool-result-pairing.js"; +import { + classifyToolUseResultPairing, + isSyntheticMissingToolResult, + SYNTHETIC_MISSING_TOOL_RESULT_DETAIL_KEY, +} from "../../../packages/agent-core/src/harness/session/tool-result-pairing.js"; import { isCompactionReplayCheckpoint } from "../../../packages/ai/src/transports/provider-compaction-checkpoint.js"; import { executeSqliteQueryTakeFirstSync, @@ -580,7 +584,18 @@ function withTranscriptContextSnapshot( .where("seq", "in", [...bySeq.keys()]); for (const row of iterateSqliteQuerySync(database.db, query)) { const entry = bySeq.get(row.seq)!; - payloads.set(entry, hydrateContextEntry(row.event_json, entry)); + const hydrated = hydrateContextEntry(row.event_json, entry); + if ( + entry.type === "message" && + entry.message.role === "toolResult" && + isSyntheticMissingToolResult(entry.message) && + hydrated.type === "message" && + hydrated.message.role === "toolResult" + ) { + // Retain pairing provenance after SQL removes opaque tool details. + hydrated.message.details = { [SYNTHETIC_MISSING_TOOL_RESULT_DETAIL_KEY]: true }; + } + payloads.set(entry, hydrated); } } return payloads; diff --git a/src/config/sessions/session-model-context-projection.ts b/src/config/sessions/session-model-context-projection.ts index 15db35ee6874..60ff365993cd 100644 --- a/src/config/sessions/session-model-context-projection.ts +++ b/src/config/sessions/session-model-context-projection.ts @@ -1,6 +1,6 @@ import { sql, type Expression, type RawBuilder } from "kysely"; import { - DEFAULT_MISSING_TOOL_RESULT_TEXT, + LEGACY_MISSING_TOOL_RESULT_TEXT, SYNTHETIC_MISSING_TOOL_RESULT_DETAIL_KEY, } from "../../../packages/agent-core/src/harness/session/tool-result-pairing.js"; import { supportsNodeSqliteJsonb } from "../../infra/node-sqlite.js"; @@ -169,7 +169,7 @@ export function projectModelContextNavigationSql( AND ${contentPropertySql(event, "type")} IN ('toolCall', 'toolUse', 'functionCall'))`; const synthetic = /* kysely-allow-raw: pairing prefers real results over synthetic missing-result placeholders. */ sql`COALESCE(json_extract(${event}, ${`$.message.details.${SYNTHETIC_MISSING_TOOL_RESULT_DETAIL_KEY}`}), 0) = 1 OR EXISTS ( SELECT 1 FROM json_each(${event}, '$.message.content') WHERE type = 'object' - AND ${contentPropertySql(event, "type")} = 'text' AND ${contentPropertySql(event, "text")} = ${DEFAULT_MISSING_TOOL_RESULT_TEXT})`; + AND ${contentPropertySql(event, "type")} = 'text' AND ${contentPropertySql(event, "text")} = ${LEGACY_MISSING_TOOL_RESULT_TEXT})`; return /* kysely-allow-raw: retain readable empty bodies only for navigation outside the model window. */ sql`CASE json_extract(${event}, '$.type') WHEN 'message' THEN json_set(${entry}, '$.message', json_set(${messageFacts}, '$.content', json(${calls}), '$.command', '', '$.output', '', diff --git a/src/context-engine/types.ts b/src/context-engine/types.ts index 0b045060e2ea..be2b3c74816f 100644 --- a/src/context-engine/types.ts +++ b/src/context-engine/types.ts @@ -249,6 +249,11 @@ type ContextEnginePromptCacheUsage = { }; type ContextEnginePromptCacheObservationChangeCode = + | "historyRewrite" + | "compaction" + | "pruning" + | "runtimeContextCarrier" + | "imageCleanup" | "aggregateToolResultTruncation" | "cacheRetention" | "model" @@ -258,7 +263,7 @@ type ContextEnginePromptCacheObservationChangeCode = | "tools" | "transport"; -type ContextEnginePromptCacheObservationChange = { +export type ContextEnginePromptCacheObservationChange = { code: ContextEnginePromptCacheObservationChangeCode; detail: string; }; diff --git a/src/gateway/server-methods/chat-send-synthetic-repair.integration.test.ts b/src/gateway/server-methods/chat-send-synthetic-repair.integration.test.ts index b2c4a58fd307..860d5bc25a39 100644 --- a/src/gateway/server-methods/chat-send-synthetic-repair.integration.test.ts +++ b/src/gateway/server-methods/chat-send-synthetic-repair.integration.test.ts @@ -6,9 +6,10 @@ import path from "node:path"; import { DatabaseSync } from "node:sqlite"; import { text as readText } from "node:stream/consumers"; import type { MessageCreateParamsStreaming } from "@anthropic-ai/sdk/resources/messages"; +import { DEFAULT_MISSING_TOOL_RESULT_TEXT } from "@openclaw/llm-core/types"; import { expect, it } from "vitest"; import { - DEFAULT_MISSING_TOOL_RESULT_TEXT, + LEGACY_MISSING_TOOL_RESULT_TEXT, makeMissingToolResult, } from "../../../packages/agent-core/src/harness/session/tool-result-pairing.js"; import { makeUserMessage } from "../../../test/helpers/user-message.js"; @@ -202,6 +203,7 @@ it("chat.send replays synthetic repairs through session history and the register const missing = makeMissingToolResult({ toolCallId: "callmissing", toolName: "read" }); if (legacy) { delete missing.details; + missing.content = [{ type: "text", text: LEGACY_MISSING_TOOL_RESULT_TEXT }]; } manager.appendMessage(missing); if (late) { @@ -256,7 +258,7 @@ it("chat.send replays synthetic repairs through session history and the register expect.soft(results, scenario).toEqual([ expect.objectContaining({ tool_use_id: "callmissing", - content: late ? "LATE_ACTUAL_RESULT" : "No result provided", + content: late ? "LATE_ACTUAL_RESULT" : DEFAULT_MISSING_TOOL_RESULT_TEXT, is_error: !late, }), expect.objectContaining({ diff --git a/src/shared/model-context-message.ts b/src/shared/model-context-message.ts index 6459de975385..54cf50f1fc6a 100644 --- a/src/shared/model-context-message.ts +++ b/src/shared/model-context-message.ts @@ -1,5 +1,6 @@ import type { AgentMessage } from "@openclaw/agent-core"; import { asOptionalRecord } from "@openclaw/normalization-core/record-coerce"; +import { SYNTHETIC_MISSING_TOOL_RESULT_DETAIL_KEY } from "../../packages/agent-core/src/harness/session/tool-result-pairing.js"; // Native replay retains the exact submitted prompt. Model-context consumers already // have its visible content; copying this storage-only payload duplicates the prompt. @@ -15,7 +16,15 @@ export function stripToolResultDetails(messages: unknown[]): unknown[] { return message; } const sanitized = { ...record }; - delete sanitized.details; + const details = asOptionalRecord(record.details); + if (details?.[SYNTHETIC_MISSING_TOOL_RESULT_DETAIL_KEY] === true) { + if (Object.keys(details).length === 1) { + return message; + } + sanitized.details = { [SYNTHETIC_MISSING_TOOL_RESULT_DETAIL_KEY]: true }; + } else { + delete sanitized.details; + } touched = true; return sanitized; }); diff --git a/ui/vitest.config.ts b/ui/vitest.config.ts index 2608687b56ed..88d0f7deeee8 100644 --- a/ui/vitest.config.ts +++ b/ui/vitest.config.ts @@ -62,6 +62,10 @@ const workspaceSourceAliases = [ ...sharedVitestConfig.resolve.alias.filter( (alias) => typeof alias.find === "string" && alias.find.startsWith("openclaw/plugin-sdk/"), ), + { + find: "@openclaw/llm-core/types", + replacement: path.resolve(repoRoot, "packages/llm-core/src/types.ts"), + }, { find: /^@openclaw\/model-catalog-core\/(.+)$/u, replacement: path.resolve(repoRoot, "packages/model-catalog-core/src/$1.ts"),