Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 5 additions & 2 deletions src/main/codex/codex-persistent-command-retention.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,7 @@ function fixture(maxMetadataBytes?: number) {
{
sink,
maxMetadataBytes,
linkageFor: () => ({}),
schedule: (run) => {
scheduled.add(run)
return () => {
Expand Down Expand Up @@ -116,7 +117,8 @@ describe('persistent command retention', () => {
turnLifecycle: null,
sink,
streams: items.streams,
activeItems: items.activeItems
activeItems: items.activeItems,
linkageFor: () => ({})
})
).toEqual({ accepted: true })
}
Expand Down Expand Up @@ -211,7 +213,8 @@ describe('persistent command retention', () => {
turnLifecycle: null,
sink,
streams: items.streams,
activeItems: items.activeItems
activeItems: items.activeItems,
linkageFor: () => ({})
})
).toEqual({ accepted: true })
expect(items.activeItems.size).toBe(1)
Expand Down
13 changes: 11 additions & 2 deletions src/main/codex/codex-structured-item-stream-contracts.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,14 +2,18 @@ import type { AgentJournalItemIdentity } from '../../shared/agent-session-journa
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 type { codexJournalItem, CodexThreadItem } from './codex-structured-item-translation'
import type { CodexRowLinkage } from './codex-subagent-linkage'

export type CodexItemStreamDeps = {
sink: StructuredAgentSessionEventSink
/** The turn a delta-only item belongs to, read the way its first delta names it. */
turnIdFor: (threadId: string, params: unknown) => string | null
identityFor: (
threadId: string,
params: unknown,
turnId: string | null,
item: CodexThreadItem
) => AgentJournalItemIdentity
linkageFor: CodexRowLinkage
coalesceMs?: number
maxRetainedBytes?: number
maxTotalRetainedBytes?: number
Expand Down Expand Up @@ -39,7 +43,12 @@ export type CodexStructuredItemStreamHandleResult = {
export type CodexStructuredItemStreams = {
readonly persistentCount: number
canTrack: (threadId: string, item: CodexThreadItem, identity: AgentJournalItemIdentity) => boolean
track: (threadId: string, item: CodexThreadItem, identity: AgentJournalItemIdentity) => boolean
track: (
threadId: string,
turnId: string | null,
item: CodexThreadItem,
identity: AgentJournalItemIdentity
) => boolean
handle: (
threadId: string,
method: string,
Expand Down
52 changes: 27 additions & 25 deletions src/main/codex/codex-structured-item-streams.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import { agentJournalItemKey } from '../../shared/agent-session-journal-item-key'
import { createAgentSessionDeltaCoalescer } from '../native-chat/agent-session-wire/agent-session-delta-coalescer'
import { CodexItemStreamRetention } from './codex-item-stream-retention'
import { appendCodexItemAndPublish } from './codex-structured-journal-sink'
import {
codexJournalItem,
codexStreamingJournalItem,
Expand Down Expand Up @@ -50,6 +51,13 @@ export function createCodexStructuredItemStreams(
const states = new CodexItemStreamRetention(deps.maxMetadataBytes)
const checkpointLengths = new Map<string, number>()
const pendingCheckpoints = new Set<string>()
// Which thread and turn produced each stream, resolved to linkage per append
// so a parent learned after the first checkpoint still reaches the row.
const producers = new Map<string, { threadId: string; turnId: string | null }>()
const linkageOf = (key: string) => {
const producer = producers.get(key)
return producer ? deps.linkageFor(producer.threadId, producer.turnId) : {}
}
// Patch updates are authoritative item snapshots. Keep the latest rejected
// snapshot until the journal admits it; unlike streamed deltas, there is no
// coalescer timer to retry these events for us.
Expand All @@ -61,6 +69,7 @@ export function createCodexStructuredItemStreams(
states.forget(key)
checkpointLengths.delete(key)
pendingCheckpoints.delete(key)
producers.delete(key)
const pending = pendingPatches.get(key)
if (pending) {
retainedPatchBytes = Math.max(0, retainedPatchBytes - pendingPatchBytes(pending))
Expand Down Expand Up @@ -98,23 +107,15 @@ export function createCodexStructuredItemStreams(
}
}

const append = (state: CodexItemStreamState, text: string): boolean => {
const append = (key: string, state: CodexItemStreamState, text: string): boolean => {
const translated = codexStreamingJournalItem(state.item, text)
if (!translated.body) {
return true
}
const options = { coalescingKey: `checkpoint:${agentJournalItemKey(state.identity)}` }
const admission = deps.sink.tryAppendItem
? deps.sink.tryAppendItem(state.identity, translated.body, options)
: (deps.sink.appendItem(state.identity, translated.body, options),
{ accepted: true as const })
if (!admission.accepted) {
return false
}
const published = deps.sink.tryPublish
? deps.sink.tryPublish()
: (deps.sink.publish(), { accepted: true as const })
return published.accepted
return appendCodexItemAndPublish(deps.sink, state.identity, translated.body, {
coalescingKey: `checkpoint:${agentJournalItemKey(state.identity)}`,
...linkageOf(key)
}).accepted
}

const persist = (key: string, text: string, force: boolean): boolean => {
Expand All @@ -124,7 +125,7 @@ export function createCodexStructuredItemStreams(
return true
}
const state = states.get(key)
if (state && append(state, text)) {
if (state && append(key, state, text)) {
checkpointLengths.set(key, text.length)
pendingCheckpoints.delete(key)
return true
Expand Down Expand Up @@ -155,10 +156,12 @@ export function createCodexStructuredItemStreams(
return existing
}
const item = { type, id: itemId }
const state = { item, identity: deps.identityFor(threadId, params, item) }
const turnId = deps.turnIdFor(threadId, params)
const state = { item, identity: deps.identityFor(threadId, turnId, item) }
if (!states.retain(key, state)) {
return null
}
producers.set(key, { threadId, turnId })
trimStates()
return state
}
Expand All @@ -184,18 +187,15 @@ export function createCodexStructuredItemStreams(
if (!pending) {
return { accepted: true }
}
const admission = deps.sink.tryAppendItem
? deps.sink.tryAppendItem(pending.identity, pending.body)
: (deps.sink.appendItem(pending.identity, pending.body), { accepted: true as const })
const admission = appendCodexItemAndPublish(
deps.sink,
pending.identity,
pending.body,
linkageOf(key)
)
if (!admission.accepted) {
return admission
}
const published = deps.sink.tryPublish
? deps.sink.tryPublish()
: (deps.sink.publish(), { accepted: true as const })
if (!published.accepted) {
return published
}
retainedPatchBytes = Math.max(0, retainedPatchBytes - pendingPatchBytes(pending))
pendingPatches.delete(key)
return { accepted: true }
Expand All @@ -210,11 +210,12 @@ export function createCodexStructuredItemStreams(
item: boundStreamItem(item) as CodexThreadItem,
identity
}),
track: (threadId, item, identity) => {
track: (threadId, turnId, item, identity) => {
const key = codexStructuredItemKey(threadId, item.id)
if (!states.retain(key, { item: boundStreamItem(item) as CodexThreadItem, identity })) {
return false
}
producers.set(key, { threadId, turnId })
trimStates()
return true
},
Expand Down Expand Up @@ -299,6 +300,7 @@ export function createCodexStructuredItemStreams(
states.clear()
checkpointLengths.clear()
pendingCheckpoints.clear()
producers.clear()
pendingPatches.clear()
retainedPatchBytes = 0
},
Expand Down
7 changes: 5 additions & 2 deletions src/main/codex/codex-structured-journal-compactions.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,13 +8,15 @@ import {
import { MAX_CODEX_GENERIC_TURN_BUCKETS } from './codex-structured-journal-limits'
import { appendCodexLifecycleItem, publishCodexLifecycle } from './codex-structured-journal-sink'
import { readCodexTurnId } from './codex-structured-thread-facts'
import type { CodexRowLinkage } from './codex-subagent-linkage'

export class CodexJournalCompactions {
private readonly turns = new Map<string, 'item' | 'legacy'>()

constructor(
private readonly sink: StructuredAgentSessionEventSink,
private readonly activeTurn: (threadId: string) => string | null
private readonly activeTurn: (threadId: string) => string | null,
private readonly linkageFor: CodexRowLinkage
) {}

handle(event: {
Expand All @@ -41,7 +43,8 @@ export class CodexJournalCompactions {
const admission = appendCodexLifecycleItem(
this.sink,
{ provider: 'orca', clientMessageId: `codex-compaction:${key}` },
{ kind: 'status', text: 'Context compacted', presentation: 'compaction' }
{ kind: 'status', text: 'Context compacted', presentation: 'compaction' },
this.linkageFor(event.threadId, turnId)
)
if (!admission.accepted) {
return admission
Expand Down
75 changes: 45 additions & 30 deletions src/main/codex/codex-structured-journal-generic-frames.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,9 +12,16 @@ import {
MAX_CODEX_GENERIC_TURN_BUCKETS
} from './codex-structured-journal-limits'
import { readCodexTurnId } from './codex-structured-thread-facts'
import type { CodexRowLinkage } from './codex-subagent-linkage'

const OVERFLOW_BUCKET = '__codex-generic-overflow__'
type SuppressedSummary = { count: number; publishedCount: number }
/** `producer` is absent only on the overflow bucket, which pools every thread's
* evicted turns and so has no single author: it reads as the session's own. */
type SuppressedSummary = {
count: number
publishedCount: number
producer?: { threadId: string; turnId: string }
}

function boundedTurnBucket(threadId: string, turnId: string): string {
const encoded = `${encodeURIComponent(threadId)}:${encodeURIComponent(turnId)}`
Expand Down Expand Up @@ -51,7 +58,9 @@ export class CodexJournalGenericFrames {
private cancelSuppressionFlush: (() => void) | null = null

constructor(
private readonly deps: Pick<CodexJournalTranslatorDeps, 'sink' | 'schedule' | 'coalesceMs'>,
private readonly deps: Pick<CodexJournalTranslatorDeps, 'sink' | 'schedule' | 'coalesceMs'> & {
linkageFor: CodexRowLinkage
},
private readonly activeTurn: (threadId: string) => string | null
) {
this.schedule = deps.schedule ?? defaultSchedule
Expand All @@ -61,22 +70,23 @@ export class CodexJournalGenericFrames {
appendUnhandled(
kind: string,
payload: unknown,
threadId = 'session'
threadId: string
): CodexJournalTranslationAdmission {
const translated = unhandledProviderFrameJournalItem('codex', kind, payload)
// 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) {
return CODEX_JOURNAL_ADMITTED
}
const turnId = readCodexTurnId(payload) ?? this.activeTurn(threadId) ?? 'outside-turn'
const frameTurnId = readCodexTurnId(payload) ?? this.activeTurn(threadId)
const turnId = frameTurnId ?? 'outside-turn'
const bucket = this.bucketFor(threadId, turnId)
const rowCount = this.genericRowsByTurn.get(bucket) ?? 0
// The cap bounds noise, never evidence: an error frame is always journaled, and
// capped frames stay countable through one summary row per turn.
const isError = translated.classification === 'error-surface'
if (!isError && rowCount >= MAX_CODEX_GENERIC_ROWS_PER_TURN) {
this.addSuppressed(bucket, 1)
this.addSuppressed(bucket, 1, { threadId, turnId })
this.recordBucket(bucket)
this.scheduleSuppressedRows()
return CODEX_JOURNAL_ADMITTED
Expand All @@ -88,16 +98,14 @@ export class CodexJournalGenericFrames {
}
}
this.fallbackSequence += 1
const identity = {
provider: 'orca' as const,
clientMessageId: `provider-frame:codex:${this.fallbackSequence}`
}
const linkage = this.deps.linkageFor(threadId, frameTurnId)
const admission = this.deps.sink.tryAppendItem
? this.deps.sink.tryAppendItem(
{ provider: 'orca', clientMessageId: `provider-frame:codex:${this.fallbackSequence}` },
translated.body
)
: (this.deps.sink.appendItem(
{ provider: 'orca', clientMessageId: `provider-frame:codex:${this.fallbackSequence}` },
translated.body
),
CODEX_JOURNAL_ADMITTED)
? this.deps.sink.tryAppendItem(identity, translated.body, linkage)
: (this.deps.sink.appendItem(identity, translated.body, linkage), CODEX_JOURNAL_ADMITTED)
if (!admission.accepted) {
this.fallbackSequence -= 1
return admission
Expand All @@ -109,7 +117,7 @@ export class CodexJournalGenericFrames {

suppress(threadId: string, turnId: string, count = 1): void {
const bucket = this.bucketFor(threadId, turnId)
this.addSuppressed(bucket, count)
this.addSuppressed(bucket, count, { threadId, turnId })
this.recordBucket(bucket)
}

Expand All @@ -127,20 +135,19 @@ export class CodexJournalGenericFrames {
bucket === OVERFLOW_BUCKET
? `${summary.count} more provider notification${summary.count === 1 ? '' : 's'} not shown across evicted turns`
: `${summary.count} more provider notification${summary.count === 1 ? '' : 's'} not shown for this turn`
const identity = {
provider: 'orca' as const,
clientMessageId: `provider-frame-suppressed:codex:${bucket}`
}
const options = {
coalescingKey: `provider-frame-suppressed:codex:${bucket}`,
...(summary.producer
? this.deps.linkageFor(summary.producer.threadId, summary.producer.turnId)
: {})
}
const admission = this.deps.sink.tryAppendItem
? this.deps.sink.tryAppendItem(
{ provider: 'orca', clientMessageId: `provider-frame-suppressed:codex:${bucket}` },
{
kind: 'status',
text
},
{ coalescingKey: `provider-frame-suppressed:codex:${bucket}` }
)
: (this.deps.sink.appendItem(
{ provider: 'orca', clientMessageId: `provider-frame-suppressed:codex:${bucket}` },
{ kind: 'status', text },
{ coalescingKey: `provider-frame-suppressed:codex:${bucket}` }
),
? this.deps.sink.tryAppendItem(identity, { kind: 'status', text }, options)
: (this.deps.sink.appendItem(identity, { kind: 'status', text }, options),
CODEX_JOURNAL_ADMITTED)
if (!admission.accepted) {
blocked ??= admission
Expand Down Expand Up @@ -188,8 +195,16 @@ export class CodexJournalGenericFrames {
: requested
}

private addSuppressed(bucket: string, count: number): void {
const summary = this.suppressedRowsByTurn.get(bucket) ?? { count: 0, publishedCount: 0 }
private addSuppressed(
bucket: string,
count: number,
producer?: SuppressedSummary['producer']
): void {
const summary = this.suppressedRowsByTurn.get(bucket) ?? {
count: 0,
publishedCount: 0,
...(producer && bucket !== OVERFLOW_BUCKET ? { producer } : {})
}
summary.count += count
this.suppressedRowsByTurn.set(bucket, summary)
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -137,7 +137,7 @@ describe('codex goal lifecycle admission', () => {
appendTombstone: () => {},
publish: () => {}
} satisfies StructuredAgentSessionEventSink
const goals = new CodexJournalGoals(sink)
const goals = new CodexJournalGoals(sink, () => ({}))
const update = (goal: Record<string, unknown> = {}) =>
goals.handle({ threadId: THREAD, method: 'thread/goal/updated', params: goalFrame(goal) })
const clear = () =>
Expand Down Expand Up @@ -178,7 +178,7 @@ describe('codex goal lifecycle admission', () => {
appendTombstone: () => {},
publish: () => {}
} satisfies StructuredAgentSessionEventSink
const goals = new CodexJournalGoals(sink)
const goals = new CodexJournalGoals(sink, () => ({}))
const send = (threadId: string) =>
goals.handle({ threadId, method: 'thread/goal/updated', params: goalFrame() })

Expand Down Expand Up @@ -211,7 +211,7 @@ describe('codex goal lifecycle admission', () => {
appendTombstone: () => {},
publish: () => {}
} satisfies StructuredAgentSessionEventSink
const goals = new CodexJournalGoals(sink)
const goals = new CodexJournalGoals(sink, () => ({}))
const event = { threadId: THREAD, method: 'thread/goal/updated', params: goalFrame() }

goals.handle(event)
Expand Down
Loading
Loading