Bound provider headline updates and clear activity on reconnect

This commit is contained in:
Merge Sim
2026-09-06 16:06:45 -07:00
parent a88c6fe773
commit 8cc2ad7484
5 changed files with 179 additions and 13 deletions
@@ -1,4 +1,4 @@
import { codexProviderFrameActivity } from '../native-chat/agent-session-wire/provider-frame-activity'
import { createCodexProviderActivityReader } from '../native-chat/agent-session-wire/provider-frame-activity'
import { CodexJournalGenericFrames } from './codex-structured-journal-generic-frames'
import { CodexJournalItems } from './codex-structured-journal-items'
import { CodexJournalPrompts } from './codex-structured-journal-prompts'
@@ -57,6 +57,7 @@ export function createCodexJournalTranslator(
)
const flushStreams = (): CodexJournalTranslationAdmission =>
items.streams.flush() ? CODEX_JOURNAL_ADMITTED : { accepted: false, reason: 'backpressure' }
let readActivity = createCodexProviderActivityReader()
const publishActivity = (
event: Extract<CodexStructuredSessionEvent, { type: 'notification' }>,
admission: CodexJournalTranslationAdmission
@@ -68,10 +69,7 @@ export function createCodexJournalTranslator(
if (!turnId) {
return admission
}
const params = readCodexJournalRecord(event.params)
const itemId = readCodexJournalString(params, 'itemId')
const reasoningText = itemId ? items.streams.snapshot(event.threadId, itemId)?.text : null
const text = codexProviderFrameActivity(event.method, event.params, reasoningText)
const text = readActivity(event.method, event.params)
if (text !== undefined) {
deps.sink.setActivity?.(text ? { turnId, text } : null)
}
@@ -79,8 +77,11 @@ export function createCodexJournalTranslator(
}
return {
restoreThread: (threadId, thread) =>
restoreCodexJournalThread({
restoreThread: (threadId, thread) => {
if (threadId === (deps.primaryThreadId?.() ?? null)) {
readActivity = createCodexProviderActivityReader()
}
return restoreCodexJournalThread({
threadId,
thread,
currentTurnIds: activeTurns.byThread,
@@ -92,7 +93,8 @@ export function createCodexJournalTranslator(
: { accepted: false, reason: 'untranslated' }
},
flush: items.streams.flush
}),
})
},
handle: (event) => {
if (event.type === 'ended') {
const streamAdmission = flushStreams()
@@ -116,6 +118,7 @@ export function createCodexJournalTranslator(
if (!admission.accepted) {
return admission
}
readActivity = createCodexProviderActivityReader()
deps.sink.setActivity?.(null)
items.activeItems.clear()
prompts.pending.clear()
@@ -232,6 +235,7 @@ export function createCodexJournalTranslator(
if (admission.accepted) {
activeTurns.remember(event.threadId, turnId)
if (event.threadId === (deps.primaryThreadId?.() ?? null)) {
readActivity = createCodexProviderActivityReader()
deps.sink.setActivity?.(null)
}
}
@@ -263,6 +267,7 @@ export function createCodexJournalTranslator(
items.ordinals.forgetTurn(event.threadId, turnId)
activeTurns.forget(event.threadId, turnId)
if (event.threadId === (deps.primaryThreadId?.() ?? null)) {
readActivity = createCodexProviderActivityReader()
deps.sink.setActivity?.(null)
}
}
@@ -138,3 +138,51 @@ export function claudeProviderFrameActivity(kind: string, payload: unknown): Act
}
return undefined
}
/** Retain only the current summary headline, never materialize the growing transcript. */
export function createCodexProviderActivityReader(): (
method: string,
payload: unknown
) => ActivityText {
let itemId: unknown
let summaryIndex: unknown
let headline = ''
let complete = false
const limit = MAX_PROVIDER_ACTIVITY_LENGTH * 2 + 16
return (method, payload) => {
if (
method !== 'item/reasoning/summaryTextDelta' &&
method !== 'item/reasoning/summaryPartAdded'
) {
return codexProviderFrameActivity(method, payload)
}
const source = record(payload)
if (!stringField(source, 'itemId')) {
return undefined
}
if (
source?.itemId !== itemId ||
source?.summaryIndex !== summaryIndex ||
method === 'item/reasoning/summaryPartAdded'
) {
itemId = source?.itemId
summaryIndex = source?.summaryIndex
headline = ''
complete = false
}
if (method === 'item/reasoning/summaryPartAdded') {
return null
}
if (complete || typeof source?.delta !== 'string') {
return undefined
}
headline += source.delta.slice(0, limit - headline.length)
const line = headline.trimStart().split(/\r?\n/, 1)[0]
complete =
headline.length === limit || /\r?\n/.test(headline.trimStart()) || /^\*\*.+\*\*/.test(line)
if (complete && line.startsWith('**') && !/\*\*.+\*\*/.test(line)) {
return providerActivityText(line.slice(2))
}
return codexProviderFrameActivity(method, payload, line)
}
}
@@ -7,6 +7,7 @@ import type { AgentSessionTurnActivity } from '../../../shared/agent-session-wir
import { createClaudeJournalTranslator } from '../../claude/claude-structured-journal-translation'
import { createCodexJournalTranslator } from '../../codex/codex-structured-journal-translation'
import type { CodexStructuredSessionEvent } from '../../codex/codex-structured-session-state'
import * as deltaCoalescer from './agent-session-delta-coalescer'
import type { StructuredAgentSessionEventSink } from './structured-agent-session-event-sink'
const SESSION_ID = 'session-1'
@@ -87,6 +88,109 @@ describe('provider turn activity routing', () => {
expect(state.activities.at(-1)?.text).toBe('Tracing the activity pipeline')
})
it('does not materialize full stream snapshots for activity on token deltas', () => {
const original = deltaCoalescer.createAgentSessionDeltaCoalescer
const snapshot = vi.fn()
const factory = vi
.spyOn(deltaCoalescer, 'createAgentSessionDeltaCoalescer')
.mockImplementation((deps) => {
const coalescer = original(deps)
return {
...coalescer,
snapshot: (key) => {
snapshot()
return coalescer.snapshot(key)
}
}
})
try {
const state = recordingSink()
const translator = createCodexJournalTranslator({
sink: state.sink,
primaryThreadId: () => THREAD_ID,
schedule: () => () => {}
})
translator.handle(codexNotification('turn/started', { turn: { id: TURN_ID } }))
for (const method of [
'item/agentMessage/delta',
'item/commandExecution/outputDelta',
'item/reasoning/summaryTextDelta'
]) {
for (let index = 0; index < 100; index++) {
translator.handle(
codexNotification(method, {
turnId: TURN_ID,
itemId: method,
summaryIndex: 0,
delta: index === 0 ? '**Inspecting**\n' : 'more output'
})
)
}
}
expect(snapshot).not.toHaveBeenCalled()
translator.dispose()
} finally {
factory.mockRestore()
}
})
it('uses the newest summary part and stops republishing its body', () => {
const state = recordingSink()
const translator = createCodexJournalTranslator({
sink: state.sink,
primaryThreadId: () => THREAD_ID,
schedule: () => () => {}
})
translator.handle(codexNotification('turn/started', { turn: { id: TURN_ID } }))
const params = { turnId: TURN_ID, itemId: 'reasoning-1' }
for (const [summaryIndex, headline] of ['First headline', 'Newest headline'].entries()) {
translator.handle(
codexNotification('item/reasoning/summaryPartAdded', { ...params, summaryIndex })
)
translator.handle(
codexNotification('item/reasoning/summaryTextDelta', {
...params,
summaryIndex,
delta: `**${headline}`
})
)
expect(state.activities.at(-1)).toBeNull()
translator.handle(
codexNotification('item/reasoning/summaryTextDelta', {
...params,
summaryIndex,
delta: '**\n\nBody'
})
)
expect(state.activities.at(-1)?.text).toBe(headline)
}
const publications = state.activities.length
for (let index = 0; index < 100; index++) {
translator.handle(
codexNotification('item/reasoning/summaryTextDelta', {
...params,
summaryIndex: 1,
delta: ' more body'
})
)
}
expect(state.activities).toHaveLength(publications)
translator.handle(
codexNotification('turn/completed', { turn: { id: TURN_ID, status: 'completed' } })
)
translator.handle(codexNotification('turn/started', { turn: { id: 'turn-2' } }))
translator.handle(
codexNotification('item/reasoning/summaryTextDelta', {
...params,
turnId: 'turn-2',
summaryIndex: 1,
delta: '**Next turn**'
})
)
expect(state.activities.at(-1)).toEqual({ turnId: 'turn-2', text: 'Next turn' })
translator.dispose()
})
it('keeps Codex tool rows singular and the activity free of tool labels', () => {
const state = recordingSink()
const translator = createCodexJournalTranslator({
@@ -67,7 +67,8 @@ describe('AgentSessionSubscribers', () => {
removedItemIds: [],
submissions: []
},
fence: 7
fence: 7,
activity: null
}
])
})
@@ -290,7 +291,16 @@ describe('AgentSessionSubscribers', () => {
activity: { turnId: 'turn-1', text: 'Inspecting the session wire' }
})
subscribers.close(SESSION, 'subscriber-1')
subscribers.publish(SESSION, journal, null)
subscribers.open({
id: 'reconnected',
sessionId: SESSION,
journal,
fence: 1,
cursor,
emit: (event) => events.push(event)
})
expect(journal.cursor()).toEqual(cursor)
expect(events.at(-1)).toMatchObject({ activity: null })
})
@@ -230,7 +230,7 @@ export class AgentSessionSubscribers {
activity !== undefined
? activity
: emitCheckpoint
? this.activityBySession.get(subscriber.sessionId)
? (this.activityBySession.get(subscriber.sessionId) ?? null)
: undefined
while (true) {
const result = readAgentSessionHistory(journal, {
@@ -318,8 +318,7 @@ export class AgentSessionSubscribers {
}
}
private activityField(sessionId: string): { activity?: AgentSessionTurnActivity } {
const activity = this.activityBySession.get(sessionId)
return activity ? { activity } : {}
private activityField(sessionId: string): { activity: AgentSessionTurnActivity | null } {
return { activity: this.activityBySession.get(sessionId) ?? null }
}
}