mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 01:29:56 +00:00
improve(agents): explain unknown tool outcomes and assert append-only request history (#163379)
* fix(agents): describe unrecorded tool results as unknown outcomes
Use shared unknown-outcome guidance while preserving the Responses-family aborted convention. Retain synthetic provenance in model context so late real results can replace repair placeholders without changing existing detection or stored transcript rows.
* improve(agents): assert provider request history stays append-only
Track converted message prefixes and record declared compaction, pruning, runtime-context, and image rewrites. Notify every cache-affinity baseline for the same session identity. Memoize message/content and schema fingerprints, with strict assertions enabled only by an opt-in environment flag.
* fix(agents): pin the shared tool-result text for the request wrapper
Include packages/llm-core/src/types.ts in the PR wrapper's extracted source inventory. Main commit 93625aa3d1 pulled tool-result-pairing.ts into that graph, so its existing shared tool-result text import must be included for standalone wrapper execution.
This commit is contained in:
parent
a1566b2dbf
commit
cee0906f38
41 changed files with 962 additions and 527 deletions
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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`
|
||||
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
})
|
||||
);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
},
|
||||
]);
|
||||
|
|
|
|||
|
|
@ -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<TApi extends Api>(
|
|||
)
|
||||
: 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<string, ToolCall>();
|
||||
let pendingToolCalls: ToolCall[] = [];
|
||||
|
|
@ -160,7 +165,8 @@ export function transformMessages<TApi extends Api>(
|
|||
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(),
|
||||
});
|
||||
|
|
|
|||
|
|
@ -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");
|
||||
|
|
|
|||
|
|
@ -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";
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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<typeof captureEnv> | undefined;
|
||||
|
||||
type UserContent = Extract<Message, { role: "user" }>["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);
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
|
|
|
|||
|
|
@ -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<ObservationParams, "sessionId"> & Partial<ObservationParams>,
|
||||
) {
|
||||
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<string, unknown> = { type: "object" };
|
||||
circular.self = circular;
|
||||
it("memoizes cycle-safe schemas and skips unreadable tools", () => {
|
||||
const parameters: Record<string, unknown> = { 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<string, unknown> = {};
|
||||
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",
|
||||
|
|
|
|||
|
|
@ -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<PromptHistoryRewriteReason>;
|
||||
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<string, PromptCacheTracker>();
|
||||
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<object, string>();
|
||||
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<object> },
|
||||
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<string, unknown>;
|
||||
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<string, unknown>;
|
||||
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;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -24,7 +24,7 @@ export type PromptCacheRequestObservation = {
|
|||
export function createPromptCacheRequestObserver(
|
||||
params: Omit<
|
||||
Parameters<typeof beginPromptCacheObservation>[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<Parameters<StreamFn>[0], "provider" | "id" | "api">,
|
||||
context: Pick<Parameters<StreamFn>[1], "systemPrompt" | "tools">,
|
||||
context: Pick<Parameters<StreamFn>[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 });
|
||||
},
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
|
|||
|
|
@ -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 = [];
|
||||
|
|
|
|||
|
|
@ -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<typeof beginPromptCacheObservation>["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,
|
||||
]);
|
||||
});
|
||||
});
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<typeof beginPromptCacheObservation>[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,
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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) => {
|
||||
|
|
|
|||
|
|
@ -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<AgentMessage[]>;
|
||||
convertToLlm: Agent["convertToLlm"];
|
||||
prompt?: (...args: unknown[]) => Promise<unknown>;
|
||||
streamFn?: (...args: Parameters<StreamFn>) => Promise<unknown>;
|
||||
transport?: string;
|
||||
|
|
@ -862,12 +863,7 @@ type SessionPromptOverride = (
|
|||
options?: { images?: unknown[]; preflightResult?: (submitted: boolean) => void },
|
||||
) => Promise<void>;
|
||||
|
||||
type TestAgentStream = {
|
||||
result: () => Promise<unknown>;
|
||||
[Symbol.asyncIterator]: () => AsyncIterator<unknown>;
|
||||
};
|
||||
|
||||
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),
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
};
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<typeof captureEnv>;
|
||||
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",
|
||||
|
|
|
|||
|
|
@ -98,6 +98,7 @@ function toolDigest(
|
|||
session = { sessionId: "embedded-session", sessionKey: "agent:main:main" },
|
||||
) {
|
||||
return beginPromptCacheObservation({
|
||||
messages: [],
|
||||
...session,
|
||||
provider: "openai",
|
||||
modelId: "gpt-test",
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
|
|
|
|||
|
|
@ -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 };
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -299,18 +299,31 @@ export function pruneProcessedHistoryImages(messages: AgentMessage[]): AgentMess
|
|||
export function installHistoryImagePruneContextTransform(
|
||||
agent: PrunableContextAgent,
|
||||
mediaOptions?: Parameters<typeof hydratePromptMediaMessages>[1],
|
||||
onPruned?: (messages: ReadonlyMap<number, AgentMessage>) => void,
|
||||
): () => void {
|
||||
const originalTransformContext = agent.transformContext;
|
||||
agent.transformContext = async (messages: AgentMessage[], signal?: AbortSignal) => {
|
||||
const prunedInput = pruneProcessedHistoryImages(messages) ?? messages;
|
||||
const pruned = new Map<number, AgentMessage>();
|
||||
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;
|
||||
|
|
|
|||
|
|
@ -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<string>;
|
||||
prunedImageMessages?: Set<string>;
|
||||
removedRuntimeContextKeys?: Set<string>;
|
||||
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<typeof serializeCacheTtlToolResultProjections>,
|
||||
): string {
|
||||
|
|
|
|||
|
|
@ -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<Record<string, unknown>>;
|
||||
|
|
@ -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. */
|
||||
|
|
|
|||
|
|
@ -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') = ?",
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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<string>([
|
|||
...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<string, string>();
|
||||
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;
|
||||
|
||||
|
|
|
|||
|
|
@ -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<T>(
|
|||
.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;
|
||||
|
|
|
|||
|
|
@ -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<number>`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<string>`CASE json_extract(${event}, '$.type')
|
||||
WHEN 'message' THEN json_set(${entry}, '$.message', json_set(${messageFacts},
|
||||
'$.content', json(${calls}), '$.command', '', '$.output', '',
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
};
|
||||
|
|
|
|||
|
|
@ -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({
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
});
|
||||
|
|
|
|||
|
|
@ -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"),
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue