diff --git a/src/main/claude/claude-structured-item-translation.ts b/src/main/claude/claude-structured-item-translation.ts index 1c1673b59cd..3b43a0c9301 100644 --- a/src/main/claude/claude-structured-item-translation.ts +++ b/src/main/claude/claude-structured-item-translation.ts @@ -6,7 +6,7 @@ import type { import type { NativeChatBlock } from '../../shared/native-chat-types' import { boundInlineText, - DEFAULT_JOURNAL_PAYLOAD_LIMITS + type JournalPayloadLimits } from '../native-chat/agent-session-journal/journal-payload-bounds' export type ClaudeMessageEnvelope = { @@ -171,18 +171,19 @@ export function claudeThinkingText(envelope: ClaudeMessageEnvelope): string | nu return parts.length > 0 ? parts.join('\n') : null } -export function claudeToolBody(input: { - tool: ClaudeToolUse - result?: ClaudeToolResult -}): AgentJournalItemBody { +export function claudeToolBody( + input: { + tool: ClaudeToolUse + result?: ClaudeToolResult + }, + limits: JournalPayloadLimits +): AgentJournalItemBody { return { kind: 'tool-call', name: input.tool.name, input: input.tool.input, state: input.result ? (input.result.failed ? 'failed' : 'completed') : 'running', - ...(input.result - ? { output: boundInlineText(input.result.output, DEFAULT_JOURNAL_PAYLOAD_LIMITS).bounded } - : {}) + ...(input.result ? { output: boundInlineText(input.result.output, limits).bounded } : {}) } } diff --git a/src/main/claude/claude-structured-journal-translation.test.ts b/src/main/claude/claude-structured-journal-translation.test.ts index 951872f0c69..07d786c2842 100644 --- a/src/main/claude/claude-structured-journal-translation.test.ts +++ b/src/main/claude/claude-structured-journal-translation.test.ts @@ -17,7 +17,7 @@ import { } from '../native-chat/agent-session-wire/structured-agent-session-event-sink' import { boundInlineText, - DEFAULT_JOURNAL_PAYLOAD_LIMITS + UNRETAINED_JOURNAL_PAYLOAD_LIMITS } from '../native-chat/agent-session-journal/journal-payload-bounds' import type { ClaudePendingPrompt } from './claude-structured-prompt-replies' import { createClaudeJournalTranslator } from './claude-structured-journal-translation' @@ -567,7 +567,7 @@ describe('Claude structured journal translation', () => { expect(state.items.at(-1)?.body).toEqual({ kind: 'status', - text: boundInlineText(thinking, DEFAULT_JOURNAL_PAYLOAD_LIMITS).text + text: boundInlineText(thinking, UNRETAINED_JOURNAL_PAYLOAD_LIMITS).text }) }) diff --git a/src/main/claude/claude-structured-journal-translation.ts b/src/main/claude/claude-structured-journal-translation.ts index 759c2926763..9d26f89f60d 100644 --- a/src/main/claude/claude-structured-journal-translation.ts +++ b/src/main/claude/claude-structured-journal-translation.ts @@ -1,11 +1,11 @@ import type { AgentJournalItemIdentity } from '../../shared/agent-session-journal-types' import { agentJournalItemKey } from '../../shared/agent-session-journal-item-key' import type { AgentSessionDeltaCoalescerDeps } from '../native-chat/agent-session-wire/agent-session-delta-coalescer' -import type { StructuredAgentSessionEventSink } from '../native-chat/agent-session-wire/structured-agent-session-event-sink' import { - boundInlineText, - DEFAULT_JOURNAL_PAYLOAD_LIMITS -} from '../native-chat/agent-session-journal/journal-payload-bounds' + structuredAgentSessionPayloadLimits, + type StructuredAgentSessionEventSink +} from '../native-chat/agent-session-wire/structured-agent-session-event-sink' +import { boundInlineText } from '../native-chat/agent-session-journal/journal-payload-bounds' import type { ClaudeStructuredSessionEvent } from './claude-structured-session-state' import { claudeMessageBody, @@ -166,7 +166,7 @@ export function createClaudeJournalTranslator( tools.set(tool.id, tool) deps.sink.appendItem( claudeToolIdentity(envelope.sessionId, tool.id), - claudeToolBody({ tool }) + claudeToolBody({ tool }, structuredAgentSessionPayloadLimits(deps.sink)) ) changed = true } @@ -178,7 +178,7 @@ export function createClaudeJournalTranslator( } deps.sink.appendItem( claudeToolIdentity(envelope.sessionId, result.toolUseId), - claudeToolBody({ tool, result }) + claudeToolBody({ tool, result }, structuredAgentSessionPayloadLimits(deps.sink)) ) // A spawn call's result is the parent turn's evidence its child finished. subagents.observeToolResult(result.toolUseId, result.failed) @@ -190,7 +190,7 @@ export function createClaudeJournalTranslator( if (thinking) { deps.sink.appendItem(claudeThinkingIdentity(envelope.sessionId, envelope.uuid), { kind: 'status', - text: boundInlineText(thinking, DEFAULT_JOURNAL_PAYLOAD_LIMITS).text + text: boundInlineText(thinking, structuredAgentSessionPayloadLimits(deps.sink)).text }) changed = true } @@ -234,7 +234,10 @@ export function createClaudeJournalTranslator( promptKey: event.prompt.promptKey }) identities.push(identity) - deps.sink.appendItem(identity, claudeApprovalItem(event.prompt)) + deps.sink.appendItem( + identity, + claudeApprovalItem(event.prompt, structuredAgentSessionPayloadLimits(deps.sink)) + ) deps.bindPromptItemId?.(agentJournalItemKey(identity), event.prompt.promptKey) } promptItems.set(event.prompt.promptKey, identities) diff --git a/src/main/claude/claude-structured-prompt-items.ts b/src/main/claude/claude-structured-prompt-items.ts index 3bdf8ab6091..ed2ea85d712 100644 --- a/src/main/claude/claude-structured-prompt-items.ts +++ b/src/main/claude/claude-structured-prompt-items.ts @@ -7,7 +7,7 @@ import type { } from '../../shared/agent-session-journal-types' import { boundInlineText, - DEFAULT_JOURNAL_PAYLOAD_LIMITS + type JournalPayloadLimits } from '../native-chat/agent-session-journal/journal-payload-bounds' import { claudeRecord, claudeText } from './claude-structured-item-translation' import { @@ -43,12 +43,15 @@ export function claudePromptIdentity(input: { } } -export function claudeApprovalItem(prompt: ClaudePendingPrompt): AgentJournalApprovalItem { +export function claudeApprovalItem( + prompt: ClaudePendingPrompt, + limits: JournalPayloadLimits +): AgentJournalApprovalItem { const serialized = JSON.stringify(prompt.input) return { kind: 'approval', title: `Allow ${prompt.toolName}?`, - detail: serialized ? boundInlineText(serialized, DEFAULT_JOURNAL_PAYLOAD_LIMITS).text : null, + detail: serialized ? boundInlineText(serialized, limits).text : null, options: CLAUDE_APPROVAL_DECISIONS.map((decision) => ({ id: decision, label: APPROVAL_LABELS[decision] diff --git a/src/main/claude/claude-structured-provider-fallback.ts b/src/main/claude/claude-structured-provider-fallback.ts index 68aec07976b..d919242b766 100644 --- a/src/main/claude/claude-structured-provider-fallback.ts +++ b/src/main/claude/claude-structured-provider-fallback.ts @@ -1,8 +1,8 @@ -import type { StructuredAgentSessionEventSink } from '../native-chat/agent-session-wire/structured-agent-session-event-sink' import { - boundInlineText, - DEFAULT_JOURNAL_PAYLOAD_LIMITS -} from '../native-chat/agent-session-journal/journal-payload-bounds' + structuredAgentSessionPayloadLimits, + type StructuredAgentSessionEventSink +} from '../native-chat/agent-session-wire/structured-agent-session-event-sink' +import { boundInlineText } from '../native-chat/agent-session-journal/journal-payload-bounds' import { CLAUDE_STREAM_JSON_FRAME_KINDS } from '../native-chat/agent-session-wire/claude-stream-json-frame-schema' import { readableProviderFrameText, @@ -112,13 +112,12 @@ export function createClaudeProviderFrameFallback( return { append: (kind, payload, displayText) => { sequence += 1 - const translated = unhandledProviderFrameJournalItem('claude', kind, payload) + const limits = structuredAgentSessionPayloadLimits(sink) + const translated = unhandledProviderFrameJournalItem('claude', kind, payload, limits) if (!translated) { return } - const bounded = displayText - ? boundInlineText(displayText, DEFAULT_JOURNAL_PAYLOAD_LIMITS).text - : null + const bounded = displayText ? boundInlineText(displayText, limits).text : null sink.appendItem( { provider: 'orca', diff --git a/src/main/codex/codex-image-item-translation.ts b/src/main/codex/codex-image-item-translation.ts index 7d5e70a0c51..2bf5effcf16 100644 --- a/src/main/codex/codex-image-item-translation.ts +++ b/src/main/codex/codex-image-item-translation.ts @@ -1,12 +1,12 @@ import type { AgentJournalMessageItem } from '../../shared/agent-session-journal-types' import type { NativeChatBlock } from '../../shared/native-chat-types' import { buildImageDataUri } from '../../shared/image-data-uri' -import { DEFAULT_JOURNAL_PAYLOAD_LIMITS } from '../native-chat/agent-session-journal/journal-payload-bounds' +import { DEFAULT_JOURNAL_INLINE_HEAD_BYTES } from '../native-chat/agent-session-journal/journal-payload-bounds' import { readString } from './codex-item-field-readers' import type { CodexThreadItem } from './codex-thread-item-identity' // Leave room for operation text and the journal envelope beside inline image bytes. -const MAX_IMAGE_REFERENCE_BYTES = DEFAULT_JOURNAL_PAYLOAD_LIMITS.inlineHeadBytes / 2 +const MAX_IMAGE_REFERENCE_BYTES = DEFAULT_JOURNAL_INLINE_HEAD_BYTES / 2 function imagePath(item: CodexThreadItem, key: string): string | null { const value = readString(item, key) diff --git a/src/main/codex/codex-notice-item-translation.test.ts b/src/main/codex/codex-notice-item-translation.test.ts index 15f638c4015..be7313d588f 100644 --- a/src/main/codex/codex-notice-item-translation.test.ts +++ b/src/main/codex/codex-notice-item-translation.test.ts @@ -2,30 +2,44 @@ import { describe, expect, it } from 'vitest' import { codexItemBody, codexStreamingJournalItem } from './codex-structured-item-translation' import { AgentJournalItemBodySchema } from '../../shared/agent-session-journal-schemas' import { projectStructuredItemsToNativeChat } from '../../shared/structured-agent-session-projection' +import { UNRETAINED_JOURNAL_PAYLOAD_LIMITS } from '../native-chat/agent-session-journal/journal-payload-bounds' describe('plan document translation', () => { it('marks both complete documents and streaming snapshots', () => { const item = { id: 'plan-1', type: 'plan', text: '# Plan\n\nReadable prose.' } - expect(codexItemBody(item)).toEqual({ + expect(codexItemBody(item, UNRETAINED_JOURNAL_PAYLOAD_LIMITS)).toEqual({ kind: 'status', text: item.text, presentation: 'plan-document' }) - expect(codexStreamingJournalItem(item, '# Plan\n\nPartial').body).toEqual({ + expect( + codexStreamingJournalItem(item, '# Plan\n\nPartial', UNRETAINED_JOURNAL_PAYLOAD_LIMITS).body + ).toEqual({ kind: 'status', text: '# Plan\n\nPartial', presentation: 'plan-document' }) - expect(codexItemBody({ id: 'plan-1', type: 'plan' })).toBeNull() + expect( + codexItemBody({ id: 'plan-1', type: 'plan' }, UNRETAINED_JOURNAL_PAYLOAD_LIMITS) + ).toBeNull() }) it('preserves the full existing reasoning body byte for byte', () => { expect( - codexItemBody({ id: 'r', type: 'reasoning', summary: ['Thinking through the problem.'] }) + codexItemBody( + { id: 'r', type: 'reasoning', summary: ['Thinking through the problem.'] }, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) ).toEqual({ kind: 'status', text: 'Thinking through the problem.' }) - expect(codexStreamingJournalItem({ id: 'r', type: 'reasoning' }, 'Thinking…')).toEqual({ + expect( + codexStreamingJournalItem( + { id: 'r', type: 'reasoning' }, + 'Thinking…', + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) + ).toEqual({ body: { kind: 'status', text: 'Thinking…' }, handled: true }) @@ -36,7 +50,9 @@ describe('image item translation', () => { it.each(['/remote/work/image.png', 'C:\\work\\image.png'])( 'preserves the execution-host path %s and old-reader operation text', (path) => { - expect(codexItemBody({ id: 'view', type: 'imageView', path })).toEqual({ + expect( + codexItemBody({ id: 'view', type: 'imageView', path }, UNRETAINED_JOURNAL_PAYLOAD_LIMITS) + ).toEqual({ kind: 'message', role: 'assistant', blocks: [ @@ -47,13 +63,16 @@ describe('image item translation', () => { } ) it('prefers saved paths over unbounded inline image data', () => { - const body = codexItemBody({ - id: 'gen', - type: 'imageGeneration', - status: 'completed', - savedPath: '/remote/image.png', - result: 'A'.repeat(100_000) - }) + const body = codexItemBody( + { + id: 'gen', + type: 'imageGeneration', + status: 'completed', + savedPath: '/remote/image.png', + result: 'A'.repeat(100_000) + }, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) expect(body).toEqual({ kind: 'message', role: 'assistant', @@ -72,7 +91,10 @@ describe('image item translation', () => { 'maps bounded image data onto a meaningful existing image-ref: %s', (result) => { expect( - codexItemBody({ id: 'gen', type: 'imageGeneration', status: 'completed', result }) + codexItemBody( + { id: 'gen', type: 'imageGeneration', status: 'completed', result }, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) ).toMatchObject({ kind: 'message', role: 'assistant', @@ -86,12 +108,15 @@ describe('image item translation', () => { it.each(['A'.repeat(20_000), 'not valid image bytes', 'data:text/html;base64,AAAA', ''])( 'keeps unavailable results bounded and readable', (result) => { - const body = codexItemBody({ - id: 'gen', - type: 'imageGeneration', - status: 'completed', - result - }) + const body = codexItemBody( + { + id: 'gen', + type: 'imageGeneration', + status: 'completed', + result + }, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) expect(body).toEqual({ kind: 'message', role: 'assistant', @@ -105,7 +130,12 @@ describe('image item translation', () => { [{ status: 'completed', failure: { type: 'usageLimitExceeded' } }, 'Image generation failed'], [{ status: 'future-state' }, 'Image generation: preview unavailable'] ])('does not invent completed output for %j', (fields, text) => { - expect(codexItemBody({ id: 'gen', type: 'imageGeneration', ...fields })).toEqual({ + expect( + codexItemBody( + { id: 'gen', type: 'imageGeneration', ...fields }, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) + ).toEqual({ kind: 'message', role: 'assistant', blocks: [{ type: 'text', text }] diff --git a/src/main/codex/codex-prompt-registry-bounds.ts b/src/main/codex/codex-prompt-registry-bounds.ts index e5f79fe3a31..54a7ce640a3 100644 --- a/src/main/codex/codex-prompt-registry-bounds.ts +++ b/src/main/codex/codex-prompt-registry-bounds.ts @@ -1,3 +1,4 @@ +import { JOURNAL_OVERFLOW_NOT_RETAINED } from '../native-chat/agent-session-journal/journal-overflow-store' import { boundPayload, digestPayload @@ -19,8 +20,11 @@ export function codexJournalPromptIdPart(value: string): string { return value } const suffix = `#${digestPayload(value).slice(0, 32)}` + // A prompt id, not user content: the clipped tail is replaced by the digest + // suffix that keeps it unique, so there is nothing to retain. const bounded = boundPayload(value, { - inlineHeadBytes: CODEX_JOURNAL_PROMPT_ID_COMPONENT_MAX_BYTES - suffix.length + inlineHeadBytes: CODEX_JOURNAL_PROMPT_ID_COMPONENT_MAX_BYTES - suffix.length, + overflow: JOURNAL_OVERFLOW_NOT_RETAINED }) return `${bounded.head}${suffix}` } diff --git a/src/main/codex/codex-structured-item-streams.ts b/src/main/codex/codex-structured-item-streams.ts index e67263f133d..1e3e243e2b6 100644 --- a/src/main/codex/codex-structured-item-streams.ts +++ b/src/main/codex/codex-structured-item-streams.ts @@ -1,4 +1,5 @@ import { agentJournalItemKey } from '../../shared/agent-session-journal-item-key' +import { structuredAgentSessionPayloadLimits } from '../native-chat/agent-session-wire/structured-agent-session-event-sink' import { createAgentSessionDeltaCoalescer } from '../native-chat/agent-session-wire/agent-session-delta-coalescer' import { CodexItemStreamRetention } from './codex-item-stream-retention' import { @@ -99,7 +100,8 @@ export function createCodexStructuredItemStreams( } const append = (state: CodexItemStreamState, text: string): boolean => { - const translated = codexStreamingJournalItem(state.item, text) + const limits = structuredAgentSessionPayloadLimits(deps.sink) + const translated = codexStreamingJournalItem(state.item, text, limits) if (!translated.body) { return true } @@ -237,7 +239,8 @@ export function createCodexStructuredItemStreams( return { handled: true, admission: { accepted: false, reason: 'failed' } } } state.item = { ...state.item, changes: paramsRecord.changes } - const translated = codexJournalItem(state.item) + const limits = structuredAgentSessionPayloadLimits(deps.sink) + const translated = codexJournalItem(state.item, limits) if (translated.body) { const nextPending: CodexPendingItemPatch = { identity: state.identity, diff --git a/src/main/codex/codex-structured-item-translation.test.ts b/src/main/codex/codex-structured-item-translation.test.ts index 201df58bc90..9c398db4f6e 100644 --- a/src/main/codex/codex-structured-item-translation.test.ts +++ b/src/main/codex/codex-structured-item-translation.test.ts @@ -6,10 +6,10 @@ import { describeToolInput } from '../../shared/native-chat-tool-summary' import { - codexItemBody, + codexItemBody as boundCodexItemBody, codexItemIdentity, - codexJournalItem, - codexMessageBlocks, + codexJournalItem as boundCodexJournalItem, + codexMessageBlocks as boundCodexMessageBlocks, CodexTurnOrdinals, MAX_CODEX_TURN_ORDINAL_BYTES, MAX_CODEX_TURN_ORDINAL_ENTRIES, @@ -17,6 +17,16 @@ import { readCodexThreadItem, type CodexThreadItem } from './codex-structured-item-translation' +import { UNRETAINED_JOURNAL_PAYLOAD_LIMITS } from '../native-chat/agent-session-journal/journal-payload-bounds' + +// These cases assert the clip and the shape of a translated row. Retention of +// what the clip drops has its own suite, so they all bound the same way. +const codexItemBody = (item: CodexThreadItem) => + boundCodexItemBody(item, UNRETAINED_JOURNAL_PAYLOAD_LIMITS) +const codexJournalItem = (item: CodexThreadItem) => + boundCodexJournalItem(item, UNRETAINED_JOURNAL_PAYLOAD_LIMITS) +const codexMessageBlocks = (item: CodexThreadItem) => + boundCodexMessageBlocks(item, UNRETAINED_JOURNAL_PAYLOAD_LIMITS) /** The tool-call input a Codex item lands on, which is what the row label and * the collapsed run header are both derived from. */ diff --git a/src/main/codex/codex-structured-item-translation.ts b/src/main/codex/codex-structured-item-translation.ts index 576f3fb19ec..051b214440c 100644 --- a/src/main/codex/codex-structured-item-translation.ts +++ b/src/main/codex/codex-structured-item-translation.ts @@ -4,7 +4,7 @@ import type { NativeChatBlock } from '../../shared/native-chat-types' import { boundInlineText, boundToolInput, - DEFAULT_JOURNAL_PAYLOAD_LIMITS + type JournalPayloadLimits } from '../native-chat/agent-session-journal/journal-payload-bounds' import { unhandledProviderFrameJournalItem } from '../native-chat/agent-session-wire/unhandled-provider-frame' import { codexImageItemBody } from './codex-image-item-translation' @@ -31,13 +31,16 @@ export { // Codex thread items → journal item bodies. /** `userMessage` carries structured content parts; `agentMessage` a flat text. */ -export function codexMessageBlocks(item: CodexThreadItem): NativeChatBlock[] { +export function codexMessageBlocks( + item: CodexThreadItem, + limits: JournalPayloadLimits +): NativeChatBlock[] { const text = item.type === 'agentMessage' ? (readString(item, 'text') ?? readTextContent(item, 'content')) : readString(item, 'text') if (text !== null) { - return [{ type: 'text', text: boundInlineText(text, DEFAULT_JOURNAL_PAYLOAD_LIMITS).text }] + return [{ type: 'text', text: boundInlineText(text, limits).text }] } const content = item.content if (!Array.isArray(content)) { @@ -52,7 +55,7 @@ export function codexMessageBlocks(item: CodexThreadItem): NativeChatBlock[] { if (partText !== null) { blocks.push({ type: 'text', - text: boundInlineText(partText, DEFAULT_JOURNAL_PAYLOAD_LIMITS).text + text: boundInlineText(partText, limits).text }) continue } @@ -85,9 +88,9 @@ export type CodexJournalItem = { handled: boolean } -function commandItem(item: CodexThreadItem): CodexJournalItem { +function commandItem(item: CodexThreadItem, limits: JournalPayloadLimits): CodexJournalItem { const output = readFirstString(item, ['aggregatedOutput', 'aggregated_output']) - const bounded = output === null ? null : boundInlineText(output, DEFAULT_JOURNAL_PAYLOAD_LIMITS) + const bounded = output === null ? null : boundInlineText(output, limits) const parsed = commandActionFacts(item) return { body: { @@ -96,7 +99,7 @@ function commandItem(item: CodexThreadItem): CodexJournalItem { // Raw command and cwd stay so the expanded view still shows what ran. input: boundToolInput( { command: item.command ?? null, cwd: item.cwd ?? null, ...parsed?.fields }, - DEFAULT_JOURNAL_PAYLOAD_LIMITS + limits ), state: commandState(item), ...toolExecutionMetadata(item), @@ -106,7 +109,7 @@ function commandItem(item: CodexThreadItem): CodexJournalItem { } } -function fileChangeItem(item: CodexThreadItem): CodexJournalItem { +function fileChangeItem(item: CodexThreadItem, limits: JournalPayloadLimits): CodexJournalItem { const changes = Array.isArray(item.changes) ? item.changes.flatMap((change) => { const record = typeof change === 'object' && change !== null ? readRecord(change) : {} @@ -120,14 +123,14 @@ function fileChangeItem(item: CodexThreadItem): CodexJournalItem { body: { kind: 'tool-call', name: 'apply_patch', - input: boundToolInput({ changes: item.changes ?? null }, DEFAULT_JOURNAL_PAYLOAD_LIMITS), + input: boundToolInput({ changes: item.changes ?? null }, limits), state: commandState(item) }, handled: true } } const patch = changes.map((change) => change.diff).join('\n') - const bounded = boundInlineText(patch, DEFAULT_JOURNAL_PAYLOAD_LIMITS).bounded + const bounded = boundInlineText(patch, limits).bounded return { body: { kind: 'diff', @@ -161,18 +164,18 @@ function mcpToolArguments(value: unknown): unknown { return Array.isArray(value) ? { arguments: value } : Object.keys(value).length > 0 ? value : null } -function mcpToolCallItem(item: CodexThreadItem): CodexJournalItem { +function mcpToolCallItem(item: CodexThreadItem, limits: JournalPayloadLimits): CodexJournalItem { const server = readString(item, 'server') const tool = readString(item, 'tool') const failure = readString(readRecord(item.error), 'message') const text = failure ?? readTextContent(readRecord(item.result), 'content') - const bounded = text === null ? null : boundInlineText(text, DEFAULT_JOURNAL_PAYLOAD_LIMITS) + const bounded = text === null ? null : boundInlineText(text, limits) return { body: { kind: 'tool-call', name: mcpToolCallName(item), ...(server && tool ? { mcpIdentity: { server, tool } } : {}), - input: boundToolInput(mcpToolArguments(item.arguments), DEFAULT_JOURNAL_PAYLOAD_LIMITS), + input: boundToolInput(mcpToolArguments(item.arguments), limits), state: failure === null ? commandState(item) : 'failed', ...(bounded === null ? {} : { output: bounded.bounded }) }, @@ -205,16 +208,16 @@ function webSearchInput(item: CodexThreadItem): Record | null { * action, then sends the action, so `action` is the completion signal — a * completed item's own `query` is routinely still empty. The hits arrive on * `results` and are the call's output. */ -function webSearchItem(item: CodexThreadItem): CodexJournalItem { +function webSearchItem(item: CodexThreadItem, limits: JournalPayloadLimits): CodexJournalItem { const results = toolWebSearchResults(item.results) const hits = Array.isArray(item.results) && item.results.length > 0 ? item.results : null - const bounded = hits && boundInlineText(JSON.stringify(hits), DEFAULT_JOURNAL_PAYLOAD_LIMITS) + const bounded = hits && boundInlineText(JSON.stringify(hits), limits) return { body: { kind: 'tool-call', name: 'web_search', ...(results.length > 0 ? { webSearchResults: results } : {}), - input: boundToolInput(webSearchInput(item), DEFAULT_JOURNAL_PAYLOAD_LIMITS), + input: boundToolInput(webSearchInput(item), limits), state: item.action === null || item.action === undefined ? 'running' : 'completed', ...(bounded === null ? {} : { output: bounded.bounded }) }, @@ -228,9 +231,12 @@ function webSearchItem(item: CodexThreadItem): CodexJournalItem { * Known empty items wait for later deltas. Unknown types become bounded status * rows so a provider release cannot make new activity invisible. */ -export function codexJournalItem(item: CodexThreadItem): CodexJournalItem { +export function codexJournalItem( + item: CodexThreadItem, + limits: JournalPayloadLimits +): CodexJournalItem { if (item.type === 'userMessage' || item.type === 'agentMessage') { - const blocks = codexMessageBlocks(item) + const blocks = codexMessageBlocks(item, limits) return { body: blocks.length === 0 @@ -240,16 +246,16 @@ export function codexJournalItem(item: CodexThreadItem): CodexJournalItem { } } if (item.type === 'commandExecution') { - return commandItem(item) + return commandItem(item, limits) } if (item.type === 'fileChange') { - return fileChangeItem(item) + return fileChangeItem(item, limits) } if (item.type === 'mcpToolCall') { - return mcpToolCallItem(item) + return mcpToolCallItem(item, limits) } if (item.type === 'webSearch') { - return webSearchItem(item) + return webSearchItem(item, limits) } if (item.type === 'imageView' || item.type === 'imageGeneration') { return { body: codexImageItemBody(item), handled: true } @@ -262,7 +268,7 @@ export function codexJournalItem(item: CodexThreadItem): CodexJournalItem { ? null : { kind: 'status', - text: boundInlineText(text, DEFAULT_JOURNAL_PAYLOAD_LIMITS).text, + text: boundInlineText(text, limits).text, presentation: 'plan-document' }, handled: true @@ -274,43 +280,50 @@ export function codexJournalItem(item: CodexThreadItem): CodexJournalItem { readTextContent(item, 'summary') ?? readTextContent(item, 'content') return { - body: - text === null - ? null - : { kind: 'status', text: boundInlineText(text, DEFAULT_JOURNAL_PAYLOAD_LIMITS).text }, + body: text === null ? null : { kind: 'status', text: boundInlineText(text, limits).text }, handled: true } } - const unhandled = unhandledProviderFrameJournalItem('codex', `item:${item.type}`, item) + const unhandled = unhandledProviderFrameJournalItem('codex', `item:${item.type}`, item, limits) return unhandled ? { body: unhandled.body, handled: false } : { body: null, handled: true } } -export function codexItemBody(item: CodexThreadItem): AgentJournalItemBody | null { - return codexJournalItem(item).body +export function codexItemBody( + item: CodexThreadItem, + limits: JournalPayloadLimits +): AgentJournalItemBody | null { + return codexJournalItem(item, limits).body } /** Snapshot body for text still streaming, before its item completes. */ -export function codexStreamingMessageBody(text: string): AgentJournalItemBody { +export function codexStreamingMessageBody( + text: string, + limits: JournalPayloadLimits +): AgentJournalItemBody { return { kind: 'message', role: 'assistant', - blocks: [{ type: 'text', text: boundInlineText(text, DEFAULT_JOURNAL_PAYLOAD_LIMITS).text }] + blocks: [{ type: 'text', text: boundInlineText(text, limits).text }] } } /** Snapshot body for any item-level stream, keyed onto its parent item. */ -export function codexStreamingJournalItem(item: CodexThreadItem, text: string): CodexJournalItem { +export function codexStreamingJournalItem( + item: CodexThreadItem, + text: string, + limits: JournalPayloadLimits +): CodexJournalItem { if (item.type === 'agentMessage') { - return { body: codexStreamingMessageBody(text), handled: true } + return { body: codexStreamingMessageBody(text, limits), handled: true } } if (item.type === 'commandExecution') { - return commandItem({ ...item, aggregatedOutput: text }) + return commandItem({ ...item, aggregatedOutput: text }, limits) } if (item.type === 'fileChange') { const path = Array.isArray(item.changes) ? readString(readRecord(item.changes[0]), 'path') : null - const bounded = boundInlineText(text, DEFAULT_JOURNAL_PAYLOAD_LIMITS).bounded + const bounded = boundInlineText(text, limits).bounded return { body: { kind: 'diff', path: path ?? 'pending patch', patch: bounded }, handled: true @@ -320,12 +333,12 @@ export function codexStreamingJournalItem(item: CodexThreadItem, text: string): return { body: { kind: 'status', - text: boundInlineText(text, DEFAULT_JOURNAL_PAYLOAD_LIMITS).text, + text: boundInlineText(text, limits).text, presentation: 'plan-document' }, handled: true } } - const bounded = boundInlineText(text, DEFAULT_JOURNAL_PAYLOAD_LIMITS) + const bounded = boundInlineText(text, limits) return { body: { kind: 'status', text: bounded.text }, handled: true } } diff --git a/src/main/codex/codex-structured-journal-generic-frames.ts b/src/main/codex/codex-structured-journal-generic-frames.ts index 6c211a8f5a4..121cd2ba78a 100644 --- a/src/main/codex/codex-structured-journal-generic-frames.ts +++ b/src/main/codex/codex-structured-journal-generic-frames.ts @@ -12,6 +12,7 @@ import { MAX_CODEX_GENERIC_TURN_BUCKETS } from './codex-structured-journal-limits' import { readCodexTurnId } from './codex-structured-thread-facts' +import { structuredAgentSessionPayloadLimits } from '../native-chat/agent-session-wire/structured-agent-session-event-sink' const OVERFLOW_BUCKET = '__codex-generic-overflow__' type SuppressedSummary = { count: number; publishedCount: number } @@ -63,7 +64,12 @@ export class CodexJournalGenericFrames { payload: unknown, threadId = 'session' ): CodexJournalTranslationAdmission { - const translated = unhandledProviderFrameJournalItem('codex', kind, payload) + const translated = unhandledProviderFrameJournalItem( + 'codex', + kind, + payload, + structuredAgentSessionPayloadLimits(this.deps.sink) + ) // A frame the classifier declines is deliberately not journaled, which is success. // Failing admission here force-closes the provider through the retry queue. if (!translated) { diff --git a/src/main/codex/codex-structured-journal-items.ts b/src/main/codex/codex-structured-journal-items.ts index 62091580da3..fa6a78f6c8b 100644 --- a/src/main/codex/codex-structured-journal-items.ts +++ b/src/main/codex/codex-structured-journal-items.ts @@ -3,6 +3,7 @@ import type { AgentJournalItemIdentity } from '../../shared/agent-session-journal-types' import { requiresTerminalSettlement } from '../native-chat/agent-session-journal/journal-terminal-settlement' +import { structuredAgentSessionPayloadLimits } from '../native-chat/agent-session-wire/structured-agent-session-event-sink' import { codexItemIdentity, codexJournalItem, @@ -89,7 +90,7 @@ export class CodexJournalItems { ) { return { handled: true, admission: { accepted: false, reason: 'failed' } } } - const translated = codexJournalItem(item) + const translated = codexJournalItem(item, structuredAgentSessionPayloadLimits(this.deps.sink)) const command = readCodexJournalString(item, 'command') if (command) { const boundedCommand = Buffer.from(command, 'utf8') @@ -214,7 +215,10 @@ export class CodexJournalItems { } const evicted = this.activeItems.get(oldest) if (evicted) { - const translated = codexJournalItem(evicted.item).body + const translated = codexJournalItem( + evicted.item, + structuredAgentSessionPayloadLimits(this.deps.sink) + ).body if (translated) { const admission = appendCodexLifecycleItem( this.deps.sink, diff --git a/src/main/codex/codex-structured-journal-prompts.ts b/src/main/codex/codex-structured-journal-prompts.ts index 93ecec77f77..e7a61ce6483 100644 --- a/src/main/codex/codex-structured-journal-prompts.ts +++ b/src/main/codex/codex-structured-journal-prompts.ts @@ -1,5 +1,6 @@ import { agentJournalItemKey } from '../../shared/agent-session-journal-item-key' import { cancelledJournalPromptBody } from '../native-chat/agent-session-journal/journal-prompt-body-bounds' +import { structuredAgentSessionPayloadLimits } from '../native-chat/agent-session-wire/structured-agent-session-event-sink' import { codexApprovalItem, codexPromptIdentity, @@ -38,7 +39,8 @@ export class CodexJournalPrompts { const questions = codexQuestionItems({ threadId: event.threadId, promptKey: event.promptKey, - params: event.params + params: event.params, + limits: structuredAgentSessionPayloadLimits(this.deps.sink) }) const promptItems = questions.map(({ identity, body }) => ({ identity, body })) const admission = this.admit(event, promptItems) @@ -63,7 +65,8 @@ export class CodexJournalPrompts { const body = codexApprovalItem({ method: event.method, params: event.params, - detail: this.detailFor(event.threadId, event.codexItemId) + detail: this.detailFor(event.threadId, event.codexItemId), + limits: structuredAgentSessionPayloadLimits(this.deps.sink) }) const admission = this.admit(event, [{ identity, body }]) if (!admission.accepted) { diff --git a/src/main/codex/codex-structured-journal-settlement.ts b/src/main/codex/codex-structured-journal-settlement.ts index 8c8bf091839..e32de7029bd 100644 --- a/src/main/codex/codex-structured-journal-settlement.ts +++ b/src/main/codex/codex-structured-journal-settlement.ts @@ -21,6 +21,7 @@ import { import type { CodexStructuredItemStreams } from './codex-structured-item-streams' import type { CodexStructuredSessionEvent } from './codex-structured-session-adapter' import { codexCommandOutlivesTurn } from './codex-command-lifecycle' +import { structuredAgentSessionPayloadLimits } from '../native-chat/agent-session-wire/structured-agent-session-event-sink' export type CodexActiveJournalItem = { threadId: string @@ -50,9 +51,10 @@ export function settleCodexJournalSession(input: { const turnOrdinalsToForget: { threadId: string; turnId: string }[] = [] for (const active of input.activeItems.values()) { const streamed = input.streams.snapshot(active.threadId, active.item.id) + const limits = structuredAgentSessionPayloadLimits(input.sink) const translated = streamed - ? codexStreamingJournalItem(active.item, streamed.text) - : codexJournalItem(active.item) + ? codexStreamingJournalItem(active.item, streamed.text, limits) + : codexJournalItem(active.item, limits) const body = interruptedBody(translated.body) if (body) { mutations.push({ kind: 'item', identity: active.identity, body }) @@ -123,9 +125,10 @@ export function settleCodexJournalTurn(input: { continue } const streamed = input.streams.snapshot(active.threadId, active.item.id) + const limits = structuredAgentSessionPayloadLimits(input.sink) const translated = streamed - ? codexStreamingJournalItem(active.item, streamed.text) - : codexJournalItem(active.item) + ? codexStreamingJournalItem(active.item, streamed.text, limits) + : codexJournalItem(active.item, limits) const body = interruptedBody(translated.body) if (body) { mutations.push({ kind: 'item', identity: active.identity, body }) @@ -176,9 +179,10 @@ export function settleCodexOversizedNotification(input: { continue } const streamed = input.streams.snapshot(active.threadId, active.item.id) + const limits = structuredAgentSessionPayloadLimits(input.sink) const translated = streamed - ? codexStreamingJournalItem(active.item, streamed.text) - : codexJournalItem(active.item) + ? codexStreamingJournalItem(active.item, streamed.text, limits) + : codexJournalItem(active.item, limits) const body = interruptedBody(translated.body) if (body) { mutations.push({ kind: 'item', identity: active.identity, body }) diff --git a/src/main/codex/codex-structured-prompt-items.test.ts b/src/main/codex/codex-structured-prompt-items.test.ts index e5d996f05d5..7ab544ca611 100644 --- a/src/main/codex/codex-structured-prompt-items.test.ts +++ b/src/main/codex/codex-structured-prompt-items.test.ts @@ -11,6 +11,7 @@ import { encodeCodexQuestionOptionId } from './codex-structured-prompt-replies' import { MAX_JOURNAL_LIFECYCLE_BATCH_BYTES } from '../native-chat/agent-session-journal/journal-row-schema' +import { UNRETAINED_JOURNAL_PAYLOAD_LIMITS } from '../native-chat/agent-session-journal/journal-payload-bounds' const THREAD_ID = 'thread-abc' const CODEX_ITEM_ID = 'item-4' @@ -42,7 +43,8 @@ describe('codex approval items', () => { const command = codexApprovalItem({ method: CODEX_COMMAND_APPROVAL_METHOD, params: { availableDecisions: ['accept'] }, - detail: 'rm -rf build' + detail: 'rm -rf build', + limits: UNRETAINED_JOURNAL_PAYLOAD_LIMITS }) expect(command).toMatchObject({ @@ -52,10 +54,20 @@ describe('codex approval items', () => { resolution: { state: 'pending', selectedOptionId: null, resolvedBy: null, resolvedAt: null } }) expect( - codexApprovalItem({ method: CODEX_FILE_CHANGE_APPROVAL_METHOD, params: {}, detail: null }) + codexApprovalItem({ + method: CODEX_FILE_CHANGE_APPROVAL_METHOD, + params: {}, + detail: null, + limits: UNRETAINED_JOURNAL_PAYLOAD_LIMITS + }) ).toMatchObject({ title: 'Apply file changes?', detail: null }) expect( - codexApprovalItem({ method: 'item/other/requestApproval', params: {}, detail: null }) + codexApprovalItem({ + method: 'item/other/requestApproval', + params: {}, + detail: null, + limits: UNRETAINED_JOURNAL_PAYLOAD_LIMITS + }) ).toMatchObject({ title: 'Approve this action?' }) }) @@ -63,7 +75,8 @@ describe('codex approval items', () => { const item = codexApprovalItem({ method: CODEX_COMMAND_APPROVAL_METHOD, params: { reason: 'writes outside the workspace' }, - detail: 'rm -rf build' + detail: 'rm -rf build', + limits: UNRETAINED_JOURNAL_PAYLOAD_LIMITS }) expect(item.detail).toBe('writes outside the workspace') @@ -74,21 +87,24 @@ describe('codex approval items', () => { codexApprovalItem({ method: CODEX_COMMAND_APPROVAL_METHOD, params: { command: ['git', 'status'] }, - detail: 'parent command' + detail: 'parent command', + limits: UNRETAINED_JOURNAL_PAYLOAD_LIMITS }).detail ).toBe('git status') expect( codexApprovalItem({ method: CODEX_COMMAND_APPROVAL_METHOD, params: { command: ['pnpm', 'test'], reason: 'same parent reason' }, - detail: 'parent command' + detail: 'parent command', + limits: UNRETAINED_JOURNAL_PAYLOAD_LIMITS }).detail ).toBe('pnpm test') expect( codexApprovalItem({ method: CODEX_FILE_CHANGE_APPROVAL_METHOD, params: { grantRoot: '/outside' }, - detail: null + detail: null, + limits: UNRETAINED_JOURNAL_PAYLOAD_LIMITS }).detail ).toBe('"/outside"') }) @@ -97,7 +113,8 @@ describe('codex approval items', () => { const item = codexApprovalItem({ method: CODEX_COMMAND_APPROVAL_METHOD, params: { command: 'x'.repeat(2_000_000) }, - detail: null + detail: null, + limits: UNRETAINED_JOURNAL_PAYLOAD_LIMITS }) expect(item.detail).toContain('output truncated') @@ -118,7 +135,12 @@ describe('codex question items', () => { } it('makes one journal item per question, each with its own resolution', () => { - const items = codexQuestionItems({ threadId: THREAD_ID, promptKey: CODEX_ITEM_ID, params }) + const items = codexQuestionItems({ + threadId: THREAD_ID, + promptKey: CODEX_ITEM_ID, + params, + limits: UNRETAINED_JOURNAL_PAYLOAD_LIMITS + }) expect(items.map((item) => item.questionId)).toEqual(['q1', 'q2']) expect(items.map((item) => item.body.question)).toEqual(['Which branch?', 'Proceed?']) @@ -126,7 +148,12 @@ describe('codex question items', () => { }) it('keys each question separately so two answers cannot collide on one row', () => { - const items = codexQuestionItems({ threadId: THREAD_ID, promptKey: CODEX_ITEM_ID, params }) + const items = codexQuestionItems({ + threadId: THREAD_ID, + promptKey: CODEX_ITEM_ID, + params, + limits: UNRETAINED_JOURNAL_PAYLOAD_LIMITS + }) expect(items.map((item) => item.identity)).toEqual([ { provider: 'orca', clientMessageId: 'codex-prompt:thread-abc:item-4:q1' }, @@ -135,7 +162,12 @@ describe('codex question items', () => { }) it('names the question inside every option id, because codex replies by question', () => { - const items = codexQuestionItems({ threadId: THREAD_ID, promptKey: CODEX_ITEM_ID, params }) + const items = codexQuestionItems({ + threadId: THREAD_ID, + promptKey: CODEX_ITEM_ID, + params, + limits: UNRETAINED_JOURNAL_PAYLOAD_LIMITS + }) expect(items[0]?.body.options).toEqual([ { id: encodeCodexQuestionOptionId('q1', 'main'), label: 'main' }, @@ -148,7 +180,8 @@ describe('codex question items', () => { const items = codexQuestionItems({ threadId: THREAD_ID, promptKey: CODEX_ITEM_ID, - params: { questions: [{ question: 'no id' }, { id: 'q3' }, { id: 'q4', question: 'ok' }] } + params: { questions: [{ question: 'no id' }, { id: 'q3' }, { id: 'q4', question: 'ok' }] }, + limits: UNRETAINED_JOURNAL_PAYLOAD_LIMITS }) expect(items.map((item) => item.questionId)).toEqual(['q4']) @@ -156,7 +189,12 @@ describe('codex question items', () => { it('returns nothing when the request carries no questions at all', () => { expect( - codexQuestionItems({ threadId: THREAD_ID, promptKey: CODEX_ITEM_ID, params: {} }) + codexQuestionItems({ + threadId: THREAD_ID, + promptKey: CODEX_ITEM_ID, + params: {}, + limits: UNRETAINED_JOURNAL_PAYLOAD_LIMITS + }) ).toEqual([]) }) @@ -173,7 +211,8 @@ describe('codex question items', () => { options: [{ label: 'Known' }, { label: 'Other', isOther: true }] } ] - } + }, + limits: UNRETAINED_JOURNAL_PAYLOAD_LIMITS }) expect(withoutOptions?.body).toMatchObject({ options: [], freeTextQuestionId: 'q1' }) @@ -197,7 +236,8 @@ describe('codex question items', () => { options: Array.from({ length: 80 }, () => ({ label: longLabel })) } ] - } + }, + limits: UNRETAINED_JOURNAL_PAYLOAD_LIMITS }) const item = items[0] diff --git a/src/main/codex/codex-structured-prompt-items.ts b/src/main/codex/codex-structured-prompt-items.ts index d088c6684fc..2c70cf2b957 100644 --- a/src/main/codex/codex-structured-prompt-items.ts +++ b/src/main/codex/codex-structured-prompt-items.ts @@ -6,7 +6,7 @@ import type { } from '../../shared/agent-session-journal-types' import { boundInlineText, - DEFAULT_JOURNAL_PAYLOAD_LIMITS + type JournalPayloadLimits } from '../native-chat/agent-session-journal/journal-payload-bounds' import { CODEX_APPROVAL_DECISIONS, @@ -40,7 +40,7 @@ const PENDING = { const MAX_CODEX_PROMPT_QUESTIONS = 64 const MAX_CODEX_PROMPT_OPTIONS = 64 -const PROMPT_OPTION_LIMITS = { ...DEFAULT_JOURNAL_PAYLOAD_LIMITS, inlineHeadBytes: 1024 } +const PROMPT_OPTION_INLINE_HEAD_BYTES = 1024 function readParams(params: unknown): Record { return typeof params === 'object' && params !== null ? (params as Record) : {} @@ -51,12 +51,13 @@ function readString(source: Record, key: string): string | null return typeof value === 'string' && value.length > 0 ? value : null } -function boundPromptText(value: string): string { - return boundInlineText(value, DEFAULT_JOURNAL_PAYLOAD_LIMITS).text +function boundPromptText(value: string, limits: JournalPayloadLimits): string { + return boundInlineText(value, limits).text } -function boundPromptOptionLabel(value: string): string { - return boundInlineText(value, PROMPT_OPTION_LIMITS).text +function boundPromptOptionLabel(value: string, limits: JournalPayloadLimits): string { + return boundInlineText(value, { ...limits, inlineHeadBytes: PROMPT_OPTION_INLINE_HEAD_BYTES }) + .text } /** @@ -82,6 +83,7 @@ export function codexApprovalItem(input: { /** What is being approved, taken from the item Codex already announced — * the approval request itself does not repeat the command or the patch. */ detail: string | null + limits: JournalPayloadLimits }): AgentJournalApprovalItem { const params = readParams(input.params) return { @@ -92,14 +94,17 @@ export function codexApprovalItem(input: { : input.method === CODEX_COMMAND_APPROVAL_METHOD ? 'Run a command?' : 'Approve this action?', - detail: boundNullablePromptText(approvalDetail(params) ?? input.detail), + detail: boundNullablePromptText(approvalDetail(params) ?? input.detail, input.limits), options: codexApprovalOptions(input.params), resolution: { ...PENDING } } } -function boundNullablePromptText(value: string | null): string | null { - return value === null ? null : boundPromptText(value) +function boundNullablePromptText( + value: string | null, + limits: JournalPayloadLimits +): string | null { + return value === null ? null : boundPromptText(value, limits) } function approvalDetail(params: Record): string | null { @@ -134,6 +139,7 @@ export function codexQuestionItems(input: { threadId: string promptKey: string params: unknown + limits: JournalPayloadLimits }): CodexQuestionItem[] { const questions = readParams(input.params).questions if (!Array.isArray(questions)) { @@ -152,8 +158,8 @@ export function codexQuestionItems(input: { identity: codexPromptIdentity({ ...input, questionId }), body: { kind: 'question', - question: boundPromptText(prompt), - options: questionOptions(question, questionId), + question: boundPromptText(prompt, input.limits), + options: questionOptions(question, questionId, input.limits), ...(questionAllowsFreeText(question) ? { freeTextQuestionId: codexJournalPromptIdPart(questionId) } : {}), @@ -175,7 +181,8 @@ function questionAllowsFreeText(question: Record): boolean { function questionOptions( question: Record, - questionId: string + questionId: string, + limits: JournalPayloadLimits ): AgentJournalPromptOption[] { const options = question.options if (!Array.isArray(options)) { @@ -193,7 +200,7 @@ function questionOptions( // question id, and the client only ever hands back an option id. mapped.push({ id: encodeCodexJournalQuestionOptionId(questionId, label), - label: boundPromptOptionLabel(label) + label: boundPromptOptionLabel(label, limits) }) } } diff --git a/src/main/codex/codex-tool-identity-translation.test.ts b/src/main/codex/codex-tool-identity-translation.test.ts index 751886f35aa..7852c4b6666 100644 --- a/src/main/codex/codex-tool-identity-translation.test.ts +++ b/src/main/codex/codex-tool-identity-translation.test.ts @@ -1,6 +1,7 @@ import { describe, expect, it } from 'vitest' import { codexItemBody, codexStreamingJournalItem } from './codex-structured-item-translation' import { boundStreamItem } from './codex-structured-item-stream-bounds' +import { UNRETAINED_JOURNAL_PAYLOAD_LIMITS } from '../native-chat/agent-session-journal/journal-payload-bounds' const command = { type: 'commandExecution', @@ -11,7 +12,9 @@ const command = { describe('command row metadata', () => { it.each([0, 127, -1])('preserves exit %s with provider duration', (exitCode) => { - expect(codexItemBody({ ...command, exitCode, durationMs: 400 })).toMatchObject({ + expect( + codexItemBody({ ...command, exitCode, durationMs: 400 }, UNRETAINED_JOURNAL_PAYLOAD_LIMITS) + ).toMatchObject({ kind: 'tool-call', name: 'shell', exitCode, @@ -25,7 +28,7 @@ describe('command row metadata', () => { { exitCode: null, durationMs: null }, { exitCode: 1.5, durationMs: -1 } ]) { - const body = codexItemBody({ ...command, ...fields }) + const body = codexItemBody({ ...command, ...fields }, UNRETAINED_JOURNAL_PAYLOAD_LIMITS) expect(body).not.toHaveProperty('exitCode') expect(body).not.toHaveProperty('durationMs') } @@ -38,19 +41,24 @@ describe('command row metadata', () => { aggregatedOutput: 'x'.repeat(70000) } expect(boundStreamItem(source)).toMatchObject({ exitCode: 0, durationMs: 400 }) - expect(codexStreamingJournalItem(source, 'output').body).toMatchObject({ + expect( + codexStreamingJournalItem(source, 'output', UNRETAINED_JOURNAL_PAYLOAD_LIMITS).body + ).toMatchObject({ exitCode: 0, durationMs: 400 }) }) it('keeps metadata on classified exec rows', () => { expect( - codexItemBody({ - ...command, - exitCode: 0, - durationMs: 15, - commandActions: [{ type: 'read', command: 'cat a.ts', name: 'a.ts', path: 'a.ts' }] - }) + codexItemBody( + { + ...command, + exitCode: 0, + durationMs: 15, + commandActions: [{ type: 'read', command: 'cat a.ts', name: 'a.ts', path: 'a.ts' }] + }, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) ).toMatchObject({ name: 'read', exitCode: 0, durationMs: 15 }) }) }) @@ -58,13 +66,16 @@ describe('command row metadata', () => { describe('web result annotations', () => { it('adds safe result annotations and retains old-reader JSON output', () => { const results = [{ title: 'Docs', url: 'https://example.com/' }, { url: 'javascript:alert(1)' }] - const body = codexItemBody({ - type: 'webSearch', - id: 'web', - query: 'docs', - action: { type: 'search' }, - results - }) + const body = codexItemBody( + { + type: 'webSearch', + id: 'web', + query: 'docs', + action: { type: 'search' }, + results + }, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) expect(body).toMatchObject({ kind: 'tool-call', name: 'web_search', @@ -75,18 +86,26 @@ describe('web result annotations', () => { }) it('leaves legacy and malformed results without an annotation', () => { expect( - codexItemBody({ type: 'webSearch', id: 'web', query: 'docs', results: [{}] }) + codexItemBody( + { type: 'webSearch', id: 'web', query: 'docs', results: [{}] }, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) ).not.toHaveProperty('webSearchResults') }) }) it('annotates only confirmed MCP calls and retains the raw server/tool name', () => { expect( - codexItemBody({ type: 'mcpToolCall', id: 'm', server: 'my_server', tool: 'ns.tool' }) + codexItemBody( + { type: 'mcpToolCall', id: 'm', server: 'my_server', tool: 'ns.tool' }, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) ).toMatchObject({ kind: 'tool-call', name: 'my_server/ns.tool', mcpIdentity: { server: 'my_server', tool: 'ns.tool' } }) - expect(codexItemBody(command)).not.toHaveProperty('mcpIdentity') + expect(codexItemBody(command, UNRETAINED_JOURNAL_PAYLOAD_LIMITS)).not.toHaveProperty( + 'mcpIdentity' + ) }) diff --git a/src/main/codex/codex-turn-ordinals.ts b/src/main/codex/codex-turn-ordinals.ts index 89ed72666db..cdc27e97a39 100644 --- a/src/main/codex/codex-turn-ordinals.ts +++ b/src/main/codex/codex-turn-ordinals.ts @@ -1,3 +1,4 @@ +import { JOURNAL_OVERFLOW_NOT_RETAINED } from '../native-chat/agent-session-journal/journal-overflow-store' import { boundPayload, digestPayload @@ -32,7 +33,9 @@ export class CodexTurnOrdinals { const suffix = `#${digestPayload(value).slice(0, 24)}` return `${ boundPayload(encoded, { - inlineHeadBytes: 256 - Buffer.byteLength(suffix, 'utf8') + inlineHeadBytes: 256 - Buffer.byteLength(suffix, 'utf8'), + // A turn key; the digest suffix carries what the clip drops. + overflow: JOURNAL_OVERFLOW_NOT_RETAINED }).head }${suffix}` } diff --git a/src/main/native-chat/agent-session-journal/journal-legacy-import.test.ts b/src/main/native-chat/agent-session-journal/journal-legacy-import.test.ts index 13673626ebc..6a62727d317 100644 --- a/src/main/native-chat/agent-session-journal/journal-legacy-import.test.ts +++ b/src/main/native-chat/agent-session-journal/journal-legacy-import.test.ts @@ -16,7 +16,7 @@ import { appendLegacyTranscriptMessages, importLegacyTranscriptIntoJournal } from './journal-legacy-import' -import { DEFAULT_JOURNAL_PAYLOAD_LIMITS } from './journal-payload-bounds' +import { UNRETAINED_JOURNAL_PAYLOAD_LIMITS } from './journal-payload-bounds' import { openAgentSessionJournal } from './journal-store-factory' import type { AgentSessionJournal } from './journal-store' @@ -419,7 +419,10 @@ describe('payload bounds on import', () => { agent: 'claude', sessionId: CLAUDE_SESSION, fence: 1, - options: { filePath, limits: { ...DEFAULT_JOURNAL_PAYLOAD_LIMITS, inlineHeadBytes: 1_024 } } + options: { + filePath, + limits: { ...UNRETAINED_JOURNAL_PAYLOAD_LIMITS, inlineHeadBytes: 1_024 } + } }) const item = journal.snapshot().items[0] @@ -540,7 +543,7 @@ describe('import failures', () => { it('bounds oversized legacy tool-call input before journal publication', async () => { const journalDir = join(root, 'bounded-tool-input-journal') - const limits = { ...DEFAULT_JOURNAL_PAYLOAD_LIMITS, inlineHeadBytes: 64 } + const limits = { ...UNRETAINED_JOURNAL_PAYLOAD_LIMITS, inlineHeadBytes: 64 } const journal = await open('claude', CLAUDE_SESSION, { journalDir }) const filePath = await writeFixture('oversized-tool-input.jsonl', [ { @@ -650,7 +653,7 @@ describe('import failures', () => { // row that also has text, and omp's execution cells always pair the invocation // with its output — so the multi-block path carries untrusted tool input too. describe('multi-block legacy messages', () => { - const limits = { ...DEFAULT_JOURNAL_PAYLOAD_LIMITS, inlineHeadBytes: 64 } + const limits = { ...UNRETAINED_JOURNAL_PAYLOAD_LIMITS, inlineHeadBytes: 64 } const oversized = 'x'.repeat(10_000) /** The tool-call block of the first imported multi-block message. */ diff --git a/src/main/native-chat/agent-session-journal/journal-legacy-import.ts b/src/main/native-chat/agent-session-journal/journal-legacy-import.ts index 801c91fb534..d08366c5ee8 100644 --- a/src/main/native-chat/agent-session-journal/journal-legacy-import.ts +++ b/src/main/native-chat/agent-session-journal/journal-legacy-import.ts @@ -35,9 +35,13 @@ import { boundInlineText, boundPayload, boundToolInput, - DEFAULT_JOURNAL_PAYLOAD_LIMITS, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS, type JournalPayloadLimits } from './journal-payload-bounds' + +// The bounds here retain nothing on purpose: the provider's transcript file is +// the durable copy, it is what this import reads, Orca never deletes it, and the +// remnant disclosure names its path in the timeline. import type { AgentSessionJournal } from './journal-store' export type LegacyImportOptions = ResolveSessionFileOptions & { @@ -78,7 +82,7 @@ export async function appendLegacyTranscriptMessages(input: { sessionId: input.sessionId, recordId: message.id }, - legacyItemBody(message, DEFAULT_JOURNAL_PAYLOAD_LIMITS), + legacyItemBody(message, UNRETAINED_JOURNAL_PAYLOAD_LIMITS), { fence: input.fence, observedAt: message.timestamp ?? undefined } ) appended += 1 @@ -112,7 +116,7 @@ export async function prepareLegacyTranscriptImport(input: { options?: LegacyImportOptions }): Promise<{ ok: true; items: JournalReplacementItem[] } | { ok: false; error: string }> { const options = input.options ?? {} - const limits = options.limits ?? DEFAULT_JOURNAL_PAYLOAD_LIMITS + const limits = options.limits ?? UNRETAINED_JOURNAL_PAYLOAD_LIMITS const transcriptAgent = resolveNativeChatTranscriptAgent(input.agent) if (!transcriptAgent) { return { ok: false, error: `Unsupported agent for journal import: ${input.agent}` } diff --git a/src/main/native-chat/agent-session-journal/journal-overflow-store.test.ts b/src/main/native-chat/agent-session-journal/journal-overflow-store.test.ts new file mode 100644 index 00000000000..3a8c112d1be --- /dev/null +++ b/src/main/native-chat/agent-session-journal/journal-overflow-store.test.ts @@ -0,0 +1,112 @@ +// The remainder of a bounded payload is retained, not dropped. + +import { mkdtemp, readdir, rm, writeFile } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { afterEach, beforeEach, describe, expect, it } from 'vitest' +import { + JOURNAL_OVERFLOW_NOT_RETAINED, + journalOverflowDirectory, + journalOverflowSink, + readJournalOverflow +} from './journal-overflow-store' +import { boundInlineText, boundPayload, digestPayload } from './journal-payload-bounds' + +let root: string + +beforeEach(async () => { + root = await mkdtemp(join(tmpdir(), 'orca-journal-overflow-')) +}) + +afterEach(async () => { + await rm(root, { recursive: true, force: true }) +}) + +describe('bounding a payload larger than the row budget', () => { + const limits = (dir: string, quota?: number) => ({ + inlineHeadBytes: 16, + overflow: journalOverflowSink(dir, quota) + }) + + it('retains every byte the clip drops, readable back by the row digest', () => { + const payload = `${'x'.repeat(4_096)}λ tail` + const bounded = boundPayload(payload, limits(root)) + + expect(bounded.truncated).toBe(true) + expect(bounded.head).toHaveLength(16) + expect(bounded.spilled).toBe(true) + // The whole payload, not the clipped head: this is the P0. + expect(readJournalOverflow(root, bounded.digest)).toBe(payload) + expect(bounded.digest).toBe(digestPayload(payload)) + }) + + it('retains the text behind an inline truncation marker too', () => { + const payload = 'assistant output '.repeat(400) + const { text, bounded } = boundInlineText(payload, limits(root)) + + expect(text).toContain('[Orca: output truncated') + expect(readJournalOverflow(root, bounded.digest)).toBe(payload) + }) + + it('leaves a payload that fits inline with nothing retained', () => { + const bounded = boundPayload('short', limits(root)) + + expect(bounded).toMatchObject({ truncated: false, head: 'short' }) + expect(bounded.spilled).toBeUndefined() + expect(readJournalOverflow(root, bounded.digest)).toBeNull() + }) + + it('retains identical payloads once, so item revisions do not accumulate copies', async () => { + const payload = 'y'.repeat(4_096) + boundPayload(payload, limits(root)) + boundPayload(payload, limits(root)) + boundPayload(payload, limits(root)) + + expect(await readdir(journalOverflowDirectory(root))).toHaveLength(1) + }) + + it('marks the row unspilled rather than failing when retention cannot happen', async () => { + // A file where the store needs its directory: every write below fails. + await writeFile(journalOverflowDirectory(root), 'not a directory') + const payload = 'z'.repeat(4_096) + const bounded = boundPayload(payload, limits(root)) + + expect(bounded).toMatchObject({ truncated: true, byteLength: 4_096 }) + expect(bounded.spilled).toBeUndefined() + expect(readJournalOverflow(root, bounded.digest)).toBeNull() + }) + + it('sheds the oldest retained payloads to stay inside the quota', async () => { + const quota = 6_000 + const first = boundPayload('a'.repeat(4_000), limits(root, quota)) + const second = boundPayload('b'.repeat(4_000), limits(root, quota)) + + expect(second.spilled).toBe(true) + expect(readJournalOverflow(root, second.digest)).toBe('b'.repeat(4_000)) + // Two 4 KB payloads do not fit in a 6 KB budget; the older one goes. + expect(readJournalOverflow(root, first.digest)).toBeNull() + expect(await readdir(journalOverflowDirectory(root))).toHaveLength(1) + }) + + it('refuses a single payload larger than the whole quota without evicting for it', () => { + const quota = 2_000 + const kept = boundPayload('a'.repeat(1_000), limits(root, quota)) + const oversized = boundPayload('b'.repeat(9_000), limits(root, quota)) + + expect(oversized).toMatchObject({ truncated: true, byteLength: 9_000 }) + expect(oversized.spilled).toBeUndefined() + expect(readJournalOverflow(root, kept.digest)).toBe('a'.repeat(1_000)) + }) +}) + +describe('a bound that deliberately retains nothing', () => { + it('still reports the clip, its length and its digest', () => { + const bounded = boundPayload('k'.repeat(4_096), { + inlineHeadBytes: 16, + overflow: JOURNAL_OVERFLOW_NOT_RETAINED + }) + + expect(bounded).toMatchObject({ truncated: true, byteLength: 4_096 }) + expect(bounded.spilled).toBeUndefined() + }) +}) diff --git a/src/main/native-chat/agent-session-journal/journal-overflow-store.ts b/src/main/native-chat/agent-session-journal/journal-overflow-store.ts new file mode 100644 index 00000000000..2991b667eb2 --- /dev/null +++ b/src/main/native-chat/agent-session-journal/journal-overflow-store.ts @@ -0,0 +1,144 @@ +// Durable home for payload bytes the live timeline cannot hold. +// +// A row every reconnecting client replays must stay small, so an oversized tool +// result, diff, or assistant message is clipped to a head. The remainder is +// written HERE first, addressed by the sha256 the bound already computes, so the +// clip is a display bound rather than a delete. +// +// Content addressing makes retention idempotent: an item republished at a new +// revision, or two tools returning the same output, retain once. The store lives +// inside the session's journal directory, so removing a session removes its +// payloads with it and nothing has to reference-count them. +// +// Retention is best effort BY DESIGN. A failed write leaves the row exactly as +// it is bounded today and never worse, because refusing the append instead would +// turn a full disk into a lost turn. + +import { existsSync, mkdirSync, readdirSync, renameSync, rmSync, statSync } from 'node:fs' +import { closeSync, fsyncSync, openSync, readFileSync, writeSync } from 'node:fs' +import { join } from 'node:path' + +const OVERFLOW_DIR_NAME = 'overflow' + +/** Per-session budget for retained remainders. Oldest are evicted first; the + * rows keep their head, byte length and digest either way. */ +export const JOURNAL_OVERFLOW_QUOTA_BYTES = 64 * 1024 * 1024 + +/** Where the remainder of a bounded payload goes. `retain` answers whether the + * bytes are now durable, which is what the row records as `spilled`. */ +export type JournalOverflowSink = { + retain: (digest: string, payload: string) => boolean +} + +/** For bounds applied to derived identifiers — a prompt option id, a turn + * ordinal map, an image reference — where the clipped tail is a key fragment + * and not user data. Named so the choice is visible at the call site. */ +export const JOURNAL_OVERFLOW_NOT_RETAINED: JournalOverflowSink = { retain: () => false } + +export function journalOverflowDirectory(journalDir: string): string { + return join(journalDir, OVERFLOW_DIR_NAME) +} + +export function journalOverflowFile(journalDir: string, digest: string): string { + return join(journalOverflowDirectory(journalDir), `${digest}.txt`) +} + +/** Retained bytes for a digest, or null when the payload was never retained or + * has since been evicted. */ +export function readJournalOverflow(journalDir: string, digest: string): string | null { + try { + return readFileSync(journalOverflowFile(journalDir, digest), 'utf8') + } catch { + return null + } +} + +export function journalOverflowSink( + journalDir: string, + quotaBytes: number = JOURNAL_OVERFLOW_QUOTA_BYTES +): JournalOverflowSink { + return { + retain: (digest, payload) => retainJournalOverflow(journalDir, digest, payload, quotaBytes) + } +} + +function retainJournalOverflow( + journalDir: string, + digest: string, + payload: string, + quotaBytes: number +): boolean { + const bytes = Buffer.from(payload, 'utf8') + // A single payload larger than the whole budget would evict everything else + // to store one item; the row's head, length and digest still describe it. + if (bytes.byteLength > quotaBytes) { + return false + } + const file = journalOverflowFile(journalDir, digest) + try { + if (existsSync(file)) { + return true + } + const directory = journalOverflowDirectory(journalDir) + mkdirSync(directory, { recursive: true }) + evictJournalOverflow(directory, quotaBytes - bytes.byteLength) + writeOverflowFile(directory, file, digest, bytes) + return true + } catch { + return false + } +} + +/** Durable before the rename, so a reader never sees a half-written payload + * under a digest that claims to describe it. */ +function writeOverflowFile(directory: string, file: string, digest: string, bytes: Buffer): void { + const staging = join(directory, `${digest}.${process.pid}.partial`) + const handle = openSync(staging, 'w') + try { + writeSync(handle, bytes) + fsyncSync(handle) + } finally { + closeSync(handle) + } + try { + renameSync(staging, file) + } catch (error) { + rmSync(staging, { force: true }) + throw error + } +} + +/** Oldest-first, down to `budget`. Retained payloads are recoverable evidence, + * not the timeline: shedding the oldest is how the sidecar stays bounded. */ +function evictJournalOverflow(directory: string, budget: number): void { + const entries: { path: string; bytes: number; modifiedAt: number }[] = [] + let total = 0 + for (const name of readdirSync(directory)) { + const path = join(directory, name) + try { + const stats = statSync(path) + if (!stats.isFile()) { + continue + } + entries.push({ path, bytes: stats.size, modifiedAt: stats.mtimeMs }) + total += stats.size + } catch { + // Raced with another eviction; it is already gone. + } + } + if (total <= budget) { + return + } + entries.sort((a, b) => a.modifiedAt - b.modifiedAt) + for (const entry of entries) { + if (total <= budget) { + return + } + try { + rmSync(entry.path, { force: true }) + total -= entry.bytes + } catch { + // Left in place; the next retention pass tries again. + } + } +} diff --git a/src/main/native-chat/agent-session-journal/journal-payload-bounds.ts b/src/main/native-chat/agent-session-journal/journal-payload-bounds.ts index ea09ecf6df3..afc0307df65 100644 --- a/src/main/native-chat/agent-session-journal/journal-payload-bounds.ts +++ b/src/main/native-chat/agent-session-journal/journal-payload-bounds.ts @@ -2,19 +2,37 @@ // // A 40 MB tool result must not be inlined into a row that every reconnecting // client replays. A bounded payload keeps a head plus the original byte length -// and digest; the remainder is discarded. Crossing a bound is always marked — -// never a silent drop. +// and digest, and the remainder goes to the session's overflow store BEFORE the +// clip — so the bound is what a client renders, never what survives on the host. +// +// The sink is part of `JournalPayloadLimits` rather than a default, because the +// question "where does the remainder go?" has to be answered at the call site. +// A bound applied to a derived identifier answers it with +// `UNRETAINED_JOURNAL_PAYLOAD_LIMITS`, which is a visible, reviewable choice. import { createHash } from 'node:crypto' import type { AgentJournalBoundedPayload } from '../../../shared/agent-session-journal-types' +import { JOURNAL_OVERFLOW_NOT_RETAINED, type JournalOverflowSink } from './journal-overflow-store' export type JournalPayloadLimits = { /** Bytes of the payload kept inline on the row. */ inlineHeadBytes: number + /** Where the clipped remainder is retained. */ + overflow: JournalOverflowSink } -export const DEFAULT_JOURNAL_PAYLOAD_LIMITS: JournalPayloadLimits = { - inlineHeadBytes: 16 * 1024 +export const DEFAULT_JOURNAL_INLINE_HEAD_BYTES = 16 * 1024 + +/** The row budget, retaining everything it clips into `overflow`. */ +export function journalPayloadLimits(overflow: JournalOverflowSink): JournalPayloadLimits { + return { inlineHeadBytes: DEFAULT_JOURNAL_INLINE_HEAD_BYTES, overflow } +} + +/** For bounds whose clipped tail is a derived key fragment — a prompt option id, + * a turn ordinal map, an image reference — and not user data. */ +export const UNRETAINED_JOURNAL_PAYLOAD_LIMITS: JournalPayloadLimits = { + inlineHeadBytes: DEFAULT_JOURNAL_INLINE_HEAD_BYTES, + overflow: JOURNAL_OVERFLOW_NOT_RETAINED } /** Marker appended to a clipped inline string so the UI never presents a @@ -28,8 +46,9 @@ export function digestPayload(payload: string): string { return createHash('sha256').update(payload, 'utf8').digest('hex') } -/** Clip `payload` to the inline head. `truncated` means the remainder was - * discarded; `digest` and `byteLength` describe the original. */ +/** Clip `payload` to the inline head, retaining the whole of it under `digest` + * first. `spilled` says the retention succeeded; without it the row is exactly + * the head, length and digest it would have carried before. */ export function boundPayload( payload: string, limits: JournalPayloadLimits @@ -39,11 +58,14 @@ export function boundPayload( if (buffer.byteLength <= limits.inlineHeadBytes) { return { head: payload, byteLength: buffer.byteLength, digest, truncated: false } } + // Retain BEFORE the clip: this is the last point that holds the whole payload. + const spilled = limits.overflow.retain(digest, payload) return { head: clipUtf8(buffer, limits.inlineHeadBytes), byteLength: buffer.byteLength, digest, - truncated: true + truncated: true, + ...(spilled ? { spilled: true as const } : {}) } } @@ -82,7 +104,8 @@ export function boundToolInput(input: unknown, limits: JournalPayloadLimits): un truncated: true, byteLength: bounded.byteLength, digest: bounded.digest, - head: bounded.head + head: bounded.head, + ...(bounded.spilled ? { spilled: true as const } : {}) } : input } diff --git a/src/main/native-chat/agent-session-journal/journal-prompt-body-bounds.ts b/src/main/native-chat/agent-session-journal/journal-prompt-body-bounds.ts index 58d12dfcf89..a234d058fa9 100644 --- a/src/main/native-chat/agent-session-journal/journal-prompt-body-bounds.ts +++ b/src/main/native-chat/agent-session-journal/journal-prompt-body-bounds.ts @@ -7,12 +7,17 @@ import type { import { boundInlineText, boundPayload, - DEFAULT_JOURNAL_PAYLOAD_LIMITS + UNRETAINED_JOURNAL_PAYLOAD_LIMITS } from './journal-payload-bounds' +// Everything here retains nothing, and both cases are bounded already: +// `boundJournalStatusText` bounds Orca's own short status sentences, and +// `cancelledJournalPromptBody` re-bounds a body the journal bounded on the way +// in. Neither is the last holder of any bytes. + export const MAX_JOURNAL_PROMPT_OPTIONS = 64 -const JOURNAL_PROMPT_OPTION_LIMITS = { inlineHeadBytes: 1024 } +const JOURNAL_PROMPT_OPTION_LIMITS = { ...UNRETAINED_JOURNAL_PAYLOAD_LIMITS, inlineHeadBytes: 1024 } const JOURNAL_PROMPT_ID_MAX_BYTES = 1024 export function cancelledJournalPromptBody( @@ -34,7 +39,7 @@ export function cancelledJournalPromptBody( } export function boundJournalStatusText(text: string): string { - return boundInlineText(text, DEFAULT_JOURNAL_PAYLOAD_LIMITS).text + return boundInlineText(text, UNRETAINED_JOURNAL_PAYLOAD_LIMITS).text } function boundJournalPromptBody( @@ -68,13 +73,16 @@ function boundPromptOptions( } function boundPromptText(value: string): string { - return boundInlineText(value, DEFAULT_JOURNAL_PAYLOAD_LIMITS).text + return boundInlineText(value, UNRETAINED_JOURNAL_PAYLOAD_LIMITS).text } function boundPromptIdentifier(value: string): string { if (Buffer.byteLength(value, 'utf8') <= JOURNAL_PROMPT_ID_MAX_BYTES) { return value } - const bounded = boundPayload(value, { inlineHeadBytes: JOURNAL_PROMPT_ID_MAX_BYTES - 33 }) + const bounded = boundPayload(value, { + ...UNRETAINED_JOURNAL_PAYLOAD_LIMITS, + inlineHeadBytes: JOURNAL_PROMPT_ID_MAX_BYTES - 33 + }) return `${bounded.head}#${bounded.digest.slice(0, 32)}` } diff --git a/src/main/native-chat/agent-session-journal/journal-repair-marker.ts b/src/main/native-chat/agent-session-journal/journal-repair-marker.ts index f9ef09e688d..5433b0e1a5e 100644 --- a/src/main/native-chat/agent-session-journal/journal-repair-marker.ts +++ b/src/main/native-chat/agent-session-journal/journal-repair-marker.ts @@ -14,7 +14,6 @@ // sequence owns the epoch and stops the retry. import type Database from '../../sqlite/sync-database' -import { deleteJournalRowSuffix } from './journal-row-table' const SELECT_REPAIR = 'SELECT epoch, content_from FROM journal_repairs WHERE session_id = ?' const UPSERT_REPAIR = `INSERT INTO journal_repairs (session_id, epoch, content_from, repaired_at) @@ -45,25 +44,14 @@ export function clearJournalRepairMarker(db: Database.Database, sessionId: strin db.prepare(DELETE_REPAIR).run(sessionId) } -/** Drop the rejected suffix and record that it is owed, atomically. */ -export function deleteJournalRepairedSuffix(input: { - db: Database.Database - sessionId: string - epoch: string - /** First sequence of the rejected suffix. */ - fromSeq: number - /** First sequence left free once the suffix is gone. */ - contentFrom: number +/** Raise the marker inside the caller's repair transaction, so the rows and the + * record that they are owed commit together. */ +export function writeJournalRepairMarker( + db: Database.Database, + sessionId: string, + epoch: string, + contentFrom: number, now: number -}): number { - input.db.exec('BEGIN IMMEDIATE') - try { - const deleted = deleteJournalRowSuffix(input.db, input.sessionId, input.epoch, input.fromSeq) - input.db.prepare(UPSERT_REPAIR).run(input.sessionId, input.epoch, input.contentFrom, input.now) - input.db.exec('COMMIT') - return deleted - } catch (error) { - input.db.exec('ROLLBACK') - throw error - } +): void { + db.prepare(UPSERT_REPAIR).run(sessionId, epoch, contentFrom, now) } diff --git a/src/main/native-chat/agent-session-journal/journal-repair-suffix.ts b/src/main/native-chat/agent-session-journal/journal-repair-suffix.ts new file mode 100644 index 00000000000..71843fbd7b4 --- /dev/null +++ b/src/main/native-chat/agent-session-journal/journal-repair-suffix.ts @@ -0,0 +1,53 @@ +// Dropping the rejected suffix a replay found — after it is durable elsewhere. +// +// Ordering is the whole contract: quarantine, fsync, THEN the transaction that +// deletes the rows and records the rebuild they are owed. A repair that cannot +// write the copy does not delete; it reports `quarantined: false` and the caller +// latches read-only, which leaves the timeline exactly as it was found. + +import type Database from '../../sqlite/sync-database' +import { quarantineJournalRows } from './journal-row-quarantine' +import { writeJournalRepairMarker } from './journal-repair-marker' +import { deleteJournalRowSuffix, readJournalRowsAfter } from './journal-row-table' + +export type JournalRepairedSuffix = + | { quarantined: true; quarantineFile: string; deleted: number } + /** The rows are still live and untouched. */ + | { quarantined: false; error: unknown } + +export function quarantineAndDeleteJournalRepairedSuffix(input: { + db: Database.Database + journalDir: string + sessionId: string + epoch: string + /** First sequence of the rejected suffix. */ + fromSeq: number + /** First sequence left free once the suffix is gone. */ + contentFrom: number + now: number +}): JournalRepairedSuffix { + let quarantineFile: string + try { + quarantineFile = quarantineJournalRows({ + journalDir: input.journalDir, + epoch: input.epoch, + fromSeq: input.fromSeq, + now: input.now, + // `> fromSeq - 1` is `>= fromSeq`: the rejected row itself and everything + // behind it, which is exactly what the delete below removes. + rows: readJournalRowsAfter(input.db, input.sessionId, input.epoch, input.fromSeq - 1) + }) + } catch (error) { + return { quarantined: false, error } + } + input.db.exec('BEGIN IMMEDIATE') + try { + const deleted = deleteJournalRowSuffix(input.db, input.sessionId, input.epoch, input.fromSeq) + writeJournalRepairMarker(input.db, input.sessionId, input.epoch, input.contentFrom, input.now) + input.db.exec('COMMIT') + return { quarantined: true, quarantineFile, deleted } + } catch (error) { + input.db.exec('ROLLBACK') + throw error + } +} diff --git a/src/main/native-chat/agent-session-journal/journal-row-quarantine.test.ts b/src/main/native-chat/agent-session-journal/journal-row-quarantine.test.ts new file mode 100644 index 00000000000..cdc88556237 --- /dev/null +++ b/src/main/native-chat/agent-session-journal/journal-row-quarantine.test.ts @@ -0,0 +1,175 @@ +// A repair copies the rows it rejects before it drops them. +// +// The rejected row and every valid row behind it are the only copy of what the +// session wrote after the fault — a submission receipt among them, which no +// provider transcript can rebuild. These assert the copy exists, byte for byte, +// and that a repair which cannot write one does not delete anything. + +import { chmod, mkdir, mkdtemp, readdir, readFile, rm } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { afterEach, beforeEach, expect, it } from 'vitest' +import type { + AgentJournalItemBody, + AgentJournalItemIdentity, + AgentSessionJournalIdentity +} from '../../../shared/agent-session-journal-types' +import type Database from '../../sqlite/sync-database' +import { openJournalDatabase } from './journal-database' +import { journalDatabaseFile } from './journal-paths' +import { journalQuarantineDirectory } from './journal-row-quarantine' +import { createTrackedJournalOpener } from './journal-store-test-open' + +const IDENTITY: AgentSessionJournalIdentity = { + sessionId: 'session-1', + workspaceId: 'ws-1', + hostId: 'host-1', + agent: 'codex', + providerHandle: { kind: 'codex', threadId: 'thread-1' } +} + +const CORRUPT_ROW = '}{ not a row λ' + +let root: string +let clock = 1_000 +const journals = createTrackedJournalOpener() + +function tick(): number { + clock += 1 + return clock +} + +function item(ordinal: number): AgentJournalItemIdentity { + return { provider: 'codex', threadId: 'thread-1', turnId: 'turn-1', ordinal } +} + +function body(value: string): AgentJournalItemBody { + return { kind: 'message', role: 'assistant', blocks: [{ type: 'text', text: value }] } +} + +function open() { + return journals.open({ + identity: IDENTITY, + journalDir: root, + now: tick, + mintEpoch: () => `epoch-${clock}` + }) +} + +async function withJournalDatabase(run: (db: Database.Database) => T): Promise { + const opened = openJournalDatabase(journalDatabaseFile(root)) + try { + return run(opened.db) + } finally { + opened.db.close() + } +} + +function storedRows(db: Database.Database): { seq: number; row_json: string }[] { + return db.prepare('SELECT seq, row_json FROM journal_rows ORDER BY seq').all() as { + seq: number + row_json: string + }[] +} + +type QuarantinedRow = { epoch: string; seq: number; ts: number; rowJson: string } + +async function quarantinedRows(): Promise { + const directory = journalQuarantineDirectory(root) + const files = await readdir(directory) + expect(files).toHaveLength(1) + const contents = await readFile(join(directory, files[0]!), 'utf8') + return contents + .split('\n') + .filter((line) => line.length > 0) + .map((line) => JSON.parse(line) as QuarantinedRow) +} + +/** A prefix, a row this build cannot parse, and valid rows behind it — one of + * them an Orca-minted submission receipt. */ +async function seedCorruptedJournal(): Promise { + const journal = await open() + await journal.appendItem(item(0), body('readable prefix'), { fence: 1 }) + await journal.appendItem(item(1), body('the fault'), { fence: 1 }) + await journal.appendSubmission({ + clientMessageId: 'orca-only-receipt', + payloadFingerprint: 'fingerprint', + fence: 1, + body: { kind: 'message', role: 'user', blocks: [{ type: 'text', text: 'only in Orca' }] } + }) + await journal.appendItem(item(2), body('valid suffix'), { fence: 1 }) + await journal.close() + // Sequence 1 is the epoch row, so the fault lands at 3 and 4..5 are valid. + await withJournalDatabase((db) => { + db.prepare('UPDATE journal_rows SET row_json = ? WHERE seq = ?').run(CORRUPT_ROW, 3) + }) +} + +beforeEach(async () => { + root = await mkdtemp(join(tmpdir(), 'orca-journal-quarantine-')) + clock = 1_000 +}) + +afterEach(async () => { + await journals.closeAll() + await chmod(journalQuarantineDirectory(root), 0o755).catch(() => undefined) + await rm(root, { recursive: true, force: true }) +}) + +it('quarantines the rejected row and every valid row behind it, byte for byte', async () => { + await seedCorruptedJournal() + const before = await withJournalDatabase(storedRows) + + const reopened = await open() + expect(reopened.repair.malformedRows).toBe(1) + + const quarantined = await quarantinedRows() + expect(quarantined.map((row) => row.seq)).toEqual([3, 4, 5]) + // The exact stored bytes, including the row that is not valid JSON. + expect(quarantined.map((row) => row.rowJson)).toEqual( + before.filter((row) => row.seq >= 3).map((row) => row.row_json) + ) + expect(quarantined[0]?.rowJson).toBe(CORRUPT_ROW) + // The Orca-only submission receipt no provider transcript can rebuild. + expect(quarantined.some((row) => row.rowJson.includes('orca-only-receipt'))).toBe(true) + const sourceEpoch = before[0]?.row_json ?? '' + expect(quarantined.every((row) => sourceEpoch.includes(row.epoch))).toBe(true) +}) + +it('drops nothing and latches read-only when the copy cannot be written', async () => { + await seedCorruptedJournal() + const before = await withJournalDatabase(storedRows) + // A directory the store cannot create its file in. + await mkdir(journalQuarantineDirectory(root), { recursive: true }) + await chmod(journalQuarantineDirectory(root), 0o500) + + const reopened = await open() + + expect(reopened.isReadOnly).toBe(true) + await expect(reopened.appendItem(item(9), body('refused'))).rejects.toMatchObject({ + code: 'journal_read_only' + }) + // Every row the repair would have deleted is still there, unchanged. + expect(await withJournalDatabase(storedRows)).toEqual(before) +}) + +it('quarantines the rows a sequence gap orphans, which carry no malformed row', async () => { + const journal = await open() + for (let ordinal = 0; ordinal < 5; ordinal += 1) { + await journal.appendItem(item(ordinal), body(`m${ordinal}`), { fence: 1 }) + } + await journal.close() + // Items occupy 2..6; removing 4 leaves 5 and 6 valid but unanchored. + await withJournalDatabase((db) => { + db.prepare('DELETE FROM journal_rows WHERE seq = ?').run(4) + }) + const before = await withJournalDatabase(storedRows) + + await open() + + const quarantined = await quarantinedRows() + expect(quarantined.map((row) => row.seq)).toEqual([5, 6]) + expect(quarantined.map((row) => row.rowJson)).toEqual( + before.filter((row) => row.seq >= 5).map((row) => row.row_json) + ) +}) diff --git a/src/main/native-chat/agent-session-journal/journal-row-quarantine.ts b/src/main/native-chat/agent-session-journal/journal-row-quarantine.ts new file mode 100644 index 00000000000..2ab8c833c11 --- /dev/null +++ b/src/main/native-chat/agent-session-journal/journal-row-quarantine.ts @@ -0,0 +1,58 @@ +// The rejected rows a repair is about to drop, copied out first. +// +// Replay stops at the first row this build cannot represent, and the repair then +// clears from there to the tip. Those rows are the ONLY copy of everything the +// session wrote after the fault, Orca-only submission receipts included, and no +// provider transcript can rebuild those. They are written here with their exact +// stored bytes BEFORE the delete transaction opens: a crash between the two +// leaves the rows still live plus an orphan quarantine file, never the reverse. +// +// A failed quarantine is a refusal to delete, not a reason to proceed — the +// caller latches the journal read-only and the next open retries. + +import { closeSync, fsyncSync, mkdirSync, openSync, writeSync } from 'node:fs' +import { join } from 'node:path' +import type { JournalStoredRow } from './journal-row-table' + +const QUARANTINE_DIR_NAME = 'quarantine' + +export function journalQuarantineDirectory(journalDir: string): string { + return join(journalDir, QUARANTINE_DIR_NAME) +} + +/** + * Copy `rows` to a durable file and return its path. Throws when the bytes are + * not durable, so no caller can read a path here as permission to delete. + * + * The stored `row_json` is carried as a JSON string rather than embedded JSON: + * a rejected row is very often not valid JSON, which is the whole reason it is + * being quarantined. + */ +export function quarantineJournalRows(input: { + journalDir: string + epoch: string + fromSeq: number + now: number + rows: readonly JournalStoredRow[] +}): string { + const directory = journalQuarantineDirectory(input.journalDir) + mkdirSync(directory, { recursive: true }) + const file = join(directory, quarantineFileName(input.epoch, input.fromSeq, input.now)) + const lines = input.rows.map((row) => + JSON.stringify({ epoch: row.epoch, seq: row.seq, ts: row.ts, rowJson: row.rowJson }) + ) + const handle = openSync(file, 'w') + try { + writeSync(handle, Buffer.from(`${lines.join('\n')}\n`, 'utf8')) + fsyncSync(handle) + } finally { + closeSync(handle) + } + return file +} + +/** Epoch first so one session's repairs group together, and the rejected + * sequence plus the clock keep repeated repairs of one epoch distinct. */ +function quarantineFileName(epoch: string, fromSeq: number, now: number): string { + return `${epoch}-${fromSeq}-${now}.jsonl` +} diff --git a/src/main/native-chat/agent-session-journal/journal-store-open.ts b/src/main/native-chat/agent-session-journal/journal-store-open.ts index 7b5b6d0dff8..b27ab7f5a76 100644 --- a/src/main/native-chat/agent-session-journal/journal-store-open.ts +++ b/src/main/native-chat/agent-session-journal/journal-store-open.ts @@ -9,6 +9,7 @@ import { journalFileFormatRemnantDisclosure } from './journal-file-format-remnant' import type { JournalLoad } from './journal-open' +import type { JournalRepairedSuffix } from './journal-repair-suffix' import { journalRepairDisclosure, type JournalRepairDisclosure } from './journal-repair-disclosure' import { staleSubagentRosterRevisions } from './journal-subagent-liveness' @@ -32,11 +33,13 @@ export async function openJournalStoreState(input: { journalDir: string loaded: JournalLoad | null | undefined replay: () => JournalLoad | null - /** Drops the rejected suffix and records the rebuild it owes, in ONE - * transaction. Corruption is not preserved; replay keeps reporting `corrupt` + /** Copies the rejected suffix somewhere durable, then drops it and records + * the rebuild it owes in ONE transaction. Replay keeps reporting `corrupt` * until provider history republishes the epoch or the session writes past * `contentFrom`, the first sequence the repair left free. */ - deleteSuffix: (fromSeq: number, contentFrom: number) => number + repairSuffix: (fromSeq: number, contentFrom: number) => JournalRepairedSuffix + /** Latches the journal against writes for the rest of its life. */ + latchReadOnly: () => void start: () => void adopt: (loaded: JournalLoad) => void /** Republishes an anchor row for an epoch a repair emptied. */ @@ -60,7 +63,13 @@ export async function openJournalStoreState(input: { } input.adopt(loaded) if (loaded.truncateFrom !== undefined && !loaded.readOnly) { - input.deleteSuffix(loaded.truncateFrom, loaded.state.lastSequence + 1) + // The rejected rows are the only copy of everything the session wrote after + // the fault. Without a durable one this open takes no repair at all: the + // journal latches read-only with its rows intact, and a later open retries. + if (!input.repairSuffix(loaded.truncateFrom, loaded.state.lastSequence + 1).quarantined) { + input.latchReadOnly() + return + } } // A repair that took every live row leaves the epoch with no anchor. Publish // one before anything can append into it: an ordinary row at sequence 1 would diff --git a/src/main/native-chat/agent-session-journal/journal-store-restore.ts b/src/main/native-chat/agent-session-journal/journal-store-restore.ts index fc69dd339d8..645e282e8a8 100644 --- a/src/main/native-chat/agent-session-journal/journal-store-restore.ts +++ b/src/main/native-chat/agent-session-journal/journal-store-restore.ts @@ -10,7 +10,7 @@ import type { JournalEpochController } from './journal-epoch-controller' import { replayJournal } from './journal-open' import type { JournalStoreHost } from './journal-store-collaborators' import { openJournalStoreState } from './journal-store-open' -import { deleteJournalRepairedSuffix } from './journal-repair-marker' +import { quarantineAndDeleteJournalRepairedSuffix } from './journal-repair-suffix' export function restoreJournalStore( host: JournalStoreHost, @@ -23,15 +23,17 @@ export function restoreJournalStore( const opened = host.database() return replayJournal(opened.db, opened.readOnly, host.identity.sessionId) }, - deleteSuffix: (fromSeq, contentFrom) => - deleteJournalRepairedSuffix({ + repairSuffix: (fromSeq, contentFrom) => + quarantineAndDeleteJournalRepairedSuffix({ db: host.database().db, + journalDir: host.journalDir, sessionId: host.identity.sessionId, epoch: host.state().epoch, fromSeq, contentFrom, now: host.now() }), + latchReadOnly: () => host.setReadOnly(true), start: () => collaborators.epochController.start('session_created', 0), // `unreconcilable_prefix` is the durable statement that this epoch exists // because a repair emptied one: replay reads it back and keeps asking for diff --git a/src/main/native-chat/agent-session-journal/journal-store.test.ts b/src/main/native-chat/agent-session-journal/journal-store.test.ts index 97eac02fe75..22d664fe433 100644 --- a/src/main/native-chat/agent-session-journal/journal-store.test.ts +++ b/src/main/native-chat/agent-session-journal/journal-store.test.ts @@ -15,7 +15,7 @@ import { loadJournal } from './journal-open' import { boundInlineText, boundPayload, - DEFAULT_JOURNAL_PAYLOAD_LIMITS + UNRETAINED_JOURNAL_PAYLOAD_LIMITS } from './journal-payload-bounds' import { journalDatabaseFile, journalDirectoryFor, journalPathSegment } from './journal-paths' import { AgentSessionJournalError, type AgentSessionJournal } from './journal-store' @@ -220,7 +220,7 @@ describe('replay', () => { describe('bounds', () => { it('marks a clipped payload instead of dropping bytes silently', () => { - const limits = { ...DEFAULT_JOURNAL_PAYLOAD_LIMITS, inlineHeadBytes: 16 } + const limits = { ...UNRETAINED_JOURNAL_PAYLOAD_LIMITS, inlineHeadBytes: 16 } const bounded = boundPayload('x'.repeat(4_096), limits) expect(bounded.truncated).toBe(true) expect(bounded.head).toHaveLength(16) @@ -229,7 +229,7 @@ describe('bounds', () => { }) it('never splits a multi-byte character across the bound', () => { - const limits = { ...DEFAULT_JOURNAL_PAYLOAD_LIMITS, inlineHeadBytes: 4 } + const limits = { ...UNRETAINED_JOURNAL_PAYLOAD_LIMITS, inlineHeadBytes: 4 } // Each character is three bytes, so a naive slice would land mid-sequence. const bounded = boundPayload('日本語テスト', limits) expect(bounded.head).toBe('日') @@ -237,10 +237,10 @@ describe('bounds', () => { }) it('leaves a payload inside the bound untouched', () => { - const bounded = boundPayload('small', DEFAULT_JOURNAL_PAYLOAD_LIMITS) + const bounded = boundPayload('small', UNRETAINED_JOURNAL_PAYLOAD_LIMITS) expect(bounded.truncated).toBe(false) expect(bounded.head).toBe('small') - expect(boundInlineText('small', DEFAULT_JOURNAL_PAYLOAD_LIMITS).text).toBe('small') + expect(boundInlineText('small', UNRETAINED_JOURNAL_PAYLOAD_LIMITS).text).toBe('small') }) }) diff --git a/src/main/native-chat/agent-session-wire/provider-frame-disposition.test.ts b/src/main/native-chat/agent-session-wire/provider-frame-disposition.test.ts index a57aa5d9ed1..87a502c50ed 100644 --- a/src/main/native-chat/agent-session-wire/provider-frame-disposition.test.ts +++ b/src/main/native-chat/agent-session-wire/provider-frame-disposition.test.ts @@ -7,6 +7,7 @@ import { PROVIDER_FRAME_CLASSIFICATIONS } from './provider-frame-disposition' import { unhandledProviderFrameJournalItem } from './unhandled-provider-frame' +import { UNRETAINED_JOURNAL_PAYLOAD_LIMITS } from '../agent-session-journal/journal-payload-bounds' describe('provider frame classification catalog', () => { it('classifies every pinned Codex app-server notification method', () => { @@ -206,10 +207,15 @@ describe('codex subagent item disposition', () => { it('journals no fallback row for subagent activity', () => { expect( - unhandledProviderFrameJournalItem('codex', 'item:subAgentActivity', { - kind: 'completed', - agentThreadId: 'child-1' - }) + unhandledProviderFrameJournalItem( + 'codex', + 'item:subAgentActivity', + { + kind: 'completed', + agentThreadId: 'child-1' + }, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) ).toBeNull() }) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-event-sink-queue.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-event-sink-queue.ts index 510c1df2f4a..51a196ae685 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-event-sink-queue.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-event-sink-queue.ts @@ -7,6 +7,12 @@ import type { StructuredAgentSessionSinkState, StructuredAgentSessionSinkWatermarks } from './structured-agent-session-event-sink' +import { + journalPayloadLimits, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS, + type JournalPayloadLimits +} from '../agent-session-journal/journal-payload-bounds' +import { journalOverflowSink } from '../agent-session-journal/journal-overflow-store' export type StructuredAgentSessionSinkOperation = { sequence: number @@ -26,6 +32,7 @@ export type StructuredAgentSessionDrainWaiter = { export class StructuredAgentSessionSinkQueue { private readingControl: StructuredAgentSessionReadingControl | undefined private target: StructuredAgentSessionEventTarget | null = null + private limits: JournalPayloadLimits | null = null private closed = false private failure: { error: unknown } | null = null private running = false @@ -75,12 +82,26 @@ export class StructuredAgentSessionSinkQueue { bind(target: StructuredAgentSessionEventTarget): void { if (!this.closed) { this.target = target + this.limits = null this.pump() } } unbind(): void { this.target = null + this.limits = null + } + + /** The row budget for the bound journal, so a clipped payload's remainder is + * retained beside that journal rather than dropped. Unretained while nothing + * is bound: there is no session directory to retain into yet. */ + payloadLimits = (): JournalPayloadLimits => { + const journalDir = this.target?.journal.directory + if (!journalDir) { + return UNRETAINED_JOURNAL_PAYLOAD_LIMITS + } + this.limits ??= journalPayloadLimits(journalOverflowSink(journalDir)) + return this.limits } close(): void { diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-event-sink.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-event-sink.ts index 6952b6d93e6..381218372c0 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-event-sink.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-event-sink.ts @@ -6,6 +6,10 @@ import type { import type { AgentSessionTurnActivity } from '../../../shared/agent-session-wire' import type { AgentSessionJournal } from '../agent-session-journal/journal-store' import type { JournalLifecycleMutationInput } from '../agent-session-journal/journal-row-builders' +import { + UNRETAINED_JOURNAL_PAYLOAD_LIMITS, + type JournalPayloadLimits +} from '../agent-session-journal/journal-payload-bounds' import { estimateStructuredAgentSessionItemBytes } from './structured-agent-session-event-sink-estimate' import { StructuredAgentSessionSinkQueue } from './structured-agent-session-event-sink-queue' @@ -63,6 +67,18 @@ export type StructuredAgentSessionEventSink = { tryPublish?(options?: StructuredAgentSessionAppendOptions): StructuredAgentSessionSinkAdmission /** Couples durable-queue pressure to the exact provider stream producing it. */ bindReadingControl?(control: StructuredAgentSessionReadingControl): () => void + /** Row payload budget for the bound session, which retains what it clips into + * that session's overflow store. Absent on a sink with no journal behind it. */ + payloadLimits?(): JournalPayloadLimits +} + +/** The budget a translator bounds with. A sink that cannot name a journal — a + * test double, or one still unbound — retains nothing, which is exactly what a + * row written before overflow retention existed carried. */ +export function structuredAgentSessionPayloadLimits( + sink: Pick +): JournalPayloadLimits { + return sink.payloadLimits?.() ?? UNRETAINED_JOURNAL_PAYLOAD_LIMITS } export type StructuredAgentSessionEventTarget = { @@ -218,7 +234,8 @@ export function createDeferredStructuredAgentSessionEventSink( run: (bound) => bound.publish(activity) }) }, - tryPublish: publish + tryPublish: publish, + payloadLimits: queue.payloadLimits }, bind: (next) => queue.bind(next), unbind: () => queue.unbind(), diff --git a/src/main/native-chat/agent-session-wire/unhandled-provider-frame.test.ts b/src/main/native-chat/agent-session-wire/unhandled-provider-frame.test.ts index 61544fd56ab..ee9aa7da69b 100644 --- a/src/main/native-chat/agent-session-wire/unhandled-provider-frame.test.ts +++ b/src/main/native-chat/agent-session-wire/unhandled-provider-frame.test.ts @@ -1,6 +1,7 @@ import { describe, expect, it } from 'vitest' import { projectStructuredItemToNativeChat } from '../../../shared/structured-agent-session-projection' import { unhandledProviderFrameJournalItem } from './unhandled-provider-frame' +import { UNRETAINED_JOURNAL_PAYLOAD_LIMITS } from '../agent-session-journal/journal-payload-bounds' describe('unhandled provider frame journal fallback', () => { it('keeps a compact label and bounds the expandable payload without dropping it', () => { @@ -8,7 +9,7 @@ describe('unhandled provider frame journal fallback', () => { 'future-provider', 'notification:new/event', { body: 'abcdefghij' }, - { inlineHeadBytes: 8 } + { ...UNRETAINED_JOURNAL_PAYLOAD_LIMITS, inlineHeadBytes: 8 } ) expect(item).not.toBeNull() @@ -33,54 +34,113 @@ describe('unhandled provider frame journal fallback', () => { const cyclic: { warning?: unknown } = {} cyclic.warning = cyclic - const item = unhandledProviderFrameJournalItem('codex', 'frame', cyclic) + const item = unhandledProviderFrameJournalItem( + 'codex', + 'frame', + cyclic, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) expect(item?.body.text).toBe('codex · frame') expect(item?.body.providerFrame?.payload.head).toContain('unserializable payload') }) it('routes provider lifecycle, startup, and status frames away from the timeline', () => { - expect(unhandledProviderFrameJournalItem('codex', 'notification:thread/started', {})).toBeNull() expect( - unhandledProviderFrameJournalItem('codex', 'notification:mcpServer/startupStatus/updated', {}) + unhandledProviderFrameJournalItem( + 'codex', + 'notification:thread/started', + {}, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) ).toBeNull() expect( - unhandledProviderFrameJournalItem('codex', 'notification:remoteControl/status/changed', {}) + unhandledProviderFrameJournalItem( + 'codex', + 'notification:mcpServer/startupStatus/updated', + {}, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) ).toBeNull() expect( - unhandledProviderFrameJournalItem('codex', 'notification:thread/tokenUsage/updated', {}) + unhandledProviderFrameJournalItem( + 'codex', + 'notification:remoteControl/status/changed', + {}, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) ).toBeNull() expect( - unhandledProviderFrameJournalItem('codex', 'notification:thread/goal/cleared', {}) + unhandledProviderFrameJournalItem( + 'codex', + 'notification:thread/tokenUsage/updated', + {}, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) ).toBeNull() - expect(unhandledProviderFrameJournalItem('claude', 'message:system:init', {})).toBeNull() expect( - unhandledProviderFrameJournalItem('claude', 'message:result', { - subtype: 'success', - is_error: false - }) + unhandledProviderFrameJournalItem( + 'codex', + 'notification:thread/goal/cleared', + {}, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) + ).toBeNull() + expect( + unhandledProviderFrameJournalItem( + 'claude', + 'message:system:init', + {}, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) + ).toBeNull() + expect( + unhandledProviderFrameJournalItem( + 'claude', + 'message:result', + { + subtype: 'success', + is_error: false + }, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) ).toBeNull() }) it('never creates generic rows for delta-shaped frames that report no failure', () => { expect( - unhandledProviderFrameJournalItem('codex', 'notification:item/commandExecution/outputDelta', { - itemId: 'exec-1', - delta: 'x' - }) + unhandledProviderFrameJournalItem( + 'codex', + 'notification:item/commandExecution/outputDelta', + { + itemId: 'exec-1', + delta: 'x' + }, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) ).toBeNull() expect( - unhandledProviderFrameJournalItem('codex', 'notification:item/future/outputDelta', { - itemId: 'future-1', - delta: 'y' - }) + unhandledProviderFrameJournalItem( + 'codex', + 'notification:item/future/outputDelta', + { + itemId: 'future-1', + delta: 'y' + }, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) ).toBeNull() }) it('surfaces an unknown delta-shaped frame whose payload reports an error', () => { - const row = unhandledProviderFrameJournalItem('codex', 'notification:item/future/outputDelta', { - error: 'stream broke mid-item' - }) + const row = unhandledProviderFrameJournalItem( + 'codex', + 'notification:item/future/outputDelta', + { + error: 'stream broke mid-item' + }, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) expect(row).not.toBeNull() expect(row?.classification).toBe('error-surface') @@ -91,15 +151,25 @@ describe('unhandled provider frame journal fallback', () => { }) it('renders codex systemError and Claude error result variants', () => { - const codex = unhandledProviderFrameJournalItem('codex', 'notification:thread/status/changed', { - threadId: 'thread-1', - status: { type: 'systemError' } - }) - const claude = unhandledProviderFrameJournalItem('claude', 'message:result', { - subtype: 'error_during_execution', - is_error: true, - result: 'Provider request failed' - }) + const codex = unhandledProviderFrameJournalItem( + 'codex', + 'notification:thread/status/changed', + { + threadId: 'thread-1', + status: { type: 'systemError' } + }, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) + const claude = unhandledProviderFrameJournalItem( + 'claude', + 'message:result', + { + subtype: 'error_during_execution', + is_error: true, + result: 'Provider request failed' + }, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) expect(codex?.body.providerFrame).toMatchObject({ provider: 'codex', @@ -133,20 +203,30 @@ describe('unhandled provider frame journal fallback', () => { const kind = 'notification:mcpServer/startupStatus/updated' expect( - unhandledProviderFrameJournalItem('codex', kind, { - name: 'filesystem', - status: 'starting', - error: null, - failureReason: null - }) + unhandledProviderFrameJournalItem( + 'codex', + kind, + { + name: 'filesystem', + status: 'starting', + error: null, + failureReason: null + }, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) ).toBeNull() expect( - unhandledProviderFrameJournalItem('codex', kind, { - name: 'filesystem', - status: 'failed', - error: 'server exited', - failureReason: null - }) + unhandledProviderFrameJournalItem( + 'codex', + kind, + { + name: 'filesystem', + status: 'failed', + error: 'server exited', + failureReason: null + }, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) ).not.toBeNull() }) @@ -154,28 +234,55 @@ describe('unhandled provider frame journal fallback', () => { const kind = 'notification:hook/completed' expect( - unhandledProviderFrameJournalItem('codex', kind, { - run: { id: 'hook-1', status: 'completed' } - }) + unhandledProviderFrameJournalItem( + 'codex', + kind, + { + run: { id: 'hook-1', status: 'completed' } + }, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) ).toBeNull() expect( - unhandledProviderFrameJournalItem('codex', kind, { - run: { id: 'hook-1', status: 'failed' } - }) + unhandledProviderFrameJournalItem( + 'codex', + kind, + { + run: { id: 'hook-1', status: 'failed' } + }, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) ).not.toBeNull() }) it('keeps unknown substantive frames visible for both providers', () => { expect( - unhandledProviderFrameJournalItem('codex', 'notification:future/event', {}) + unhandledProviderFrameJournalItem( + 'codex', + 'notification:future/event', + {}, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) + ).not.toBeNull() + expect( + unhandledProviderFrameJournalItem( + 'claude', + 'message:future/event', + {}, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) ).not.toBeNull() - expect(unhandledProviderFrameJournalItem('claude', 'message:future/event', {})).not.toBeNull() }) it('leads with the provider sentence instead of naming the opcode', () => { - const row = unhandledProviderFrameJournalItem('codex', 'notification:warning', { - message: 'Your plan limit resets in 2 hours.' - }) + const row = unhandledProviderFrameJournalItem( + 'codex', + 'notification:warning', + { + message: 'Your plan limit resets in 2 hours.' + }, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) expect(row?.body.text).toBe('Your plan limit resets in 2 hours.') // The raw frame stays available behind the row's disclosure. expect(row?.body.providerFrame?.kind).toBe('notification:warning') @@ -187,7 +294,7 @@ describe('unhandled provider frame journal fallback', () => { 'codex', 'notification:warning', { message }, - { inlineHeadBytes: 8 } + { ...UNRETAINED_JOURNAL_PAYLOAD_LIMITS, inlineHeadBytes: 8 } ) expect(row?.body.text).toContain('abcdefgh') @@ -196,13 +303,22 @@ describe('unhandled provider frame journal fallback', () => { it('unwraps a nested sentence and falls back to the opcode when there is none', () => { expect( - unhandledProviderFrameJournalItem('codex', 'notification:warning', { - warning: { text: 'Sandbox is degraded.' } - })?.body.text + unhandledProviderFrameJournalItem( + 'codex', + 'notification:warning', + { + warning: { text: 'Sandbox is degraded.' } + }, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + )?.body.text ).toBe('Sandbox is degraded.') expect( - unhandledProviderFrameJournalItem('codex', 'notification:future/event', { count: 3 })?.body - .text + unhandledProviderFrameJournalItem( + 'codex', + 'notification:future/event', + { count: 3 }, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + )?.body.text ).toBe('codex \u00b7 notification:future/event') }) }) @@ -218,7 +334,8 @@ describe('a failed provider dependency', () => { status: 'failed', error: 'MCP client for `codex_apps` failed to start: authentication token invalidated', failureReason: 'reauthenticationRequired' - } + }, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS ) expect(item?.classification).toBe('error-surface') expect(item?.body.text).toContain('failed to start') @@ -227,11 +344,16 @@ describe('a failed provider dependency', () => { it('stays out of the timeline while the dependency is merely starting', () => { expect( - unhandledProviderFrameJournalItem('codex', 'notification:mcpServer/startupStatus/updated', { - threadId: 'thread-1', - name: 'codex_apps', - status: 'starting' - }) + unhandledProviderFrameJournalItem( + 'codex', + 'notification:mcpServer/startupStatus/updated', + { + threadId: 'thread-1', + name: 'codex_apps', + status: 'starting' + }, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) ).toBeNull() }) }) @@ -239,12 +361,24 @@ describe('a failed provider dependency', () => { describe('typed notice metadata', () => { it('publishes readable compaction statuses for both provider forms', () => { expect( - unhandledProviderFrameJournalItem('codex', 'notification:thread/compacted', {}) + unhandledProviderFrameJournalItem( + 'codex', + 'notification:thread/compacted', + {}, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) ).toMatchObject({ classification: 'timeline-substantive', body: { kind: 'status', text: 'Context compacted', presentation: 'compaction' } }) - expect(unhandledProviderFrameJournalItem('codex', 'item:contextCompaction', {})).toMatchObject({ + expect( + unhandledProviderFrameJournalItem( + 'codex', + 'item:contextCompaction', + {}, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) + ).toMatchObject({ body: { kind: 'status', text: 'Context compacted', presentation: 'compaction' } }) }) @@ -266,7 +400,12 @@ describe('typed notice metadata', () => { ['error', { error: { message: 'Connection failed' } }, 'error', 'Connection failed'] ])('assigns the tone and readable text for %s', (method, payload, tone, text) => { expect( - unhandledProviderFrameJournalItem('codex', `notification:${method}`, payload) + unhandledProviderFrameJournalItem( + 'codex', + `notification:${method}`, + payload, + UNRETAINED_JOURNAL_PAYLOAD_LIMITS + ) ).toMatchObject({ classification: 'error-surface', body: { kind: 'status', text, tone } diff --git a/src/main/native-chat/agent-session-wire/unhandled-provider-frame.ts b/src/main/native-chat/agent-session-wire/unhandled-provider-frame.ts index 272e960b838..1b240d1355c 100644 --- a/src/main/native-chat/agent-session-wire/unhandled-provider-frame.ts +++ b/src/main/native-chat/agent-session-wire/unhandled-provider-frame.ts @@ -2,7 +2,6 @@ import type { AgentJournalStatusItem } from '../../../shared/agent-session-journ import { boundInlineText, boundPayload, - DEFAULT_JOURNAL_PAYLOAD_LIMITS, type JournalPayloadLimits } from '../agent-session-journal/journal-payload-bounds' import { classifyProviderFrame } from './provider-frame-disposition' @@ -75,7 +74,7 @@ export function unhandledProviderFrameJournalItem( provider: string, kind: string, payload: unknown, - limits: JournalPayloadLimits = DEFAULT_JOURNAL_PAYLOAD_LIMITS + limits: JournalPayloadLimits ): UnhandledProviderFrameJournalItem | null { const classification = classifyProviderFrame(provider, kind, payload) if ( diff --git a/src/shared/agent-session-journal-schemas.ts b/src/shared/agent-session-journal-schemas.ts index 89c2d9f7293..a80f8112c9f 100644 --- a/src/shared/agent-session-journal-schemas.ts +++ b/src/shared/agent-session-journal-schemas.ts @@ -25,7 +25,8 @@ const BoundedPayload = z.object({ head: z.string(), byteLength: z.number(), digest: z.string(), - truncated: z.boolean() + truncated: z.boolean(), + spilled: z.literal(true).optional() }) const ProviderFrame = z.object({ diff --git a/src/shared/agent-session-journal-types.ts b/src/shared/agent-session-journal-types.ts index c0074e015f5..1c5cad61996 100644 --- a/src/shared/agent-session-journal-types.ts +++ b/src/shared/agent-session-journal-types.ts @@ -59,17 +59,22 @@ export type AgentJournalItemIdentity = // ─── Bounded payloads ─────────────────────────────────────────────────────── -/** A tool output or diff body clipped to a head. The remainder is DISCARDED, - * never stored: crossing a bound sets `truncated` and the two fields below - * describe what was dropped, so it is marked rather than silently lost. */ +/** A tool output or diff body clipped to a head. `truncated` marks the clip and + * the two fields below describe the original, so a bounded body is never a + * silently short one. */ export type AgentJournalBoundedPayload = { head: string /** Byte length of the ORIGINAL payload, not of `head`. */ byteLength: number - /** sha256 of the original payload — identification only; nothing stores or - * retrieves the discarded remainder by it. */ + /** sha256 of the original payload. Also the key the full bytes are retained + * under on the execution host when `spilled` is set. */ digest: string truncated: boolean + /** The full payload was retained on the execution host under `digest` before + * the clip. Absent on a row written by a build that could not retain it, and + * on one whose retention failed — those still carry head, length and digest. + * Additive: an older client ignores it and renders the head as it always did. */ + spilled?: true } // ─── Render-model items ───────────────────────────────────────────────────── diff --git a/src/shared/native-chat-types.ts b/src/shared/native-chat-types.ts index 3a018d75e6b..3848523b4c1 100644 --- a/src/shared/native-chat-types.ts +++ b/src/shared/native-chat-types.ts @@ -44,6 +44,7 @@ export type NativeChatTextBlock = { byteLength: number digest: string truncated: boolean + spilled?: true } } }