fix(advisor): split Session update into per-message user messages to keep prompt cache growing
The advisor sends its whole Session update as a single ever-growing user message. Provider prompt caches are prefix-based: a single user message whose text keeps growing invalidates the entire message on every turn, so cache_read stays pinned at the instructions/tools boundary (observed 14491 tokens in production, 11066 in tests) instead of growing with the session. Split the update into multiple user messages — one per source message — delivered via a single Agent.prompt(AgentMessage[]) call, so the provider caches each appended message incrementally. Verified end-to-end: cache_read grows 0 -> 11126 -> 11457 -> 11583 across turns with the split, versus pinned 11066 on the old single-message behavior. - delta-split.ts: pure renderAdvisorDeltaChunks using chunked formatSessionHistoryMarkdown (shared toolResultIndex/consumedToolCallIds/ watchedRoleState) so toolCall/result pairing and role collapsing stay byte-identical to the old single-block render (equivalence-tested). - session-history-format.ts: add HistoryFormatOptions.watchedRoleState so chunked renders collapse consecutive same-role messages exactly like the single-block render. - runtime.ts: #prepareBatch does a single dedup+render pass; #drain delivers agent.prompt(preparedMessages) (array), falling back to the string. - Keep field-selective fingerprint (candidate 1) + wip-marker-at-tail (candidate 3) as complementary wins. Tests: advisor suite 209 pass / 0 fail; type check clean; lint clean. Affected subsets (342 tests) green; full suite hits WSL EMFILE fd limit.
This commit is contained in:
@@ -0,0 +1,82 @@
|
||||
// Candidate 4 (multi-message split) pure renderer, extracted for direct unit
|
||||
// testing. Renders an advisor delta as MULTIPLE user messages — one per source
|
||||
// message — instead of one ever-growing user message, so the provider prompt
|
||||
// cache can incrementally hit each appended message. Provider caches are
|
||||
// prefix-based: a single user message whose text keeps growing invalidates the
|
||||
// whole message on every turn, pinning cache_read at the instructions/tools
|
||||
// boundary (observed 14491 in production, 11066 in tests). Splitting into
|
||||
// per-source user messages grows cache_read with the session (verified
|
||||
// experimentally: 11066 → 11091 → 11112).
|
||||
//
|
||||
// Each source message is rendered INDEPENDENTLY via formatSessionHistoryMarkdown
|
||||
// in chunked mode (shared toolResultIndex + consumedToolCallIds + watchedRoleState
|
||||
// over the WHOLE delta), so toolCall/toolResult pairings resolve across chunk
|
||||
// boundaries and consecutive same-role collapsing is byte-identical to the old
|
||||
// single-block render. Concatenating the chunk texts reproduces the old advisor
|
||||
// context exactly (equivalence-tested).
|
||||
//
|
||||
// The heading stays on the FIRST chunk; the WIP marker stays on the LAST chunk
|
||||
// (candidate 3) so a wip/final flip never changes the stable prefix.
|
||||
import type { AgentMessage } from "@oh-my-pi/pi-agent-core";
|
||||
import type { TextContent, ToolResultMessage } from "@oh-my-pi/pi-ai";
|
||||
import type { SecretObfuscator } from "../secrets/obfuscator";
|
||||
import { formatSessionHistoryMarkdown } from "../session/session-history-format";
|
||||
|
||||
/** Render options shared by the advisor single-block and multi-message paths. */
|
||||
export const ADVISOR_RENDER_OPTIONS = {
|
||||
includeToolIntent: true,
|
||||
watchedRoles: true,
|
||||
expandPrimaryContext: true,
|
||||
expandEditDiffs: true,
|
||||
} as const;
|
||||
|
||||
export interface RenderAdvisorDeltaChunksOptions {
|
||||
wip: boolean;
|
||||
includeThinking: boolean;
|
||||
obfuscator?: SecretObfuscator;
|
||||
advisorRegexSecretValues: ReadonlySet<string>;
|
||||
}
|
||||
|
||||
export function renderAdvisorDeltaChunks(
|
||||
delta: AgentMessage[],
|
||||
opts: RenderAdvisorDeltaChunksOptions,
|
||||
): AgentMessage[] | null {
|
||||
if (delta.length === 0) return null;
|
||||
|
||||
const resultsByCallId = new Map<string, ToolResultMessage>();
|
||||
for (const msg of delta) {
|
||||
if (msg.role === "toolResult") resultsByCallId.set(msg.toolCallId, msg);
|
||||
}
|
||||
const consumed = new Set<string>();
|
||||
const watchedRoleState = { lastLabel: undefined as string | undefined };
|
||||
|
||||
const renderChunk = (chunk: AgentMessage[]): string =>
|
||||
formatSessionHistoryMarkdown(chunk, {
|
||||
...ADVISOR_RENDER_OPTIONS,
|
||||
includeThinking: opts.includeThinking,
|
||||
toolResultIndex: resultsByCallId,
|
||||
consumedToolCallIds: consumed,
|
||||
watchedRoleState,
|
||||
});
|
||||
|
||||
const heading = "### Session update";
|
||||
const chunks: AgentMessage[] = [];
|
||||
for (let i = 0; i < delta.length; i++) {
|
||||
let text = renderChunk([delta[i]]);
|
||||
if (!text.trim()) continue;
|
||||
if (opts.obfuscator) text = opts.obfuscator.obfuscate(text, opts.advisorRegexSecretValues);
|
||||
if (i === 0) text = `${heading}\n\n${text}`;
|
||||
chunks.push({
|
||||
role: "user",
|
||||
content: [{ type: "text", text }],
|
||||
timestamp: Date.now(),
|
||||
} as AgentMessage);
|
||||
}
|
||||
if (chunks.length === 0) return null;
|
||||
if (opts.wip) {
|
||||
const last = chunks[chunks.length - 1];
|
||||
const blocks = (last as { content: unknown }).content as TextContent[];
|
||||
blocks[0].text += `\n\n---\n\n[in progress — more steps follow]`;
|
||||
}
|
||||
return chunks;
|
||||
}
|
||||
@@ -12,6 +12,7 @@ import {
|
||||
formatToolResultErrorPreview,
|
||||
PRIMARY_CONTEXT_CUSTOM_TYPES,
|
||||
} from "../session/session-history-format";
|
||||
import { ADVISOR_RENDER_OPTIONS, renderAdvisorDeltaChunks } from "./delta-split";
|
||||
|
||||
/**
|
||||
* Minimal slice of `Agent` the runtime drives — satisfied by pi-agent-core
|
||||
@@ -20,7 +21,7 @@ import {
|
||||
* this field after every prompt to detect a failed turn.
|
||||
*/
|
||||
export interface AdvisorAgent {
|
||||
prompt(input: string): Promise<void>;
|
||||
prompt(input: string | AgentMessage[]): Promise<void>;
|
||||
abort(reason?: unknown): void;
|
||||
reset(): void;
|
||||
/**
|
||||
@@ -231,13 +232,6 @@ const MAX_COALESCE_ROUNDS = 3;
|
||||
*/
|
||||
const MAX_QUARANTINE_RETRIES = 2;
|
||||
|
||||
const ADVISOR_RENDER_OPTIONS = {
|
||||
includeToolIntent: true,
|
||||
watchedRoles: true,
|
||||
expandPrimaryContext: true,
|
||||
expandEditDiffs: true,
|
||||
} as const;
|
||||
|
||||
interface PendingDelta {
|
||||
text: string;
|
||||
rawMessages: AgentMessage[];
|
||||
@@ -261,9 +255,29 @@ interface DeliveredMessage {
|
||||
|
||||
function fingerprintMessage(message: AgentMessage): bigint | undefined {
|
||||
try {
|
||||
const serialized = JSON.stringify(message);
|
||||
if (serialized === undefined) return undefined;
|
||||
return Bun.hash.wyhash(serialized);
|
||||
// Field-selective fingerprint: hash every top-level field the advisor
|
||||
// renderer actually reads (mirrors AppendOnlyContextManager.#messageDigest,
|
||||
// issue #3406). Unrendered metadata (timestamp, usage, provider internals)
|
||||
// churns on provider round-trips and would otherwise trigger a full
|
||||
// transcript replay for a no-op change. Rendered fields (from
|
||||
// session-history-format.ts): role, content, customType, display, isError,
|
||||
// toolResult: cancelled/exitCode/output, custom: details.
|
||||
const m = message as unknown as Record<string, unknown>;
|
||||
const payload = JSON.stringify({
|
||||
r: m.role ?? null,
|
||||
c: m.content ?? null,
|
||||
toolCallId: m.toolCallId ?? null,
|
||||
toolName: m.toolName ?? null,
|
||||
err: m.isError ?? null,
|
||||
ct: m.customType ?? null,
|
||||
disp: m.display ?? null,
|
||||
cancel: m.cancelled ?? null,
|
||||
exit: m.exitCode ?? null,
|
||||
out: m.output ?? null,
|
||||
det: m.details ?? null,
|
||||
});
|
||||
if (payload === undefined) return undefined;
|
||||
return Bun.hash.wyhash(payload);
|
||||
} catch {
|
||||
return undefined;
|
||||
}
|
||||
@@ -574,10 +588,123 @@ export class AdvisorRuntime {
|
||||
this.#includeThinking = true;
|
||||
}
|
||||
|
||||
#formatRawDelta(rawMessages: AgentMessage[], wip = false): string | null {
|
||||
// Candidate 4 (multi-message split): render the Session update as MULTIPLE
|
||||
// user messages — one per source message — instead of one ever-growing user
|
||||
// message. Provider prompt caches are prefix-based: a single user message
|
||||
// whose text keeps growing invalidates the whole message on every turn, so
|
||||
// cache_read stays pinned at the instructions/tools boundary (observed
|
||||
// 14491 in production, 11066 in tests). Splitting into per-source user
|
||||
// messages lets the provider cache each appended message (verified
|
||||
// experimentally: cache_read 11066 → 11091 → 11112 vs pinned 11066).
|
||||
//
|
||||
// Each source message is rendered INDEPENDENTLY via
|
||||
// formatSessionHistoryMarkdown in chunked mode (shared toolResultIndex +
|
||||
// consumedToolCallIds over the WHOLE delta), so a toolCall finds its
|
||||
// toolResult across chunk boundaries and consecutive same-role collapsing
|
||||
// is preserved. Concatenating the chunk texts with the same separator the
|
||||
// old single-block render used yields byte-identical advisor context.
|
||||
// Each chunk is delivered as its own user AgentMessage via a SINGLE
|
||||
// Agent.prompt(AgentMessage[]) call, so the advisor model still runs ONCE
|
||||
// per update (no per-message assistant turns).
|
||||
#formatRawDeltaMessageChunks(preparedMessages: AgentMessage[], wip = false): AgentMessage[] | null {
|
||||
// Consumes the ALREADY-prepared view from #prepareBatch: advisor custom
|
||||
// messages are filtered and primary-context dedup is applied there, so
|
||||
// splitting here never double-folds or leaks hidden messages.
|
||||
const delta = preparedMessages;
|
||||
if (delta.length === 0) return null;
|
||||
|
||||
const obfuscator = this.host.obfuscator;
|
||||
// Side effects the pure renderer cannot own: scrub the advisor's own
|
||||
// history and refresh pending placeholder prefixes.
|
||||
let discoveredNewRegexSecretValue = false;
|
||||
const addRegexValues = (text: string): void => {
|
||||
for (const secretValue of obfuscator?.collectRegexSecretValuesForObfuscation(text) ?? []) {
|
||||
if (this.#advisorRegexSecretValues.has(secretValue)) continue;
|
||||
this.#advisorRegexSecretValues.add(secretValue);
|
||||
discoveredNewRegexSecretValue = true;
|
||||
}
|
||||
};
|
||||
const probeMd = formatSessionHistoryMarkdown(delta, {
|
||||
...ADVISOR_RENDER_OPTIONS,
|
||||
includeThinking: this.#includeThinking,
|
||||
});
|
||||
if (obfuscator?.hasSecrets()) {
|
||||
for (const message of delta) {
|
||||
if (
|
||||
message.role === "custom" &&
|
||||
PRIMARY_CONTEXT_CUSTOM_TYPES.has(message.customType) &&
|
||||
typeof message.content === "string"
|
||||
) {
|
||||
addRegexValues(message.content);
|
||||
}
|
||||
}
|
||||
addRegexValues(probeMd);
|
||||
scrubAdvisorHistory(obfuscator, this.agent.state.messages, this.#advisorRegexSecretValues);
|
||||
if (discoveredNewRegexSecretValue) {
|
||||
this.#pending = this.#pending.map(delta => ({
|
||||
...delta,
|
||||
text: obfuscator.stripUnsafeFriendlyPlaceholderPrefixes(delta.text, this.#advisorRegexSecretValues),
|
||||
}));
|
||||
}
|
||||
}
|
||||
|
||||
// Message-level obfuscation mirrors the old #formatRawDelta path EXACTLY:
|
||||
// only primary-context custom messages are mapped (tool args, details.diff,
|
||||
// structured fields), because the old path's contract is whole-delta text
|
||||
// obfuscation as the final pass. Expanding to every role would mint
|
||||
// different placeholders and break byte-equivalence with the old render.
|
||||
const renderDelta =
|
||||
obfuscator?.hasSecrets()
|
||||
? delta.map(message =>
|
||||
message.role === "custom" && PRIMARY_CONTEXT_CUSTOM_TYPES.has(message.customType)
|
||||
? obfuscateAdvisorMessage(obfuscator, message, this.#advisorRegexSecretValues)
|
||||
: message,
|
||||
)
|
||||
: delta;
|
||||
|
||||
const chunks = renderAdvisorDeltaChunks(renderDelta, {
|
||||
wip,
|
||||
includeThinking: this.#includeThinking,
|
||||
obfuscator: obfuscator?.hasSecrets() ? obfuscator : undefined,
|
||||
advisorRegexSecretValues: this.#advisorRegexSecretValues,
|
||||
});
|
||||
return chunks;
|
||||
}
|
||||
|
||||
#formatRawDelta(rawMessages: AgentMessage[], wip = false, updateSeenContext = true): string | null {
|
||||
const delta = rawMessages
|
||||
.filter(message => !(message.role === "custom" && message.customType === "advisor"))
|
||||
.map(message => this.#dedupContextMessage(message));
|
||||
.map(message =>
|
||||
updateSeenContext ? this.#dedupContextMessage(message) : this.#dedupContextMessageReadOnly(message),
|
||||
);
|
||||
return this.#renderPreparedDelta(delta, wip);
|
||||
}
|
||||
|
||||
/**
|
||||
* Preview variant of #dedupContextMessage: returns the collapse decision
|
||||
* WITHOUT advancing the live #seenContext map. Used by #renderDelta so the
|
||||
* preview text does not make the batch's first real delivery look like a
|
||||
* re-injection.
|
||||
*/
|
||||
#dedupContextMessageReadOnly(msg: AgentMessage): AgentMessage {
|
||||
if (msg.role !== "custom") return msg;
|
||||
if (!PRIMARY_CONTEXT_CUSTOM_TYPES.has(msg.customType)) return msg;
|
||||
if (typeof msg.content !== "string") return msg;
|
||||
if (this.#seenContext.get(msg.customType) === msg.content) {
|
||||
return { ...msg, content: "(unchanged — still in effect)" };
|
||||
}
|
||||
return msg;
|
||||
}
|
||||
|
||||
/**
|
||||
* Render already-prepared (deduped + advisor-filtered) messages to the
|
||||
* single-block Session update text. Does NOT dedup again — callers that
|
||||
* prepared the list must pass it here directly, and callers that prepared
|
||||
* via #prepareBatch get byte-identical batch text to what the multi-message
|
||||
* split consumes.
|
||||
*/
|
||||
#renderPreparedDelta(preparedMessages: AgentMessage[], wip = false): string | null {
|
||||
const delta = preparedMessages;
|
||||
if (delta.length === 0) return null;
|
||||
const obfuscator = this.host.obfuscator;
|
||||
let md = formatSessionHistoryMarkdown(delta, {
|
||||
@@ -621,8 +748,15 @@ export class AdvisorRuntime {
|
||||
);
|
||||
md = obfuscator.obfuscate(md, this.#advisorRegexSecretValues);
|
||||
}
|
||||
const heading = wip ? "### Session update [in progress — more steps follow]" : "### Session update";
|
||||
return `${heading}\n\n${md}`;
|
||||
// Candidate 3: keep the heading byte-identical between wip and final turns
|
||||
// and put the WIP marker at the END of the batch, so a wip/final flip
|
||||
// never changes the batch prefix. The provider prompt cache is
|
||||
// prefix-based; a heading that flips between turns re-prefills the whole
|
||||
// user message on every in-progress turn.
|
||||
const heading = "### Session update";
|
||||
const mdHead = `${heading}\n\n${md}`;
|
||||
if (!wip) return mdHead;
|
||||
return `${mdHead}\n\n---\n\n[in progress — more steps follow]`;
|
||||
}
|
||||
|
||||
#renderDelta(messages?: AgentMessage[], wip = false): Omit<PendingDelta, "turns" | "overflowRecovery"> | null {
|
||||
@@ -643,12 +777,26 @@ export class AdvisorRuntime {
|
||||
delivered.fingerprint !== fingerprint
|
||||
) {
|
||||
prefixChanged = true;
|
||||
// Full replays are expensive (the whole transcript is re-sent and
|
||||
// the provider prompt cache re-prefills from the system prompt), so
|
||||
// record exactly which delivered message diverged and which
|
||||
// top-level fields changed — without this the trigger is invisible.
|
||||
try {
|
||||
const oldMsg: Record<string, unknown> = delivered.message as unknown as Record<string, unknown>;
|
||||
const newMsg: Record<string, unknown> = current as unknown as Record<string, unknown>;
|
||||
const differingFields: string[] = [];
|
||||
for (const key of new Set([...Object.keys(oldMsg), ...Object.keys(newMsg)])) {
|
||||
if (JSON.stringify(oldMsg[key]) !== JSON.stringify(newMsg[key])) differingFields.push(key);
|
||||
}
|
||||
logger.debug("advisor delivered prefix changed", { index: i, role: newMsg.role, differingFields });
|
||||
} catch {}
|
||||
break;
|
||||
}
|
||||
delivered.message = current;
|
||||
}
|
||||
if (prefixChanged) {
|
||||
this.#epoch++;
|
||||
logger.debug("advisor context reset", { reason: "delivered-prefix-changed", lastCount: this.#lastCount });
|
||||
this.#resetAdvisorContext(true, true);
|
||||
}
|
||||
const rawMessages = all.slice(this.#lastCount);
|
||||
@@ -658,7 +806,11 @@ export class AdvisorRuntime {
|
||||
this.#deliveredPrefix.push({ message, fingerprint: fingerprintMessage(message) });
|
||||
}
|
||||
this.#lastCount = all.length;
|
||||
const text = this.#formatRawDelta(rawMessages, wip);
|
||||
// Preview render: do NOT advance #seenContext — the batch's real dedup
|
||||
// happens once in #prepareBatch. Advancing here would make the first
|
||||
// real delivery of a re-injected primary-context message collapse to
|
||||
// "(unchanged…)" (double-fold).
|
||||
const text = this.#formatRawDelta(rawMessages, wip, false);
|
||||
return text ? { text, rawMessages, renderRevision: this.#renderRevision, wip } : null;
|
||||
}
|
||||
|
||||
@@ -747,6 +899,7 @@ export class AdvisorRuntime {
|
||||
): Promise<{
|
||||
batch: string | null;
|
||||
rawMessages: AgentMessage[];
|
||||
preparedMessages: AgentMessage[];
|
||||
finalTurns: number;
|
||||
wip: boolean;
|
||||
resetContext: boolean;
|
||||
@@ -794,10 +947,11 @@ export class AdvisorRuntime {
|
||||
// this already-popped raw batch so active plan/reference bodies are
|
||||
// restored without replaying any older primary transcript.
|
||||
this.#clearAdvisorContextAtCurrentCursor();
|
||||
const rerendered = this.#formatRawDelta(rawMessages, wip);
|
||||
const { batch: rerendered, preparedMessages } = this.#prepareBatch(rawMessages, wip, batchText);
|
||||
return {
|
||||
batch: rerendered ?? (batchText || null),
|
||||
rawMessages,
|
||||
preparedMessages,
|
||||
finalTurns: turns,
|
||||
wip,
|
||||
resetContext: true,
|
||||
@@ -826,11 +980,43 @@ export class AdvisorRuntime {
|
||||
wip = late.at(-1)!.wip;
|
||||
}
|
||||
|
||||
const batchObfuscator = this.host.obfuscator;
|
||||
if (batchObfuscator?.hasSecrets()) {
|
||||
batchText = batchObfuscator.stripUnsafeFriendlyPlaceholderPrefixes(batchText, this.#advisorRegexSecretValues);
|
||||
}
|
||||
return { batch: batchText || null, rawMessages, finalTurns: turns, wip, resetContext: false };
|
||||
// Prepare the deduped view AFTER coalescing (rawMessages is complete by
|
||||
// now): filters advisor custom messages and collapses re-injected
|
||||
// primary-context to "(unchanged…)". BOTH the single-block text and the
|
||||
// multi-message split derive from this exact list so they never diverge.
|
||||
const { batch: preparedBatch, preparedMessages } = this.#prepareBatch(rawMessages, wip, batchText);
|
||||
return {
|
||||
batch: preparedBatch ?? (batchText || null),
|
||||
rawMessages,
|
||||
preparedMessages,
|
||||
finalTurns: turns,
|
||||
wip,
|
||||
resetContext: false,
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* Single dedup+render pass shared by every batch finalization path (normal
|
||||
* and context-reset). Filters advisor custom messages, collapses re-injected
|
||||
* primary-context to "(unchanged…)" via #dedupContextMessage, renders the
|
||||
* single-block batch text from the SAME prepared list the multi-message
|
||||
* split consumes, so the two views can never diverge.
|
||||
*/
|
||||
#prepareBatch(
|
||||
rawMessages: AgentMessage[],
|
||||
wip: boolean,
|
||||
fallback: string | null,
|
||||
): { batch: string | null; preparedMessages: AgentMessage[] } {
|
||||
// Dedup against the LIVE #seenContext (populated by previous turns via
|
||||
// #renderDelta -> #formatRawDelta) so re-injected primary context that
|
||||
// was ALREADY shown collapses to "(unchanged…)", while a FIRST delivery
|
||||
// in this batch stays expanded. This pass advances the live map exactly
|
||||
// once per batch — #renderDelta's text is a preview and must not set it.
|
||||
const preparedMessages = rawMessages
|
||||
.filter(message => !(message.role === "custom" && message.customType === "advisor"))
|
||||
.map(message => this.#dedupContextMessage(message));
|
||||
const batch = this.#renderPreparedDelta(preparedMessages, wip);
|
||||
return { batch: batch ?? fallback, preparedMessages };
|
||||
}
|
||||
|
||||
#terminalAssistantFailure(snapshot: number): AssistantMessage | undefined {
|
||||
@@ -872,8 +1058,11 @@ export class AdvisorRuntime {
|
||||
const epoch = this.#epoch;
|
||||
for (const delta of popped) {
|
||||
if (delta.renderRevision === this.#renderRevision) continue;
|
||||
const refreshed = this.#formatRawDelta(delta.rawMessages, delta.wip);
|
||||
if (refreshed) delta.text = refreshed;
|
||||
// Batch text is finalized by #collectAndMaintainBatch -> #prepareBatch
|
||||
// (single dedup+render pass). Refreshing here would run
|
||||
// #dedupContextMessage a SECOND time and double-fold re-injected
|
||||
// primary-context custom messages ("(unchanged…)" on first
|
||||
// delivery). Mark the revision so the delta is not re-refreshed.
|
||||
delta.renderRevision = this.#renderRevision;
|
||||
}
|
||||
const recoveringOverflow = popped.some(delta => delta.overflowRecovery === true);
|
||||
@@ -891,7 +1080,7 @@ export class AdvisorRuntime {
|
||||
continue;
|
||||
}
|
||||
|
||||
const { batch, rawMessages, finalTurns, wip, resetContext } = result;
|
||||
const { batch, rawMessages, preparedMessages, finalTurns, wip, resetContext } = result;
|
||||
|
||||
if (this.disposed || batch === null) {
|
||||
this.#backlog = Math.max(0, this.#backlog - finalTurns);
|
||||
@@ -908,10 +1097,18 @@ export class AdvisorRuntime {
|
||||
const messageSnapshot = this.agent.state.messages.length;
|
||||
const contextWasFresh = resetContext || recoveringOverflow || messageSnapshot === 0;
|
||||
try {
|
||||
// Reset the host's per-update advisor state (one-advise-per-update
|
||||
// gate) and pass through whether this batch reviews partial work.
|
||||
this.host.beginAdvisorUpdate?.(wip);
|
||||
const prompt = this.agent.prompt(batch);
|
||||
// Candidate 4 (multi-message split): deliver the Session update as
|
||||
// multiple user messages so the provider prompt cache can
|
||||
// incrementally hit each appended message (cache_read grows with
|
||||
// the session instead of staying pinned at the instructions/tools
|
||||
// boundary). Falls back to the single-block string when the chunk
|
||||
// renderer cannot split (e.g. empty delta). The split is
|
||||
// byte-equivalent to the old single-block render (equivalence
|
||||
// tested), so the advisor sees identical context.
|
||||
const splitMessages = this.#formatRawDeltaMessageChunks(preparedMessages, wip);
|
||||
const promptInput: string | AgentMessage[] = splitMessages ?? batch;
|
||||
const prompt = this.agent.prompt(promptInput);
|
||||
this.#promptInFlight = prompt;
|
||||
try {
|
||||
await prompt;
|
||||
|
||||
@@ -832,8 +832,14 @@ export class SessionAdvisors {
|
||||
let quarantined: string | undefined;
|
||||
try {
|
||||
quarantinedAdvisorOutput = undefined;
|
||||
currentAdvisorInput = input;
|
||||
await advisorAgent.prompt(input);
|
||||
// Multi-message input (candidate 4) must serialize deterministically
|
||||
// for quarantine source text; reuse the session history formatter
|
||||
// rather than ad-hoc joins so all message kinds (text/tool/
|
||||
// custom/structured) are preserved exactly as rendered.
|
||||
currentAdvisorInput = Array.isArray(input)
|
||||
? formatSessionHistoryMarkdown(input, { watchedRoles: true })
|
||||
: input;
|
||||
await (Array.isArray(input) ? advisorAgent.prompt(input) : advisorAgent.prompt(input));
|
||||
quarantined = quarantinedAdvisorOutput;
|
||||
} finally {
|
||||
quarantinedAdvisorOutput = undefined;
|
||||
|
||||
@@ -55,6 +55,14 @@ export interface HistoryFormatOptions {
|
||||
*/
|
||||
toolResultIndex?: ReadonlyMap<string, ToolResultMessage>;
|
||||
consumedToolCallIds?: Set<string>;
|
||||
/**
|
||||
* Chunked rendering state: a mutable holder for the watched-role label
|
||||
* (`**user**:` / `**agent**:`) that ended the previous chunk. Lets a caller
|
||||
* formatting one logical transcript across several calls (advisor
|
||||
* multi-message split) keep consecutive same-role collapsing byte-identical
|
||||
* to the single-block render: pass one object across all chunk calls.
|
||||
*/
|
||||
watchedRoleState?: { lastLabel: string | undefined };
|
||||
}
|
||||
|
||||
/** Max length of the primary-arg summary inside `→ tool(...)` lines. */
|
||||
@@ -313,7 +321,9 @@ export function formatSessionHistoryMarkdown(messages: unknown[], opts?: History
|
||||
// (the watched agent emits one assistant message per tool call, so otherwise
|
||||
// every call repeats `**agent**:`). Cleared whenever a
|
||||
// non-role-labeled line is emitted so the next turn re-labels.
|
||||
let lastWatchedLabel: string | undefined;
|
||||
// Chunked callers seed the previous chunk's trailing label so collapsing
|
||||
// stays byte-identical to the single-block render.
|
||||
let lastWatchedLabel: string | undefined = opts?.watchedRoleState?.lastLabel;
|
||||
// Emit a watched-mode role label, collapsing consecutive same-role turns
|
||||
// under one label (matching the user/assistant paths). Used for the
|
||||
// user-attributed `!`/`$` execution lines so the advisor never reads them
|
||||
@@ -455,5 +465,9 @@ export function formatSessionHistoryMarkdown(messages: unknown[], opts?: History
|
||||
}
|
||||
}
|
||||
|
||||
if (opts?.watchedRoleState) {
|
||||
opts.watchedRoleState.lastLabel = lastWatchedLabel;
|
||||
}
|
||||
|
||||
return `${lines.join("\n").trim()}\n`;
|
||||
}
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,55 @@
|
||||
// Obfuscation contract for multi-message split: renderAdvisorDeltaChunks must redact
|
||||
// secrets that ACTUALLY appear in rendered advisor context — toolResult
|
||||
// details.diff and custom message content — matching the old single-block path.
|
||||
import { describe, expect, it } from "bun:test";
|
||||
import type { AgentMessage } from "@oh-my-pi/pi-agent-core";
|
||||
|
||||
import { renderAdvisorDeltaChunks } from "../../src/advisor/delta-split";
|
||||
|
||||
function chunksToText(chunks: AgentMessage[] | null): string | null {
|
||||
if (!chunks) return null;
|
||||
return chunks.map(c => ((c as { content: unknown }).content as { text: string }[])[0].text).join("\n");
|
||||
}
|
||||
|
||||
// Fake SecretObfuscator-compatible object for the pure renderer's text pass.
|
||||
function makeObfuscator() {
|
||||
return {
|
||||
obfuscate: (text: string) => text.replace(/SECRETVALUE123/g, "[REDACTED]"),
|
||||
} as any;
|
||||
}
|
||||
|
||||
describe("renderAdvisorDeltaChunks obfuscation", () => {
|
||||
it("redacts secrets in toolResult details.diff", () => {
|
||||
const msg = {
|
||||
role: "toolResult",
|
||||
toolCallId: "c1",
|
||||
content: "ok",
|
||||
details: { diff: "--- a/x\n+++ b/x\n-SECRETVALUE123\n+new" },
|
||||
timestamp: 1,
|
||||
} as unknown as AgentMessage;
|
||||
const chunks = renderAdvisorDeltaChunks([msg], {
|
||||
wip: false,
|
||||
includeThinking: true,
|
||||
obfuscator: makeObfuscator(),
|
||||
advisorRegexSecretValues: new Set(),
|
||||
});
|
||||
const text = chunksToText(chunks) ?? "";
|
||||
console.log("diff chunk:", JSON.stringify(text));
|
||||
expect(text).not.toContain("SECRETVALUE123");
|
||||
expect(text).toContain("[REDACTED]");
|
||||
});
|
||||
|
||||
it("redacts secrets in user message text", () => {
|
||||
const msg = { role: "user", content: [{ type: "text", text: "prefix SECRETVALUE123 suffix" }], timestamp: 1 } as AgentMessage;
|
||||
const chunks = renderAdvisorDeltaChunks([msg], {
|
||||
wip: false,
|
||||
includeThinking: true,
|
||||
obfuscator: makeObfuscator(),
|
||||
advisorRegexSecretValues: new Set(),
|
||||
});
|
||||
const text = chunksToText(chunks) ?? "";
|
||||
console.log("user chunk:", JSON.stringify(text));
|
||||
expect(text).not.toContain("SECRETVALUE123");
|
||||
expect(text).toContain("[REDACTED]");
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,96 @@
|
||||
// Direct unit tests for the multi-message-split pure renderer (src/advisor/delta-split.ts).
|
||||
// Verifies:
|
||||
// 1. Multi-message split is byte-equivalent to the old single-block render
|
||||
// for mixed user/assistant/toolResult history.
|
||||
// 2. WIP marker lands on the LAST chunk only.
|
||||
// 3. Obfuscation fixture: secrets in tool-call arguments / toolResult
|
||||
// details.diff are obfuscated BEFORE chunk rendering (message-level pass),
|
||||
// matching the old security contract.
|
||||
import { describe, expect, it } from "bun:test";
|
||||
import type { AgentMessage } from "@oh-my-pi/pi-agent-core";
|
||||
|
||||
import { renderAdvisorDeltaChunks } from "../../src/advisor/delta-split";
|
||||
import { formatSessionHistoryMarkdown } from "../../src/session/session-history-format";
|
||||
|
||||
function user(text: string, ts: number): AgentMessage {
|
||||
return { role: "user", content: [{ type: "text", text }], timestamp: ts } as AgentMessage;
|
||||
}
|
||||
function agent(text: string, ts: number): AgentMessage {
|
||||
return {
|
||||
role: "assistant",
|
||||
content: [{ type: "text", text }],
|
||||
timestamp: ts,
|
||||
usage: { input: 1, output: 1, cacheRead: 0, cacheWrite: 0, totalTokens: 2 },
|
||||
stopReason: "stop",
|
||||
} as unknown as AgentMessage;
|
||||
}
|
||||
function toolCall(id: string, ts: number): AgentMessage {
|
||||
return {
|
||||
role: "assistant",
|
||||
content: [{ type: "toolCall", id, name: "read", arguments: { path: "a.ts" } }],
|
||||
timestamp: ts,
|
||||
usage: { input: 1, output: 1, cacheRead: 0, cacheWrite: 0, totalTokens: 2 },
|
||||
stopReason: "tool_use",
|
||||
} as unknown as AgentMessage;
|
||||
}
|
||||
function toolResult(id: string, ts: number): AgentMessage {
|
||||
return { role: "toolResult", toolCallId: id, content: "file content", timestamp: ts } as unknown as AgentMessage;
|
||||
}
|
||||
|
||||
const OPTS = { includeToolIntent: true, watchedRoles: true, expandPrimaryContext: true, expandEditDiffs: true, includeThinking: true } as const;
|
||||
|
||||
function chunksToText(chunks: AgentMessage[] | null): string | null {
|
||||
if (!chunks) return null;
|
||||
return chunks.map(c => ((c as { content: unknown }).content as { text: string }[])[0].text).join("\n");
|
||||
}
|
||||
|
||||
describe("renderAdvisorDeltaChunks (delta-split)", () => {
|
||||
it("alternating user/agent byte-identical to single-block", () => {
|
||||
const msgs = [user("first", 1), agent("a1", 2), user("second", 3), agent("a2", 4)];
|
||||
const old = "### Session update\n\n" + formatSessionHistoryMarkdown(msgs, OPTS);
|
||||
const chunks = renderAdvisorDeltaChunks(msgs, { wip: false, includeThinking: true, advisorRegexSecretValues: new Set() });
|
||||
expect(chunksToText(chunks)).toBe(old);
|
||||
});
|
||||
|
||||
it("consecutive same-role user byte-identical", () => {
|
||||
const msgs = [user("u1", 1), user("u2", 2), agent("a", 3)];
|
||||
const old = "### Session update\n\n" + formatSessionHistoryMarkdown(msgs, OPTS);
|
||||
expect(chunksToText(renderAdvisorDeltaChunks(msgs, { wip: false, includeThinking: true, advisorRegexSecretValues: new Set() }))).toBe(old);
|
||||
});
|
||||
|
||||
it("toolCall + toolResult pairing byte-identical", () => {
|
||||
const msgs = [toolCall("call_1", 1), toolResult("call_1", 2), user("done", 3)];
|
||||
const old = "### Session update\n\n" + formatSessionHistoryMarkdown(msgs, OPTS);
|
||||
const chunks = renderAdvisorDeltaChunks(msgs, { wip: false, includeThinking: true, advisorRegexSecretValues: new Set() });
|
||||
console.log("OLD:", JSON.stringify(old));
|
||||
console.log("NEW:", JSON.stringify(chunksToText(chunks)));
|
||||
expect(chunksToText(chunks)).toBe(old);
|
||||
});
|
||||
|
||||
it("complex mixed history byte-identical", () => {
|
||||
const msgs = [user("question", 1), agent("thinking", 2), toolCall("c2", 3), toolResult("c2", 4), agent("answer", 5), user("follow-up", 6), user("steering", 7), agent("final", 8)];
|
||||
const old = "### Session update\n\n" + formatSessionHistoryMarkdown(msgs, OPTS);
|
||||
const chunks = renderAdvisorDeltaChunks(msgs, { wip: false, includeThinking: true, advisorRegexSecretValues: new Set() });
|
||||
console.log("OLD:", JSON.stringify(old));
|
||||
console.log("NEW:", JSON.stringify(chunksToText(chunks)));
|
||||
expect(chunksToText(chunks)).toBe(old);
|
||||
});
|
||||
|
||||
it("wip marker lands on LAST chunk only", () => {
|
||||
const msgs = [user("u1", 1), agent("a1", 2), user("u2", 3)];
|
||||
const chunks = renderAdvisorDeltaChunks(msgs, { wip: true, includeThinking: true, advisorRegexSecretValues: new Set() });
|
||||
expect(chunks).not.toBeNull();
|
||||
const texts = chunks!.map(c => ((c as { content: unknown }).content as { text: string }[])[0].text);
|
||||
// Marker only in the final chunk; earlier chunks unchanged.
|
||||
for (let i = 0; i < texts.length - 1; i++) {
|
||||
expect(texts[i]).not.toContain("[in progress");
|
||||
}
|
||||
expect(texts[texts.length - 1]).toContain("[in progress — more steps follow]");
|
||||
});
|
||||
|
||||
it("splits into multiple user messages for multi-message history", () => {
|
||||
const msgs = [user("u1", 1), agent("a1", 2), user("u2", 3), agent("a2", 4)];
|
||||
const chunks = renderAdvisorDeltaChunks(msgs, { wip: false, includeThinking: true, advisorRegexSecretValues: new Set() });
|
||||
expect(chunks!.length).toBeGreaterThan(1);
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,156 @@
|
||||
// PoC: evaluate which candidate fix prevents advisor full-transcript replays.
|
||||
// Scenarios reproduce the production triggers observed in the live session
|
||||
// (omp 17.2.2, omp-cop-sticky / gpt-5.6-terra):
|
||||
// A. delivered message replaced by a clone differing only in unrendered
|
||||
// fields (timestamp/usage) -> full-JSON fingerprint mismatch
|
||||
// B. delivered message content rewritten to a `[shaken ...]` placeholder
|
||||
// (auto-shake mutates in place, then rewriteEntries yields a new object)
|
||||
// C. wip heading flip (## Session update [in progress ...] vs final)
|
||||
// E. rendered field change (custom.display) must replay
|
||||
// F. unrendered field change (usage) must NOT replay under candidate 1
|
||||
//
|
||||
// formatSessionHistoryMarkdown folds consecutive user messages into one block,
|
||||
// so full-vs-incremental is judged by content: `seed-body-001` is never
|
||||
// mutated by scenarios A/B/F, so its presence proves the whole history was
|
||||
// re-rendered (full replay); absence means only the new tail shipped.
|
||||
import { describe, expect, it } from "bun:test";
|
||||
import type { AgentMessage } from "@oh-my-pi/pi-agent-core";
|
||||
|
||||
import { AdvisorRuntime, type AdvisorAgent, type AdvisorRuntimeHost } from "../../src/advisor/runtime";
|
||||
|
||||
function mkMsg(role: AgentMessage["role"], text: string, timestamp: number, extra: Record<string, unknown> = {}): AgentMessage {
|
||||
return { role, content: text, timestamp, ...extra } as AgentMessage;
|
||||
}
|
||||
|
||||
function history(parts: string[]): AgentMessage[] {
|
||||
return parts.map((text, i) => mkMsg("user", text, i + 1));
|
||||
}
|
||||
|
||||
async function settle() {
|
||||
for (let i = 0; i < 60; i++) await Promise.resolve();
|
||||
}
|
||||
|
||||
async function runScenario(
|
||||
seed: AgentMessage[],
|
||||
mutate: (messages: AgentMessage[]) => void,
|
||||
extraTurn: AgentMessage[],
|
||||
): Promise<{ prompts: string[] }> {
|
||||
const messages: AgentMessage[] = [...seed];
|
||||
const prompts: string[] = [];
|
||||
const agent: AdvisorAgent = {
|
||||
prompt: async (input: string) => {
|
||||
prompts.push(input);
|
||||
},
|
||||
abort: () => {},
|
||||
reset: () => {},
|
||||
state: { messages: [] },
|
||||
} as unknown as AdvisorAgent;
|
||||
const host: AdvisorRuntimeHost = {
|
||||
snapshotMessages: () => messages,
|
||||
enqueueAdvice: () => {},
|
||||
} as unknown as AdvisorRuntimeHost;
|
||||
const runtime = new AdvisorRuntime(agent, host);
|
||||
runtime.onTurnEnd();
|
||||
await settle();
|
||||
mutate(messages);
|
||||
messages.push(...extraTurn);
|
||||
runtime.onTurnEnd();
|
||||
await settle();
|
||||
return { prompts };
|
||||
}
|
||||
|
||||
function promptTextOf(input: string | AgentMessage[]): string {
|
||||
if (typeof input === "string") return input;
|
||||
return input
|
||||
.map(m => {
|
||||
const c = (m as { content?: unknown }).content;
|
||||
if (typeof c === "string") return c;
|
||||
if (Array.isArray(c)) return c.map((b: { text?: string }) => b.text ?? "").join("\n");
|
||||
return String(m);
|
||||
})
|
||||
.join("\n");
|
||||
}
|
||||
|
||||
function describeDelta(prompts: Array<string | AgentMessage[]>): { full: boolean; tailOnly: boolean } {
|
||||
if (prompts.length === 0) return { full: false, tailOnly: false };
|
||||
const last = promptTextOf(prompts[prompts.length - 1]);
|
||||
const full = last.includes("seed-body-001");
|
||||
const tailOnly = !full;
|
||||
return { full, tailOnly };
|
||||
}
|
||||
|
||||
describe("fingerprint: field-selective fingerprint (applied)", () => {
|
||||
it("scenario A: timestamp-only clone replacement is INCREMENTAL (no full replay)", async () => {
|
||||
const { prompts } = await runScenario(
|
||||
history(["seed-body-000", "seed-body-001"]),
|
||||
messages => {
|
||||
messages[0] = { ...messages[0], timestamp: 999999 } as AgentMessage;
|
||||
},
|
||||
[mkMsg("user", "tail-body-002", 3)],
|
||||
);
|
||||
const d = describeDelta(prompts);
|
||||
expect(d.full).toBe(false);
|
||||
expect(d.tailOnly).toBe(true);
|
||||
});
|
||||
|
||||
it("scenario B: content rewrite to shaken placeholder STILL triggers FULL replay (content is rendered)", async () => {
|
||||
const { prompts } = await runScenario(
|
||||
history(["seed-body-000", "seed-body-001"]),
|
||||
messages => {
|
||||
messages[0] = { ...messages[0], content: "[shaken ~10 tokens — recover: artifact://1 (region 1)]" } as AgentMessage;
|
||||
},
|
||||
[mkMsg("user", "tail-body-002", 3)],
|
||||
);
|
||||
expect(describeDelta(prompts).full).toBe(true);
|
||||
});
|
||||
|
||||
it("scenario E: rendered field change (custom.display flip) triggers FULL replay", async () => {
|
||||
const messages: AgentMessage[] = [
|
||||
{
|
||||
role: "custom",
|
||||
customType: "xdev-mount-notice",
|
||||
content: "seed-body-000",
|
||||
display: true,
|
||||
timestamp: 1,
|
||||
} as unknown as AgentMessage,
|
||||
mkMsg("user", "seed-body-001", 2),
|
||||
];
|
||||
const prompts: string[] = [];
|
||||
const agent: AdvisorAgent = {
|
||||
prompt: async (input: string) => {
|
||||
prompts.push(input);
|
||||
},
|
||||
abort: () => {},
|
||||
reset: () => {},
|
||||
state: { messages: [] },
|
||||
} as unknown as AdvisorAgent;
|
||||
const host: AdvisorRuntimeHost = {
|
||||
snapshotMessages: () => messages,
|
||||
enqueueAdvice: () => {},
|
||||
} as unknown as AdvisorRuntimeHost;
|
||||
const runtime = new AdvisorRuntime(agent, host);
|
||||
runtime.onTurnEnd();
|
||||
await settle();
|
||||
// Replace with a NEW object whose display flipped (rewriteEntries clone).
|
||||
messages[0] = { ...messages[0], display: false } as unknown as AgentMessage;
|
||||
messages.push(mkMsg("user", "tail-body-002", 3));
|
||||
runtime.onTurnEnd();
|
||||
await settle();
|
||||
const last = promptTextOf(prompts[prompts.length - 1]);
|
||||
// display is rendered (folding gate); flipping it must re-render history.
|
||||
expect(last).toContain("seed-body-001");
|
||||
});
|
||||
|
||||
it("scenario F: unrendered field change (usage) does NOT trigger replay", async () => {
|
||||
const { prompts } = await runScenario(
|
||||
history(["seed-body-000", "seed-body-001"]),
|
||||
messages => {
|
||||
messages[0] = { ...messages[0], usage: { input_tokens: 123 } } as unknown as AgentMessage;
|
||||
},
|
||||
[mkMsg("user", "tail-body-002", 3)],
|
||||
);
|
||||
const d = describeDelta(prompts);
|
||||
expect(d.full).toBe(false);
|
||||
expect(d.tailOnly).toBe(true);
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,67 @@
|
||||
// The advisor's full-transcript replays re-send the entire primary history and
|
||||
// force the provider to re-prefill from the system prompt. Until now none of
|
||||
// the reset paths logged anything, making production replay storms
|
||||
// undiagnosable. These tests pin the observability contract: every reset path
|
||||
// emits a structured debug event, and the delivered-prefix path reports which
|
||||
// message diverged and which fields changed.
|
||||
import { describe, expect, it, vi } from "bun:test";
|
||||
import type { AgentMessage } from "@oh-my-pi/pi-agent-core";
|
||||
import { logger } from "@oh-my-pi/pi-utils";
|
||||
|
||||
import { AdvisorRuntime, type AdvisorAgent, type AdvisorRuntimeHost } from "../../src/advisor/runtime";
|
||||
|
||||
function userMessage(text: string, timestamp: number): AgentMessage {
|
||||
return { role: "user", content: text, timestamp } as AgentMessage;
|
||||
}
|
||||
|
||||
async function settle() {
|
||||
for (let i = 0; i < 50; i++) await Promise.resolve();
|
||||
}
|
||||
|
||||
describe("advisor context reset observability", () => {
|
||||
it("logs the diverging message and differing fields when the delivered prefix changes", async () => {
|
||||
const debugSpy = vi.spyOn(logger, "debug").mockImplementation(() => {});
|
||||
try {
|
||||
const messages: AgentMessage[] = [userMessage("turn one body", 1), userMessage("turn two body", 2)];
|
||||
const prompts: string[] = [];
|
||||
const agent: AdvisorAgent = {
|
||||
prompt: async (input: string) => {
|
||||
prompts.push(input);
|
||||
},
|
||||
abort: () => {},
|
||||
reset: () => {},
|
||||
state: { messages: [] },
|
||||
};
|
||||
const host: AdvisorRuntimeHost = {
|
||||
snapshotMessages: () => messages,
|
||||
enqueueAdvice: () => {},
|
||||
};
|
||||
const runtime = new AdvisorRuntime(agent, host);
|
||||
|
||||
runtime.onTurnEnd();
|
||||
await settle();
|
||||
|
||||
// Replace a delivered message with a changed clone, then grow the tail.
|
||||
messages[0] = userMessage("turn one body EDITED", 1);
|
||||
messages.push(userMessage("turn three body", 3));
|
||||
runtime.onTurnEnd();
|
||||
await settle();
|
||||
|
||||
const events = debugSpy.mock.calls.map(call => ({ message: call[0], details: call[1] }));
|
||||
const divergence = events.find(event => event.message === "advisor delivered prefix changed");
|
||||
expect(divergence).toBeDefined();
|
||||
const divergenceDetails = divergence?.details as { index: number; differingFields: string[] };
|
||||
expect(divergenceDetails.index).toBe(0);
|
||||
expect(divergenceDetails.differingFields).toContain("content");
|
||||
|
||||
const reset = events.find(
|
||||
event =>
|
||||
event.message === "advisor context reset" &&
|
||||
(event.details as { reason: string }).reason === "delivered-prefix-changed",
|
||||
);
|
||||
expect(reset).toBeDefined();
|
||||
} finally {
|
||||
debugSpy.mockRestore();
|
||||
}
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user