mirror of
https://github.com/stablyai/orca.git
synced 2026-09-29 08:03:20 +00:00
refactor(codex): move session teardown out of the structured adapter
Merging main crossed the 300-line cap on `codex-structured-session-adapter.ts`: the rewind backend (#19235) and this branch's close-time strip clear both landed in it. The four close paths move verbatim into `codex-structured-session-teardown.ts`, where they funnel through one `settled` helper instead of repeating the notification-retry and background-task cleanup at each call site. No ratchet bump. Also normalize a background task's description once at receipt rather than on every projection; the roster is re-projected on each observed frame.
This commit is contained in:
@@ -46,18 +46,20 @@ const AGENT_TASK_PREFIX = 'codex-agent:'
|
||||
const COMMAND_TASK_PREFIX = 'codex-command:'
|
||||
|
||||
type TrackedChild = {
|
||||
label: string | null
|
||||
label: string | undefined
|
||||
state: NativeChatSubagentState
|
||||
/** Parent turn the child was spawned in; null when Codex named none. */
|
||||
turnId: string | null
|
||||
}
|
||||
|
||||
type TrackedCommand = {
|
||||
label: string | null
|
||||
label: string | undefined
|
||||
running: boolean
|
||||
turnId: string | null
|
||||
}
|
||||
|
||||
/** Normalized once at receipt, not per projection: the roster is re-projected on
|
||||
* every observed frame, and the bound is a property of the stored value. */
|
||||
function description(value: string | null): string | undefined {
|
||||
if (value === null) {
|
||||
return undefined
|
||||
@@ -66,6 +68,14 @@ function description(value: string | null): string | undefined {
|
||||
return collapsed.length > 0 ? collapsed.slice(0, MAX_TASK_DESCRIPTION_CHARS) : undefined
|
||||
}
|
||||
|
||||
function task(
|
||||
id: string,
|
||||
kind: AgentSessionBackgroundTask['kind'],
|
||||
label: string | undefined
|
||||
): AgentSessionBackgroundTask {
|
||||
return { id, kind, ...(label === undefined ? {} : { description: label }) }
|
||||
}
|
||||
|
||||
export class CodexBackgroundTaskTracker {
|
||||
private readonly children = new Map<string, TrackedChild>()
|
||||
private readonly commands = new Map<string, TrackedCommand>()
|
||||
@@ -129,7 +139,7 @@ export class CodexBackgroundTaskTracker {
|
||||
// transition (`item/started` and `item/completed`), so a settled child must
|
||||
// not be resurrected by the duplicate.
|
||||
this.children.set(agentThreadId, {
|
||||
label: existing.label ?? label,
|
||||
label: existing.label ?? description(label),
|
||||
state: isTerminalSubagentState(existing.state) ? existing.state : state,
|
||||
turnId: existing.turnId ?? turnId
|
||||
})
|
||||
@@ -142,7 +152,7 @@ export class CodexBackgroundTaskTracker {
|
||||
) {
|
||||
return
|
||||
}
|
||||
this.children.set(agentThreadId, { label, state, turnId })
|
||||
this.children.set(agentThreadId, { label: description(label), state, turnId })
|
||||
}
|
||||
|
||||
private upsertCommand(
|
||||
@@ -154,7 +164,7 @@ export class CodexBackgroundTaskTracker {
|
||||
const existing = this.commands.get(itemId)
|
||||
if (existing) {
|
||||
this.commands.set(itemId, {
|
||||
label: label ?? existing.label,
|
||||
label: description(label) ?? existing.label,
|
||||
running,
|
||||
turnId: existing.turnId ?? turnId
|
||||
})
|
||||
@@ -163,7 +173,7 @@ export class CodexBackgroundTaskTracker {
|
||||
if (!this.makeRoom(this.commands, MAX_TRACKED_COMMANDS, (command) => !command.running)) {
|
||||
return
|
||||
}
|
||||
this.commands.set(itemId, { label, running, turnId })
|
||||
this.commands.set(itemId, { label: description(label), running, turnId })
|
||||
}
|
||||
|
||||
/** Frees a slot by dropping the oldest settled entry. Refuses to evict a live
|
||||
@@ -210,21 +220,13 @@ export class CodexBackgroundTaskTracker {
|
||||
if (isTerminalSubagentState(child.state) || !this.outlivedItsTurn(child.turnId)) {
|
||||
continue
|
||||
}
|
||||
tasks.push({
|
||||
id: `${AGENT_TASK_PREFIX}${agentThreadId}`,
|
||||
kind: 'agent',
|
||||
...(description(child.label) ? { description: description(child.label) } : {})
|
||||
})
|
||||
tasks.push(task(`${AGENT_TASK_PREFIX}${agentThreadId}`, 'agent', child.label))
|
||||
}
|
||||
for (const [itemId, command] of this.commands) {
|
||||
if (!command.running || !this.outlivedItsTurn(command.turnId)) {
|
||||
continue
|
||||
}
|
||||
tasks.push({
|
||||
id: `${COMMAND_TASK_PREFIX}${itemId}`,
|
||||
kind: 'command',
|
||||
...(description(command.label) ? { description: description(command.label) } : {})
|
||||
})
|
||||
tasks.push(task(`${COMMAND_TASK_PREFIX}${itemId}`, 'command', command.label))
|
||||
}
|
||||
return tasks
|
||||
}
|
||||
|
||||
@@ -16,11 +16,7 @@ import type { CodexJournalTranslationAdmission } from './codex-structured-journa
|
||||
import { answerCodexPrompt } from './codex-structured-prompt-replies'
|
||||
import { dispatchCodexTurn, isCodexTurnOptionKey } from './codex-structured-turn-start'
|
||||
import { supportsCodexStructuredLocation } from './codex-structured-location-support'
|
||||
import {
|
||||
closeAllCodexSessions,
|
||||
closeCodexPublishedSession,
|
||||
closeCodexSession
|
||||
} from './codex-structured-session-close'
|
||||
import { CodexStructuredSessionTeardown } from './codex-structured-session-teardown'
|
||||
import {
|
||||
applyCodexStructuredSessionOption,
|
||||
readLiveCodexSessionOptions
|
||||
@@ -54,6 +50,7 @@ export class CodexStructuredSessionAdapter implements StructuredAgentSessionAdap
|
||||
private readonly acquisitions = new CodexAcquisitionRegistry()
|
||||
private readonly turnCancellation: CodexStructuredTurnCancellation
|
||||
private readonly notificationRetries: ReturnType<typeof createCodexStructuredNotificationRetry>
|
||||
private readonly teardown: CodexStructuredSessionTeardown
|
||||
|
||||
constructor(private readonly deps: CodexStructuredSessionAdapterDeps) {
|
||||
this.notificationRetries = createCodexStructuredNotificationRetry({
|
||||
@@ -61,6 +58,15 @@ export class CodexStructuredSessionAdapter implements StructuredAgentSessionAdap
|
||||
translate: (sessionId, session, method, params) =>
|
||||
this.translateNotification(sessionId, session, method, params)
|
||||
})
|
||||
this.teardown = new CodexStructuredSessionTeardown({
|
||||
sessions: this.sessions,
|
||||
acquisitions: this.acquisitions,
|
||||
...(deps.onEvent ? { onEvent: deps.onEvent } : {}),
|
||||
...(deps.onBackgroundTasksChanged
|
||||
? { onBackgroundTasksChanged: deps.onBackgroundTasksChanged }
|
||||
: {}),
|
||||
forgetNotificationRetries: (sessionId) => this.notificationRetries.clear(sessionId, null)
|
||||
})
|
||||
this.turnCancellation = new CodexStructuredTurnCancellation({
|
||||
captureTurnProcesses: deps.captureTurnProcesses,
|
||||
terminateTurnProcesses: deps.terminateTurnProcesses,
|
||||
@@ -92,7 +98,7 @@ export class CodexStructuredSessionAdapter implements StructuredAgentSessionAdap
|
||||
handleUnhandledFrame: (sessionId, kind, payload) =>
|
||||
this.handleUnhandledFrame(sessionId, kind, payload),
|
||||
forceCloseUnexpected: (sessionId, fence, acquisitionGeneration, reason) =>
|
||||
this.forceCloseUnexpected(sessionId, fence, acquisitionGeneration, reason)
|
||||
this.teardown.forceCloseUnexpected(sessionId, fence, acquisitionGeneration, reason)
|
||||
})
|
||||
|
||||
/** Buffers pre-publication events and drops events from superseded children. */
|
||||
@@ -176,16 +182,6 @@ export class CodexStructuredSessionAdapter implements StructuredAgentSessionAdap
|
||||
sessionId
|
||||
) => this.sessions.get(sessionId)?.backgroundTasks.state
|
||||
|
||||
/** A closed session monitors nothing. Published explicitly because the state
|
||||
* reader answers `undefined` once the session leaves the map, which every
|
||||
* channel reads as "unchanged" and would leave the last roster on screen. */
|
||||
private clearBackgroundTasks(sessionId: string, closed: boolean): boolean {
|
||||
if (closed) {
|
||||
this.deps.onBackgroundTasksChanged?.(sessionId, null)
|
||||
}
|
||||
return closed
|
||||
}
|
||||
|
||||
bindPromptItemId = (sessionId: string, journalItemId: string, promptKey: string): void =>
|
||||
this.sessions
|
||||
.get(sessionId)
|
||||
@@ -286,59 +282,12 @@ export class CodexStructuredSessionAdapter implements StructuredAgentSessionAdap
|
||||
identity: AgentSessionJournalIdentity
|
||||
}): Promise<string | null> => this.sessions.get(input.identity.sessionId)?.historyPath ?? null
|
||||
|
||||
closeSession = async (sessionId: string): Promise<boolean> => {
|
||||
const closed = await closeCodexSession(
|
||||
sessionId,
|
||||
this.sessions,
|
||||
this.acquisitions,
|
||||
this.deps.onEvent
|
||||
)
|
||||
if (closed) {
|
||||
this.notificationRetries.clear(sessionId, null)
|
||||
}
|
||||
return this.clearBackgroundTasks(sessionId, closed)
|
||||
}
|
||||
forceCloseSession = async (sessionId: string): Promise<boolean> => {
|
||||
const closed = await closeCodexPublishedSession(this.sessions, sessionId, this.deps.onEvent, {
|
||||
allowFailedSettlement: true,
|
||||
requestedClose: false
|
||||
})
|
||||
if (closed) {
|
||||
this.notificationRetries.clear(sessionId, null)
|
||||
}
|
||||
return this.clearBackgroundTasks(sessionId, closed)
|
||||
}
|
||||
|
||||
private forceCloseUnexpected(
|
||||
sessionId: string,
|
||||
fence: number,
|
||||
acquisitionGeneration: string,
|
||||
reason: Error
|
||||
): Promise<boolean> {
|
||||
const session = this.sessions.get(sessionId)
|
||||
if (
|
||||
!session ||
|
||||
session.ended ||
|
||||
session.fence !== fence ||
|
||||
session.acquisitionGeneration !== acquisitionGeneration
|
||||
) {
|
||||
return Promise.resolve(false)
|
||||
}
|
||||
return closeCodexPublishedSession(this.sessions, sessionId, this.deps.onEvent, {
|
||||
allowFailedSettlement: true,
|
||||
requestedClose: false,
|
||||
expectedFence: fence,
|
||||
expectedAcquisitionGeneration: acquisitionGeneration,
|
||||
unexpectedReason: reason
|
||||
}).then((closed) => this.clearBackgroundTasks(sessionId, closed))
|
||||
}
|
||||
disposeSession = (sessionId: string): Promise<boolean> => this.closeSession(sessionId)
|
||||
closeAll = (): Promise<void> =>
|
||||
closeAllCodexSessions(this.sessions, this.acquisitions, (sessionId) =>
|
||||
this.disposeSession(sessionId)
|
||||
)
|
||||
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.closeSession(input.sessionId)
|
||||
this.teardown.close(input.sessionId)
|
||||
|
||||
private session(sessionId: string): CodexSession {
|
||||
return requireLiveCodexSession(this.sessions, sessionId)
|
||||
|
||||
@@ -0,0 +1,94 @@
|
||||
// Stopping one Codex app-server child, in the four ways the host asks for it.
|
||||
//
|
||||
// Every path funnels through `settled` so the ephemeral surfaces a closed
|
||||
// session owns are cleared exactly once, and only when the child was actually
|
||||
// proven stopped — a refused close leaves the session indexed for a retry.
|
||||
|
||||
import type { AgentSessionBackgroundTaskState } from '../../shared/agent-session-wire'
|
||||
import {
|
||||
closeAllCodexSessions,
|
||||
closeCodexPublishedSession,
|
||||
closeCodexSession
|
||||
} from './codex-structured-session-close'
|
||||
import type {
|
||||
CodexAcquisitionRegistry,
|
||||
CodexSession,
|
||||
CodexStructuredSessionEvent
|
||||
} from './codex-structured-session-state'
|
||||
|
||||
export type CodexStructuredSessionTeardownDeps = {
|
||||
sessions: Map<string, CodexSession>
|
||||
acquisitions: CodexAcquisitionRegistry
|
||||
onEvent?: (event: CodexStructuredSessionEvent) => void
|
||||
onBackgroundTasksChanged?: (
|
||||
sessionId: string,
|
||||
state: AgentSessionBackgroundTaskState | null
|
||||
) => void
|
||||
forgetNotificationRetries: (sessionId: string) => void
|
||||
}
|
||||
|
||||
export class CodexStructuredSessionTeardown {
|
||||
constructor(private readonly deps: CodexStructuredSessionTeardownDeps) {}
|
||||
|
||||
close = async (sessionId: string): Promise<boolean> => {
|
||||
const closed = await closeCodexSession(
|
||||
sessionId,
|
||||
this.deps.sessions,
|
||||
this.deps.acquisitions,
|
||||
this.deps.onEvent
|
||||
)
|
||||
return this.settled(sessionId, closed)
|
||||
}
|
||||
|
||||
forceClose = async (sessionId: string): Promise<boolean> => {
|
||||
const closed = await closeCodexPublishedSession(
|
||||
this.deps.sessions,
|
||||
sessionId,
|
||||
this.deps.onEvent,
|
||||
{ allowFailedSettlement: true, requestedClose: false }
|
||||
)
|
||||
return this.settled(sessionId, closed)
|
||||
}
|
||||
|
||||
/** Terminates this exact child as an unexpected death. Every ownership check
|
||||
* stays here so a stale caller cannot close a replacement child. */
|
||||
forceCloseUnexpected = (
|
||||
sessionId: string,
|
||||
fence: number,
|
||||
acquisitionGeneration: string,
|
||||
reason: Error
|
||||
): Promise<boolean> => {
|
||||
const session = this.deps.sessions.get(sessionId)
|
||||
if (
|
||||
!session ||
|
||||
session.ended ||
|
||||
session.fence !== fence ||
|
||||
session.acquisitionGeneration !== acquisitionGeneration
|
||||
) {
|
||||
return Promise.resolve(false)
|
||||
}
|
||||
return closeCodexPublishedSession(this.deps.sessions, sessionId, this.deps.onEvent, {
|
||||
allowFailedSettlement: true,
|
||||
requestedClose: false,
|
||||
expectedFence: fence,
|
||||
expectedAcquisitionGeneration: acquisitionGeneration,
|
||||
unexpectedReason: reason
|
||||
}).then((closed) => this.settled(sessionId, closed))
|
||||
}
|
||||
|
||||
closeAll = (): Promise<void> =>
|
||||
closeAllCodexSessions(this.deps.sessions, this.deps.acquisitions, (sessionId) =>
|
||||
this.close(sessionId)
|
||||
)
|
||||
|
||||
private settled(sessionId: string, closed: boolean): boolean {
|
||||
if (closed) {
|
||||
this.deps.forgetNotificationRetries(sessionId)
|
||||
// Explicit null, not silence: the state reader answers `undefined` once
|
||||
// the session leaves the map, which every channel reads as "unchanged"
|
||||
// and would leave the last roster on screen.
|
||||
this.deps.onBackgroundTasksChanged?.(sessionId, null)
|
||||
}
|
||||
return closed
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user