mirror of
https://github.com/stablyai/orca.git
synced 2026-09-29 16:02:50 +00:00
fix(native-chat): stop a foreign frame forfeiting a Codex conversation name
During the thread/start window the naming thread has no id yet, so the broad divert rule can catch a SUB-AGENT frame. A bare turn/completed caught that way settled the collector as a decline, which durably marks the conversation attempted — it could then never be named again. Only a frame matched to the retained throwaway id may now act on the turn; an unattributed one is still kept out of the journal and otherwise dropped. That makes the design-intent comment true rather than aspirational. The fake app-server also emitted the naming turn's frames at thread/start; a real app-server cannot emit for a turn that was never started, so they now arrive at turn/start, which is what makes them attributable.
This commit is contained in:
@@ -72,10 +72,10 @@ function run(
|
||||
// `null` leaves the turn unanswered entirely; 'DECLINE' completes it with
|
||||
// no message, which is a different fact the collector must distinguish.
|
||||
if (answer !== null && answer !== 'DECLINE') {
|
||||
collector?.handle('item/completed', { item: { type: 'agentMessage', text: answer } })
|
||||
collector?.handle('item/completed', { item: { type: 'agentMessage', text: answer } }, true)
|
||||
}
|
||||
if (answer !== null) {
|
||||
collector?.handle('turn/completed', {})
|
||||
collector?.handle('turn/completed', {}, true)
|
||||
}
|
||||
})
|
||||
})
|
||||
@@ -168,7 +168,7 @@ describe('createCodexNamingTurnCollector', () => {
|
||||
it('reports a completed turn that said nothing as a DECLINE', async () => {
|
||||
const collector = createCodexNamingTurnCollector(60_000)
|
||||
|
||||
collector.handle('turn/completed', {})
|
||||
collector.handle('turn/completed', {}, true)
|
||||
|
||||
// A model that completed and said nothing has answered. Classifying this as
|
||||
// a host failure would re-ask, and pay, on every future acquisition.
|
||||
@@ -178,10 +178,12 @@ describe('createCodexNamingTurnCollector', () => {
|
||||
it('reports the message a completed turn produced', async () => {
|
||||
const collector = createCodexNamingTurnCollector(60_000)
|
||||
|
||||
collector.handle('item/completed', {
|
||||
item: { type: 'agentMessage', text: '{"title":"Fix probe"}' }
|
||||
})
|
||||
collector.handle('turn/completed', {})
|
||||
collector.handle(
|
||||
'item/completed',
|
||||
{ item: { type: 'agentMessage', text: '{"title":"Fix probe"}' } },
|
||||
true
|
||||
)
|
||||
collector.handle('turn/completed', {}, true)
|
||||
|
||||
await expect(collector.answer).resolves.toEqual({
|
||||
outcome: 'answered',
|
||||
@@ -194,7 +196,7 @@ describe('createCodexNamingTurnCollector', () => {
|
||||
|
||||
// There is no `turn/failed` notification; a rate-limited or rejected turn
|
||||
// arrives as `error`. Without it this would hold for the whole timeout.
|
||||
collector.handle('error', { message: 'rate limit exceeded' })
|
||||
collector.handle('error', { message: 'rate limit exceeded' }, true)
|
||||
|
||||
await expect(collector.answer).resolves.toEqual({ outcome: 'failed' })
|
||||
})
|
||||
@@ -202,16 +204,33 @@ describe('createCodexNamingTurnCollector', () => {
|
||||
it('reports a failure even when the turn had already said something', async () => {
|
||||
const collector = createCodexNamingTurnCollector(60_000)
|
||||
|
||||
collector.handle('item/completed', {
|
||||
item: { type: 'agentMessage', text: '{"title":"Fix probe"}' }
|
||||
})
|
||||
collector.handle('error', { message: 'stream closed' })
|
||||
collector.handle(
|
||||
'item/completed',
|
||||
{ item: { type: 'agentMessage', text: '{"title":"Fix probe"}' } },
|
||||
true
|
||||
)
|
||||
collector.handle('error', { message: 'stream closed' }, true)
|
||||
|
||||
// The host is why there is no title; a partial answer does not make it a
|
||||
// decline the conversation should be marked for.
|
||||
await expect(collector.answer).resolves.toEqual({ outcome: 'failed' })
|
||||
})
|
||||
|
||||
it.each([
|
||||
['a bare completion', 'turn/completed', {}],
|
||||
['an agent message', 'item/completed', { item: { type: 'agentMessage', text: 'hi' } }],
|
||||
['a terminal error', 'error', { message: 'rate limit exceeded' }]
|
||||
])('ignores %s that could not be attributed to the naming thread', async (_l, method, params) => {
|
||||
const collector = createCodexNamingTurnCollector(1)
|
||||
|
||||
collector.handle(method, params, false)
|
||||
|
||||
// Inside the `thread/start` window the broad divert rule can catch a
|
||||
// SUB-AGENT frame. Settling on a bare completion would report a DECLINE,
|
||||
// which is durable — the conversation could never be named again.
|
||||
await expect(collector.answer).resolves.toEqual({ outcome: 'timed-out' })
|
||||
})
|
||||
|
||||
it('reports a turn that never answered as TIMED OUT', async () => {
|
||||
const collector = createCodexNamingTurnCollector(1)
|
||||
|
||||
@@ -227,7 +246,7 @@ describe('createCodexNamingTurnCollector', () => {
|
||||
const collector = createCodexNamingTurnCollector(60_000)
|
||||
expect(vi.getTimerCount()).toBe(1)
|
||||
|
||||
collector.handle(method, params)
|
||||
collector.handle(method, params, true)
|
||||
await collector.answer
|
||||
|
||||
// A pending timer holds the collector's closure for the whole timeout
|
||||
|
||||
@@ -98,10 +98,12 @@ export type CodexNamingState = {
|
||||
*
|
||||
* That window is one `thread/start` round trip on a healthy app-server, but it is
|
||||
* bounded by the request deadline, not instantaneous — a hung server stretches
|
||||
* it. A sub-agent frame arriving inside it is still diverted and can still be
|
||||
* recorded as the naming answer; that costs one wasted naming attempt rather
|
||||
* than a permanent forfeiture, because unparseable output classifies as
|
||||
* unsettled.
|
||||
* it. A sub-agent frame arriving inside it is diverted from the journal and then
|
||||
* DROPPED: the collector acts only on a frame matched to the retained throwaway
|
||||
* id, so nothing this rule alone caught can become the naming answer or settle
|
||||
* the turn. The cost is the sub-agent's own rows missing from the transcript for
|
||||
* that window — never a forfeited name, which a bare `turn/completed` diverted
|
||||
* here would otherwise cause by settling as a decline.
|
||||
*
|
||||
* Exact matching means leak protection now RELIES ON ID STABILITY: an app-server
|
||||
* that tagged naming frames with any id other than the one `thread/start`
|
||||
@@ -162,7 +164,9 @@ export type CodexNamingTurnResult =
|
||||
|
||||
/** Collects one ephemeral naming turn's frames and reports how it ended. */
|
||||
export type CodexNamingTurnCollector = {
|
||||
handle: (method: string, params: unknown) => void
|
||||
/** `attributed` is false for a frame only the broad pre-id divert rule caught;
|
||||
* it is kept out of the journal but may never speak for this turn. */
|
||||
handle: (method: string, params: unknown, attributed: boolean) => void
|
||||
answer: Promise<CodexNamingTurnResult>
|
||||
/** For a turn abandoned before `answer` is awaited, which would otherwise hold
|
||||
* the timer — and the closure behind it — for the whole timeout. */
|
||||
@@ -218,7 +222,13 @@ export function createCodexNamingTurnCollector(timeoutMs: number): CodexNamingTu
|
||||
})
|
||||
let latest: string | null = null
|
||||
return {
|
||||
handle: (method, params) => {
|
||||
handle: (method, params, attributed) => {
|
||||
// Only a frame matched to the retained throwaway id may speak for this
|
||||
// turn. A bare `turn/completed` from a sub-agent caught by the broad pre-id
|
||||
// rule would otherwise settle as a decline and durably forfeit naming.
|
||||
if (!attributed) {
|
||||
return
|
||||
}
|
||||
if (method === 'item/completed') {
|
||||
latest = agentMessageText(params) ?? latest
|
||||
} else if (method === 'turn/completed') {
|
||||
|
||||
@@ -139,8 +139,10 @@ export class CodexStructuredSessionAdapter implements StructuredAgentSessionAdap
|
||||
// Before anything can journal it: a naming turn runs on a throwaway thread
|
||||
// over this same connection, and the item translator journals items from ANY
|
||||
// thread. Routed here, its prompt and its JSON answer never reach the chat.
|
||||
if (isCodexNamingFrame(session, readCodexThreadId(params))) {
|
||||
session.naming?.handle(method, params)
|
||||
const frameThreadId = readCodexThreadId(params)
|
||||
if (isCodexNamingFrame(session, frameThreadId)) {
|
||||
// Diverted either way; only the exact-id half may settle the naming turn.
|
||||
session.naming?.handle(method, params, isCodexNamingThread(session, frameThreadId))
|
||||
return { accepted: true }
|
||||
}
|
||||
captureCodexConversationName(sessionId, session, method, params, this.deps)
|
||||
@@ -194,8 +196,9 @@ export class CodexStructuredSessionAdapter implements StructuredAgentSessionAdap
|
||||
|
||||
private handleUnhandledFrame(sessionId: string, kind: string, params: unknown): void {
|
||||
const session = this.sessions.get(sessionId)
|
||||
if (session && isCodexNamingFrame(session, readCodexThreadId(params))) {
|
||||
session.naming?.handle(kind, params)
|
||||
const frameThreadId = readCodexThreadId(params)
|
||||
if (session && isCodexNamingFrame(session, frameThreadId)) {
|
||||
session.naming?.handle(kind, params, isCodexNamingThread(session, frameThreadId))
|
||||
return
|
||||
}
|
||||
deliverCodexUnhandledFrame(sessionId, session, kind, params, (current, event) =>
|
||||
|
||||
@@ -46,7 +46,8 @@ export function handleCodexSessionExit(input: {
|
||||
// A naming turn in flight otherwise holds its collector until the 60s deadline
|
||||
// and then runs its cleanup against a dead connection. Settling it as a host
|
||||
// failure — not a decline — leaves the conversation askable on reacquisition.
|
||||
session.naming?.handle('error', {})
|
||||
// Attributed: the session owning the turn is what died, not a foreign thread.
|
||||
session.naming?.handle('error', {}, true)
|
||||
session.naming = null
|
||||
session.unbindReadingControl?.()
|
||||
input.onEvent?.(event)
|
||||
|
||||
@@ -132,6 +132,9 @@ function namingCodex(
|
||||
/** Never answers the ephemeral `thread/start`, so naming is in flight with
|
||||
* NO naming thread id known: the window the broad frame rule covers. */
|
||||
hangNamingThreadStart?: boolean
|
||||
/** Holds the ephemeral `thread/start` open until the test releases it, so a
|
||||
* frame can arrive inside that window and the flow still runs to the end. */
|
||||
holdNamingThreadStart?: boolean
|
||||
/** Completes the naming turn having said nothing: a genuine model decline,
|
||||
* which is a different fact from prose that ignored the schema. */
|
||||
declineNamingTurn?: boolean
|
||||
@@ -140,6 +143,10 @@ function namingCodex(
|
||||
const connections: FakeConnection[] = []
|
||||
const calls: { method: string; params: Record<string, unknown> }[] = []
|
||||
const replies: { id: number | string; result?: unknown; code?: number; message?: string }[] = []
|
||||
let releaseNamingThreadStart = (): void => {}
|
||||
const namingThreadStartGate = new Promise<void>((resolve) => {
|
||||
releaseNamingThreadStart = resolve
|
||||
})
|
||||
const openConnection = (async (
|
||||
_launch: CodexAppServerLaunch,
|
||||
handlers: CodexAppServerConnectionHandlers = {}
|
||||
@@ -154,24 +161,14 @@ function namingCodex(
|
||||
if (options.hangNamingThreadStart) {
|
||||
return await new Promise<never>(() => {})
|
||||
}
|
||||
if (options.holdNamingThreadStart) {
|
||||
await namingThreadStartGate
|
||||
}
|
||||
// Leaves the naming turn in flight: the thread id is known, but nothing
|
||||
// ever settles the collector, which is the window sub-agents run in.
|
||||
if (options.hangNamingTurn) {
|
||||
return { thread: { id: NAMING_THREAD, ephemeral: true } }
|
||||
}
|
||||
// The naming turn's frames arrive on this same connection.
|
||||
queueMicrotask(() => {
|
||||
if (!options.declineNamingTurn) {
|
||||
handlers.onNotification?.('item/completed', {
|
||||
threadId: NAMING_THREAD,
|
||||
item: {
|
||||
type: 'agentMessage',
|
||||
text: options.answer ?? '{"title":"Fix lease probe"}'
|
||||
}
|
||||
})
|
||||
}
|
||||
handlers.onNotification?.('turn/completed', { threadId: NAMING_THREAD })
|
||||
})
|
||||
return { thread: { id: NAMING_THREAD } }
|
||||
}
|
||||
if (method === 'thread/start') {
|
||||
@@ -186,6 +183,23 @@ function namingCodex(
|
||||
}
|
||||
}
|
||||
if (method === 'turn/start') {
|
||||
// The naming turn's frames arrive on this same connection, and only
|
||||
// once the turn exists — the app-server cannot emit for a turn that
|
||||
// was never started, which is what makes them attributable.
|
||||
if (params?.threadId === NAMING_THREAD) {
|
||||
queueMicrotask(() => {
|
||||
if (!options.declineNamingTurn) {
|
||||
handlers.onNotification?.('item/completed', {
|
||||
threadId: NAMING_THREAD,
|
||||
item: {
|
||||
type: 'agentMessage',
|
||||
text: options.answer ?? '{"title":"Fix lease probe"}'
|
||||
}
|
||||
})
|
||||
}
|
||||
handlers.onNotification?.('turn/completed', { threadId: NAMING_THREAD })
|
||||
})
|
||||
}
|
||||
return { turn: { id: 'turn-1' } }
|
||||
}
|
||||
return {}
|
||||
@@ -202,7 +216,7 @@ function namingCodex(
|
||||
connections.push(connection)
|
||||
return connection
|
||||
}) as typeof openCodexAppServerConnection
|
||||
return { connections, openConnection, calls, replies }
|
||||
return { connections, openConnection, calls, replies, releaseNamingThreadStart }
|
||||
}
|
||||
|
||||
const NAMING_THREAD = 'thread-naming'
|
||||
@@ -566,6 +580,31 @@ describe('Codex sub-agent threads survive the naming window', () => {
|
||||
])
|
||||
})
|
||||
|
||||
it('does not let a foreign turn/completed inside the window forfeit naming', async () => {
|
||||
const codex = namingCodex({ holdNamingThreadStart: true })
|
||||
const markNamingAttempted = vi.fn()
|
||||
const { onConversationName } = await dispatchedAdapter(codex, { markNamingAttempted })
|
||||
await settle()
|
||||
|
||||
// A sub-agent's BARE completion, arriving before the throwaway thread has an
|
||||
// id. Settling on it reports a decline, which is durable: the conversation
|
||||
// would carry `conversationNamingAttempted` with no name and could never be
|
||||
// named again.
|
||||
codex.connections[0]!.handlers.onNotification?.('turn/completed', {
|
||||
threadId: SUBAGENT_THREAD
|
||||
})
|
||||
await settle()
|
||||
|
||||
codex.releaseNamingThreadStart()
|
||||
await settle()
|
||||
|
||||
expect(codex.calls.find((call) => call.method === 'thread/name/set')?.params).toEqual({
|
||||
threadId: THREAD_ID,
|
||||
name: 'Fix lease probe'
|
||||
})
|
||||
expect(onConversationName).toHaveBeenCalledWith(SESSION, 'Fix lease probe')
|
||||
})
|
||||
|
||||
it('still keeps the naming thread out once its id is known', async () => {
|
||||
const codex = namingCodex({ hangNamingTurn: true })
|
||||
const { events } = await dispatchedAdapter(codex)
|
||||
|
||||
Reference in New Issue
Block a user