Files
orca/src/main/codex/codex-structured-session-adapter.ts
T
Brennan Benson 9def4b9ba1 feat(native-chat): publish where each submission was sent and the turn a stopped one was answered into (#25073)
* feat(native-chat): publish each submission's journal positions

A client never learned where in the journal a send was accepted or taken back,
so drawing a stopped send had to guess from its own row and from clocks, and the
guesses misfired after restarts. Each projected submission now carries
`submittedSequence` (its submission row) and `resolvedSequence` (the dispatch row
that resolved it; absent while pending): two optional fields, computed on every
fold from rows already stored, never stored themselves. Older clients ignore
them. The failure-fact schema moves to its own module to keep the journal schema
file within its size limit.

* test(native-chat): pin a submission's resolved position on the provider-echo accept path

Also say which rows can resolve a submission, and that the submission schema, not the shared type, omits acceptedSequence.

* refactor(native-chat): publish only where each submission was sent

A rejected send's own row now moves to its rejection (#24710), so the row that
resolved a send is the item's own position and needs no second copy. Keep
submittedSequence, the one journal-order record of where a send was sent.

* feat(native-chat): publish the turn a withdrawn Codex send was answered into

A send Codex answered into a turn that then ended without taking it is
rejected by that turn's end. The rejection row now names that turn's
record, and the submission carries it as answeredInTurnItemId, so a
client can tell which turn the send belonged to without guessing from
journal order or clocks. Absent on every other send.

* feat(native-chat): say how a withdrawn Codex send joined the turn it names

The answered turn becomes { turnItemId, via }: 'start' when Codex
answered the send's turn/start with that turn, 'steer' when Orca steered
it into that running turn. A client can then tell the turn's opener from
a send steered into it even when the opener is on an older page.

* feat(native-chat): state on every rejection whether it names a turn

A rejection a host writes now always carries answeredInTurn: the turn
Codex answered the send into, or null for none. A reader can then tell
'answered into no turn' from a row written before the field existed,
which keeps no answer and is placed by the older rules.

* refactor(native-chat): give the answered-turn schema its own module

Keeps the journal schema file within its line limit.

* docs(native-chat): say an unreadable answered turn reads as none
2026-10-05 15:56:55 -07:00

296 lines
12 KiB
TypeScript

import * as codexRewind from './codex-structured-rewind'
import type { AgentJournalMessageItem } from '../../shared/agent-session-journal-types'
import type { AgentSessionBackgroundTaskState } from '../../shared/agent-session-wire'
import { isCodexAppServerRequestError } from './codex-app-server-connection'
import type {
AgentSessionAcquisition,
AgentSessionDispatchOutcome,
StructuredAgentSessionAcquireInput,
StructuredAgentSessionAdapter,
StructuredAgentSessionSetOptionInput
} from '../native-chat/agent-session-wire/structured-agent-session-adapter'
import type { CodexJournalTranslationAdmission } from './codex-structured-journal-translation'
import { dispatchCodexTurn, isCodexTurnOptionKey } from './codex-structured-turn-start'
import { agentSessionFailureFact, providerDiagnosticOf } from '../../shared/agent-session-failure'
import { supportsCodexStructuredLocation } from './codex-structured-location-support'
import { CodexStructuredSessionTeardown } from './codex-structured-session-teardown'
import {
applyCodexStructuredSessionOption,
readLiveCodexSessionOptions
} from './codex-structured-session-options'
import {
CodexAcquisitionRegistry,
requireLiveCodexSession,
type CodexAcquisitionAttempt,
type CodexSession,
type CodexStructuredSessionAdapterDeps,
type CodexStructuredSessionEvent
} from './codex-structured-session-state'
import {
deliverCodexServerRequest,
deliverCodexUnhandledFrame,
translateCodexNotification
} from './codex-structured-provider-events'
import {
codexDispatchRejection,
settleCodexSendsInEndedTurn
} from './codex-structured-turn-end-settlement'
import { createCodexStructuredNotificationRetry } from './codex-structured-notification-retry'
import { acquireCodexStructuredSession } from './codex-structured-session-acquire'
import { changeCodexThreadGoal } from './codex-structured-thread-goal'
import {
answerCodexStructuredPrompt,
cancelCodexStructuredTurn
} from './codex-structured-prompt-ownership'
export type {
CodexStructuredLaunch,
CodexStructuredSessionAdapterDeps,
CodexStructuredSessionEvent
} from './codex-structured-session-state'
export class CodexStructuredSessionAdapter implements StructuredAgentSessionAdapter {
private readonly sessions = new Map<string, CodexSession>()
private readonly acquisitions = new CodexAcquisitionRegistry()
private readonly notificationRetries: ReturnType<typeof createCodexStructuredNotificationRetry>
private readonly teardown: CodexStructuredSessionTeardown
constructor(private readonly deps: CodexStructuredSessionAdapterDeps) {
this.notificationRetries = createCodexStructuredNotificationRetry({
sessionFor: (sessionId) => this.sessions.get(sessionId),
translate: (sessionId, session, method, params, observedAt, dispatchSequenceAtReceipt) =>
translateCodexNotification({
sessionId,
session,
method,
params,
observedAt,
dispatchSequenceAtReceipt,
emit: (current, event) => this.emit(current, event)
})
})
this.teardown = new CodexStructuredSessionTeardown({
sessions: this.sessions,
acquisitions: this.acquisitions,
...(deps.onEvent ? { onEvent: deps.onEvent } : {}),
...(deps.logger ? { logger: deps.logger } : {}),
forgetNotificationRetries: (sessionId) => this.notificationRetries.clear(sessionId, null)
})
}
supportsLocation = (location: Parameters<typeof supportsCodexStructuredLocation>[0]): boolean =>
supportsCodexStructuredLocation(location, this.deps.isWindowsProcessStartTimeAvailable)
acquire = (input: StructuredAgentSessionAcquireInput): Promise<AgentSessionAcquisition> =>
acquireCodexStructuredSession({
input,
deps: this.deps,
sessions: this.sessions,
acquisitions: this.acquisitions,
notificationRetries: this.notificationRetries,
deliver: (acquisition, sessionId, event, retainedBytes) =>
this.deliver(acquisition, sessionId, event, retainedBytes),
handleServerRequest: (sessionId, request) => this.handleServerRequest(sessionId, request),
handleUnhandledFrame: (sessionId, kind, payload) =>
this.handleUnhandledFrame(sessionId, kind, payload),
forceCloseUnexpected: (sessionId, fence, acquisitionGeneration, reason) =>
this.teardown.forceCloseUnexpected(sessionId, fence, acquisitionGeneration, reason)
})
/** Buffers pre-publication events and drops events from superseded children. */
private deliver(
acquisition: CodexAcquisitionAttempt['window'],
sessionId: string,
event: () => unknown,
retainedBytes?: number
): void {
if (acquisition.buffer(event, retainedBytes)) {
return
}
if (this.sessions.get(sessionId)?.connection === acquisition.connection) {
event()
} else if (acquisition.isOverflowed) {
// Pre-publication overflow is an acquisition failure, not a dropped
// notification; tear down the child so callers retry explicitly.
void acquisition.connection?.close()
}
}
/** Journal first so observers never see an event ahead of its durable row. */
private emit(
session: CodexSession,
event: CodexStructuredSessionEvent
): CodexJournalTranslationAdmission {
if (event.type === 'notification' && !session.backgroundTasks.canObserve(event)) {
return { accepted: false, reason: 'failed' }
}
const admission = session.translator?.handle(event) ?? { accepted: true }
if (!admission.accepted) {
return admission
}
if (event.type === 'notification') {
// Only an admitted turn end settles; a refused one settles on the retry that lands.
settleCodexSendsInEndedTurn(session, event, (settlement) =>
this.deps.onDispatchSettledLate?.({ sessionId: event.sessionId, ...settlement })
)
// After the admission check, so a refused frame is observed by the child records
// only on the retry that also reaches the journal.
session.backgroundTasks.observe(event, session.prompts.takeAbandonedCommands())
// After the journal and the parent's republished row, never ahead of either.
session.backgroundTasks.publishChildWork()
}
this.deps.onEvent?.(event)
return admission
}
private handleServerRequest(
sessionId: string,
request: Parameters<typeof deliverCodexServerRequest>[2]
): void {
deliverCodexServerRequest(sessionId, this.sessions.get(sessionId), request, (session, event) =>
this.emit(session, event)
)
}
private handleUnhandledFrame(sessionId: string, kind: string, params: unknown): void {
deliverCodexUnhandledFrame(
sessionId,
this.sessions.get(sessionId),
kind,
params,
(session, event) => this.emit(session, event)
)
}
/** The tracker's own roster. No host decision reads it: the host's child records are the one
* owner of "what runs", and this stays only so tests can hold the two rule sets side by side. */
backgroundTaskState = (sessionId: string): AgentSessionBackgroundTaskState | null | undefined =>
this.sessions.get(sessionId)?.backgroundTasks.state
// Codex exposes no honest stop for a child thread or a persistent command.
backgroundTaskStops: NonNullable<StructuredAgentSessionAdapter['backgroundTaskStops']> = (
sessionId
) =>
this.sessions.has(sessionId) ? { supportsTaskStop: false, supportsStopAll: false } : undefined
bindPromptItemId = (
sessionId: string,
journalItemId: string,
promptKey: string,
turnId?: string | null,
threadId?: string
): void =>
this.sessions
.get(sessionId)
?.prompts.bindJournalItemId(
journalItemId,
threadId ?? this.session(sessionId).threadId,
promptKey,
turnId
)
async dispatch(input: {
sessionId: string
clientMessageId: string
body: AgentJournalMessageItem
fence: number
requestedAt?: number
beforeDispatch?: () => Promise<void>
}): Promise<AgentSessionDispatchOutcome> {
const session = this.session(input.sessionId)
session.dispatchPending = true
try {
await input.beforeDispatch?.()
return await dispatchCodexTurn(session, input, this.deps.requestTimeoutMs)
} finally {
session.dispatchPending = false
}
}
cancelTurn: StructuredAgentSessionAdapter['cancelTurn'] = (request) =>
cancelCodexStructuredTurn({
request,
sessions: this.sessions,
requestTimeoutMs: this.deps.requestTimeoutMs
})
rewindSupport: NonNullable<StructuredAgentSessionAdapter['rewindSupport']> = (sessionId) =>
this.sessions.get(sessionId)?.historyMode === 'legacy'
? { supported: false, reason: 'history-not-paginated' }
: { supported: true }
rewind: NonNullable<StructuredAgentSessionAdapter['rewind']> = (input) =>
codexRewind.rewindCodexSession(this.session(input.sessionId), input, this.deps.requestTimeoutMs)
recoverRewind: NonNullable<StructuredAgentSessionAdapter['recoverRewind']> = (input) =>
codexRewind.recoverCodexRewind(this.session(input.sessionId), input, this.deps.requestTimeoutMs)
/** The ack is Codex's receipt; the translator ends the command's turn from the turn it opens. */
compact: NonNullable<StructuredAgentSessionAdapter['compact']> = async (input) => {
const session = this.session(input.sessionId)
session.translator?.beginCommand(input.command)
try {
await session.connection.request(
'thread/compact/start',
{ threadId: session.threadId },
{ timeoutMs: this.deps.requestTimeoutMs }
)
return { state: 'accepted', providerIdentity: null }
} catch (error) {
session.translator?.forgetCommand(input.command.turnId)
if (isCodexAppServerRequestError(error)) {
// Codex's own words, when it gave any, are the one part of the error a person can use.
return {
state: 'rejected',
...codexDispatchRejection(
agentSessionFailureFact('providerRejected', { detail: providerDiagnosticOf(error) })
)
}
}
throw error
}
}
changeThreadGoal: NonNullable<StructuredAgentSessionAdapter['changeThreadGoal']> = (input) =>
changeCodexThreadGoal(
this.session(input.sessionId),
input.change,
input.replacesGoal,
this.deps.requestTimeoutMs
)
// Provider-level: a goal change at rest starts the agent first.
supportsThreadGoal = (): boolean => true
answerPrompt: StructuredAgentSessionAdapter['answerPrompt'] = (request) =>
answerCodexStructuredPrompt({ request, sessions: this.sessions })
async setOption(
input: StructuredAgentSessionSetOptionInput
): Promise<Readonly<Record<string, string>>> {
if (!isCodexTurnOptionKey(input.key)) {
throw new Error(`codex app-server has no thread option named ${input.key}`)
}
return applyCodexStructuredSessionOption(
this.session(input.sessionId),
input.key,
input.value,
this.deps.requestTimeoutMs
)
}
readOptions = (input: { sessionId: string; fence: number }) =>
readLiveCodexSessionOptions(this.session(input.sessionId), this.deps.requestTimeoutMs)
closeSession = (sessionId: string): Promise<boolean> => this.teardown.close(sessionId)
forceCloseSession = (sessionId: string): Promise<boolean> => this.teardown.forceClose(sessionId)
disposeSession = (sessionId: string): Promise<boolean> => this.teardown.close(sessionId)
closeAll = (): Promise<void> => this.teardown.closeAll()
releaseAcquisition = (input: { sessionId: string }): Promise<boolean> =>
this.teardown.close(input.sessionId)
private session(sessionId: string): CodexSession {
return requireLiveCodexSession(this.sessions, sessionId)
}
}