mirror of
https://github.com/stablyai/orca.git
synced 2026-09-30 00:03:15 +00:00
Isolate Codex conversation naming in an owned short-lived process
This commit is contained in:
@@ -5,6 +5,8 @@ export type CodexAppServerServerRequest = {
|
||||
}
|
||||
|
||||
export type CodexAppServerConnectionHandlers = {
|
||||
/** Publishes ownership before the handshake so cancellation can reap the child. */
|
||||
onConnection?: (connection: CodexAppServerConnection) => void
|
||||
onNotification?: (method: string, params: unknown) => void
|
||||
onServerRequest?: (request: CodexAppServerServerRequest) => void
|
||||
onUnhandledFrame?: (kind: string, payload: unknown) => void
|
||||
|
||||
@@ -165,6 +165,26 @@ function responseLine(targetBytes: number, id: number): string {
|
||||
}
|
||||
|
||||
describe('openCodexAppServerConnection', () => {
|
||||
it('publishes the child before initialize so cancellation can stop a stalled handshake', async () => {
|
||||
const { child, spawnImpl, written } = stubChild()
|
||||
let owned: CodexAppServerConnection | undefined
|
||||
const opening = openCodexAppServerConnection(
|
||||
{ command: 'codex', args: ['app-server'] },
|
||||
{
|
||||
onConnection: (connection) => {
|
||||
owned = connection
|
||||
}
|
||||
},
|
||||
spawnImpl
|
||||
)
|
||||
const failed = rejection(opening)
|
||||
expect(owned).toBeDefined()
|
||||
expect(written[0]?.method).toBe('initialize')
|
||||
await expect(owned!.close()).resolves.toBe(true)
|
||||
await expect(failed).resolves.toBeInstanceOf(Error)
|
||||
child.stdout.destroy()
|
||||
})
|
||||
|
||||
it('advertises the experimental API required for rollout-path resume', async () => {
|
||||
const { child, spawnImpl, written } = stubChild()
|
||||
answerInitialize(child)
|
||||
|
||||
@@ -281,6 +281,7 @@ export async function openCodexAppServerConnection(
|
||||
}
|
||||
|
||||
try {
|
||||
handlers.onConnection?.(connection)
|
||||
await initializeCodexAppServerConnection(connection)
|
||||
} catch (error) {
|
||||
if ((await close()) !== true) {
|
||||
|
||||
@@ -1,527 +1,186 @@
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
import type { CodexAppServerConnection } from './codex-app-server-connection'
|
||||
import { describe, expect, it, vi } from 'vitest'
|
||||
import {
|
||||
createCodexNamingTurnCollector,
|
||||
isTerminalCodexTurnError,
|
||||
generateAndSetCodexConversationName,
|
||||
readCodexGeneratedTitle
|
||||
} from './codex-conversation-name-generation'
|
||||
|
||||
const THREAD = 'thread-user'
|
||||
const NAMING = 'thread-naming'
|
||||
|
||||
/** An app-server that answers the naming flow's four requests. */
|
||||
function fakeConnection(
|
||||
options: {
|
||||
nameOnReRead?: string | null
|
||||
answer?: string | null
|
||||
/** The app-server ignored `ephemeral` and persisted the throwaway thread. */
|
||||
persistDespiteEphemeral?: boolean
|
||||
/** A `thread/read` reply shape this build may or may not understand. */
|
||||
threadReadReply?: unknown
|
||||
} = {}
|
||||
function generation(
|
||||
options: { read?: unknown; opened?: unknown; config?: unknown; answer?: string | null } = {}
|
||||
) {
|
||||
const calls: { method: string; params: Record<string, unknown> }[] = []
|
||||
const connection: Pick<CodexAppServerConnection, 'request'> = {
|
||||
request: vi.fn(async (method: string, params?: Record<string, unknown>) => {
|
||||
calls.push({ method, params: params ?? {} })
|
||||
if (method === 'thread/start') {
|
||||
return {
|
||||
thread: {
|
||||
id: NAMING,
|
||||
...(options.persistDespiteEphemeral ? {} : { ephemeral: params?.ephemeral === true })
|
||||
}
|
||||
}
|
||||
const collector = createCodexNamingTurnCollector(1000)
|
||||
const connection = {
|
||||
request: vi.fn(async (method: string) => {
|
||||
if (method === 'config/read') {
|
||||
return options.config ?? { config: { mcp_servers: { files: { command: 'server' } } } }
|
||||
}
|
||||
if (method === 'thread/read') {
|
||||
return options.threadReadReply !== undefined
|
||||
? options.threadReadReply
|
||||
: {
|
||||
thread: {
|
||||
id: THREAD,
|
||||
...(options.nameOnReRead ? { name: options.nameOnReRead } : {})
|
||||
}
|
||||
}
|
||||
if (method === 'thread/start') {
|
||||
return options.opened ?? { thread: { id: 'naming', ephemeral: true } }
|
||||
}
|
||||
if (method === 'turn/start') {
|
||||
if (options.answer !== null) {
|
||||
collector.handle('item/completed', {
|
||||
item: { type: 'agentMessage', text: options.answer ?? '{"title":"Fix lease probe"}' }
|
||||
})
|
||||
}
|
||||
collector.handle('turn/completed', {})
|
||||
}
|
||||
return {}
|
||||
})
|
||||
}
|
||||
return { connection, calls }
|
||||
const userConnection = {
|
||||
request: vi.fn(async () => options.read ?? { thread: { id: 'user', name: null } })
|
||||
}
|
||||
const run = () =>
|
||||
generateAndSetCodexConversationName({
|
||||
connection,
|
||||
userConnection,
|
||||
collector,
|
||||
cwd: '/folder',
|
||||
threadId: 'user',
|
||||
prompt: 'fix lease probe',
|
||||
model: 'selected-model',
|
||||
isCancelled: () => false
|
||||
}).finally(() => collector.dispose())
|
||||
return { run, connection, userConnection }
|
||||
}
|
||||
|
||||
function run(
|
||||
connection: Pick<CodexAppServerConnection, 'request'>,
|
||||
answer: string | null
|
||||
): Promise<{ name: string | null; settled: boolean }> {
|
||||
let collector: ReturnType<typeof createCodexNamingTurnCollector> | null = null
|
||||
const done = generateAndSetCodexConversationName({
|
||||
connection,
|
||||
cwd: '/work/repo',
|
||||
threadId: THREAD,
|
||||
prompt: 'fix the flaky lease probe',
|
||||
openNamingTurn: () => {
|
||||
collector = createCodexNamingTurnCollector(5_000)
|
||||
return collector
|
||||
},
|
||||
retainNamingThread: () => {},
|
||||
closeNamingTurn: () => {}
|
||||
describe('Codex conversation naming generation', () => {
|
||||
it('generates independently and publishes only the final name on the user connection', async () => {
|
||||
const { run, connection, userConnection } = generation()
|
||||
await expect(run()).resolves.toEqual({ name: 'Fix lease probe', settled: true })
|
||||
expect(connection.request.mock.calls.map(([method]) => method)).toEqual([
|
||||
'config/read',
|
||||
'thread/start',
|
||||
'turn/start'
|
||||
])
|
||||
expect(userConnection.request.mock.calls).toEqual([
|
||||
['thread/read', { threadId: 'user' }, { timeoutMs: undefined }],
|
||||
['thread/name/set', { threadId: 'user', name: 'Fix lease probe' }, { timeoutMs: undefined }]
|
||||
])
|
||||
})
|
||||
// Drive the turn the way the app-server would, once the flow has opened it.
|
||||
queueMicrotask(() => {
|
||||
queueMicrotask(() => {
|
||||
// `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 } }, true)
|
||||
}
|
||||
if (answer !== null) {
|
||||
collector?.handle('turn/completed', {}, true)
|
||||
}
|
||||
})
|
||||
})
|
||||
return done
|
||||
}
|
||||
|
||||
describe('readCodexGeneratedTitle', () => {
|
||||
it('reads the title out of the structured answer', () => {
|
||||
expect(readCodexGeneratedTitle('{"title":"Fix flaky lease probe"}')).toBe(
|
||||
'Fix flaky lease probe'
|
||||
it('disables effective MCP servers and tools and preserves the selected model', async () => {
|
||||
const { run, connection } = generation()
|
||||
await run()
|
||||
expect(connection.request).toHaveBeenCalledWith(
|
||||
'thread/start',
|
||||
expect.objectContaining({
|
||||
model: 'selected-model',
|
||||
ephemeral: true,
|
||||
sandbox: 'read-only',
|
||||
approvalPolicy: 'never',
|
||||
dynamicTools: [],
|
||||
environments: [],
|
||||
runtimeWorkspaceRoots: [],
|
||||
selectedCapabilityRoots: [],
|
||||
config: expect.objectContaining({
|
||||
mcp_servers: { files: { enabled: false } },
|
||||
'features.shell_tool': false,
|
||||
'features.unified_exec': false,
|
||||
'features.multi_agent': false,
|
||||
'features.plugins': false,
|
||||
web_search: 'disabled'
|
||||
})
|
||||
}),
|
||||
expect.anything()
|
||||
)
|
||||
})
|
||||
|
||||
it('reports null for prose, empty answers, and a blank title', () => {
|
||||
expect(readCodexGeneratedTitle('Sure! Here is a title.')).toBeNull()
|
||||
expect(readCodexGeneratedTitle('')).toBeNull()
|
||||
expect(readCodexGeneratedTitle(null)).toBeNull()
|
||||
expect(readCodexGeneratedTitle('{"title":" "}')).toBeNull()
|
||||
expect(readCodexGeneratedTitle('{"name":"Fix it"}')).toBeNull()
|
||||
it.each([
|
||||
{ thread: { id: 'user', name: 'My own name' } },
|
||||
{ threadId: 'user' },
|
||||
{ thread: { id: 'another', name: null } },
|
||||
'unreadable'
|
||||
])('does not overwrite a named or unreadable thread: %j', async (read) => {
|
||||
const { run, userConnection } = generation({ read })
|
||||
await expect(run()).resolves.toEqual({ name: null, settled: true })
|
||||
expect(userConnection.request).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
})
|
||||
|
||||
describe('naming thread hardening', () => {
|
||||
it('opens the throwaway thread with no approvals and no write access', async () => {
|
||||
const { connection, calls } = fakeConnection()
|
||||
it('deletes and refuses a naming thread when the host ignored ephemeral', async () => {
|
||||
const { run, connection, userConnection } = generation({ opened: { thread: { id: 'naming' } } })
|
||||
await expect(run()).resolves.toEqual({ name: null, settled: false })
|
||||
expect(connection.request).toHaveBeenCalledWith(
|
||||
'thread/delete',
|
||||
{ threadId: 'naming' },
|
||||
expect.anything()
|
||||
)
|
||||
expect(userConnection.request).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
await run(connection, '{"title":"Fix lease probe"}')
|
||||
it('never runs or deletes the user thread when start returns its identity', async () => {
|
||||
const { run, connection } = generation({ opened: { thread: { id: 'user' } } })
|
||||
await expect(run()).resolves.toEqual({ name: null, settled: false })
|
||||
expect(connection.request.mock.calls.map(([method]) => method)).toEqual([
|
||||
'config/read',
|
||||
'thread/start'
|
||||
])
|
||||
})
|
||||
|
||||
// The prompt embeds untrusted user text and this thread's frames never reach
|
||||
// the journal, so anything the host would auto-approve would run unseen.
|
||||
expect(calls.find((call) => call.method === 'thread/start')?.params).toEqual({
|
||||
cwd: '/work/repo',
|
||||
ephemeral: true,
|
||||
approvalPolicy: 'never',
|
||||
sandbox: 'read-only'
|
||||
it('fails closed when effective configuration cannot be read', async () => {
|
||||
const { run, connection } = generation({ config: {} })
|
||||
await expect(run()).rejects.toThrow('readable effective configuration')
|
||||
expect(connection.request).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it.each([
|
||||
[null, true],
|
||||
['prose', false],
|
||||
['{"title":""}', false]
|
||||
])('accounts for declines and unusable responses: %s', async (answer, settled) => {
|
||||
await expect(generation({ answer }).run()).resolves.toEqual({ name: null, settled })
|
||||
})
|
||||
|
||||
it('bounds and flattens provider names', async () => {
|
||||
const { run } = generation({
|
||||
answer: JSON.stringify({ title: `Fix\nprobe ${'x'.repeat(400)}` })
|
||||
})
|
||||
})
|
||||
|
||||
it('releases the throwaway thread with the protocol cleanup', async () => {
|
||||
const { connection, calls } = fakeConnection()
|
||||
|
||||
await run(connection, '{"title":"Fix lease probe"}')
|
||||
|
||||
expect(calls.find((call) => call.method === 'thread/unsubscribe')?.params).toEqual({
|
||||
threadId: NAMING
|
||||
})
|
||||
})
|
||||
|
||||
it('never unsubscribes the user own thread when the reply named it', async () => {
|
||||
const calls: { method: string; params: Record<string, unknown> }[] = []
|
||||
const connection: Pick<CodexAppServerConnection, 'request'> = {
|
||||
request: vi.fn(async (method: string, params?: Record<string, unknown>) => {
|
||||
calls.push({ method, params: params ?? {} })
|
||||
// The reply names the SESSION's thread; unsubscribing it would cut the
|
||||
// user's chat off from every frame it depends on.
|
||||
return method === 'thread/start' ? { thread: { id: THREAD } } : {}
|
||||
})
|
||||
}
|
||||
|
||||
await expect(run(connection, '{"title":"Fix lease probe"}')).resolves.toEqual({
|
||||
name: null,
|
||||
settled: false
|
||||
})
|
||||
expect(calls.map((call) => call.method)).not.toContain('thread/unsubscribe')
|
||||
})
|
||||
|
||||
it('refuses to parse an answer far larger than any title', async () => {
|
||||
// The 36-character cap is a schema request to the model, not a bound the
|
||||
// host enforces on the reply.
|
||||
const huge = `{"title":"${'a'.repeat(9 * 1024)}"}`
|
||||
|
||||
expect(readCodexGeneratedTitle(huge)).toBeNull()
|
||||
})
|
||||
|
||||
it('flattens and bounds the name before it reaches the user Codex thread', async () => {
|
||||
const { connection, calls } = fakeConnection()
|
||||
const sprawling = `Fix\nthe lease probe ${'x'.repeat(400)}`
|
||||
|
||||
await run(connection, JSON.stringify({ title: sprawling }))
|
||||
|
||||
const set = calls.find((call) => call.method === 'thread/name/set')
|
||||
const name = String(set?.params.name)
|
||||
const { name } = await run()
|
||||
expect(name).not.toContain('\n')
|
||||
expect(name.length).toBeLessThanOrEqual(200)
|
||||
expect(name.startsWith('Fix the lease probe ')).toBe(true)
|
||||
expect(name!.length).toBeLessThanOrEqual(200)
|
||||
})
|
||||
|
||||
it.each([null, '', 'prose', '{"title":" "}', JSON.stringify({ title: 'a'.repeat(9000) })])(
|
||||
'rejects unusable structured answers: %s',
|
||||
(answer) => {
|
||||
expect(readCodexGeneratedTitle(answer)).toBeNull()
|
||||
}
|
||||
)
|
||||
})
|
||||
|
||||
describe('createCodexNamingTurnCollector', () => {
|
||||
it('reports a completed turn that said nothing as a DECLINE', async () => {
|
||||
const collector = createCodexNamingTurnCollector(60_000)
|
||||
describe('Codex naming response collection', () => {
|
||||
it.each(['failed', 'interrupted'])(
|
||||
'rejects a partial response from a %s turn',
|
||||
async (status) => {
|
||||
const collector = createCodexNamingTurnCollector(1000)
|
||||
collector.handle('item/completed', {
|
||||
item: { type: 'agentMessage', text: '{"title":"Partial"}' }
|
||||
})
|
||||
collector.handle('turn/completed', { turn: { status } })
|
||||
await expect(collector.answer).resolves.toEqual({ outcome: 'failed' })
|
||||
}
|
||||
)
|
||||
|
||||
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.
|
||||
await expect(collector.answer).resolves.toEqual({ outcome: 'declined' })
|
||||
})
|
||||
|
||||
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"}' } },
|
||||
true
|
||||
)
|
||||
collector.handle('turn/completed', {}, true)
|
||||
|
||||
await expect(collector.answer).resolves.toEqual({
|
||||
outcome: 'answered',
|
||||
text: '{"title":"Fix probe"}'
|
||||
it('continues after a retryable error and rejects a terminal failure', async () => {
|
||||
const collector = createCodexNamingTurnCollector(1000)
|
||||
collector.handle('error', { willRetry: true })
|
||||
collector.handle('item/completed', {
|
||||
item: { type: 'agentMessage', text: '{"title":"partial"}' }
|
||||
})
|
||||
})
|
||||
|
||||
it('reports a terminal error as a FAILURE, not a decline', async () => {
|
||||
const collector = createCodexNamingTurnCollector(60_000)
|
||||
|
||||
// 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' }, true)
|
||||
|
||||
collector.handle('error', { willRetry: false })
|
||||
await expect(collector.answer).resolves.toEqual({ outcome: 'failed' })
|
||||
})
|
||||
|
||||
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"}' } },
|
||||
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)
|
||||
|
||||
await expect(collector.answer).resolves.toEqual({ outcome: 'timed-out' })
|
||||
})
|
||||
|
||||
it.each([
|
||||
['a completed turn', 'turn/completed', {}],
|
||||
['a terminal error', 'error', { message: 'rate limit exceeded' }]
|
||||
])('drops its timeout once %s settles it', async (_label, method, params) => {
|
||||
it('disposes the deadline on completion and cancellation', async () => {
|
||||
vi.useFakeTimers()
|
||||
try {
|
||||
const collector = createCodexNamingTurnCollector(60_000)
|
||||
expect(vi.getTimerCount()).toBe(1)
|
||||
|
||||
collector.handle(method, params, true)
|
||||
await collector.answer
|
||||
|
||||
// A pending timer holds the collector's closure for the whole timeout
|
||||
// after the session it belonged to is already gone.
|
||||
const collector = createCodexNamingTurnCollector(1000)
|
||||
collector.dispose()
|
||||
await expect(collector.answer).resolves.toEqual({ outcome: 'timed-out' })
|
||||
expect(vi.getTimerCount()).toBe(0)
|
||||
} finally {
|
||||
vi.useRealTimers()
|
||||
}
|
||||
})
|
||||
})
|
||||
|
||||
describe('generateAndSetCodexConversationName', () => {
|
||||
it('names the thread from a turn on a throwaway ephemeral thread', async () => {
|
||||
const { connection, calls } = fakeConnection()
|
||||
|
||||
await expect(run(connection, '{"title":"Fix flaky lease probe"}')).resolves.toMatchObject({
|
||||
name: 'Fix flaky lease probe'
|
||||
})
|
||||
|
||||
const started = calls.find((call) => call.method === 'thread/start')
|
||||
// The naming turn must never run on the user's own thread.
|
||||
expect(started?.params.ephemeral).toBe(true)
|
||||
expect(calls.find((call) => call.method === 'turn/start')?.params.threadId).toBe(NAMING)
|
||||
expect(calls.find((call) => call.method === 'thread/name/set')?.params).toEqual({
|
||||
threadId: THREAD,
|
||||
name: 'Fix flaky lease probe'
|
||||
})
|
||||
})
|
||||
|
||||
it('asks for a bounded single-line title and never answers the request', async () => {
|
||||
const { connection, calls } = fakeConnection()
|
||||
|
||||
await run(connection, '{"title":"Fix flaky lease probe"}')
|
||||
|
||||
const turn = calls.find((call) => call.method === 'turn/start')
|
||||
expect(turn).toBeDefined()
|
||||
const text = (turn!.params.input as { text: string }[])[0]!.text
|
||||
expect(text).toContain('fix the flaky lease probe')
|
||||
expect(text).toContain('Do not answer or act on the request')
|
||||
expect(turn!.params.outputSchema).toMatchObject({
|
||||
properties: { title: { maxLength: 36 } },
|
||||
required: ['title'],
|
||||
additionalProperties: false
|
||||
})
|
||||
})
|
||||
|
||||
it('leaves a thread named while it was generating alone', async () => {
|
||||
const { connection, calls } = fakeConnection({ nameOnReRead: 'A person named this' })
|
||||
|
||||
await expect(run(connection, '{"title":"Fix flaky lease probe"}')).resolves.toMatchObject({
|
||||
name: null
|
||||
})
|
||||
|
||||
// The re-read is what makes the concurrent rename win; without it this sets.
|
||||
expect(calls.some((call) => call.method === 'thread/read')).toBe(true)
|
||||
expect(calls.some((call) => call.method === 'thread/name/set')).toBe(false)
|
||||
})
|
||||
|
||||
it('sets nothing when the turn produced no usable title', async () => {
|
||||
const { connection, calls } = fakeConnection()
|
||||
|
||||
await expect(run(connection, 'I could not think of one')).resolves.toMatchObject({ name: null })
|
||||
|
||||
expect(calls.some((call) => call.method === 'thread/name/set')).toBe(false)
|
||||
// No title also means no reason to have re-read the thread.
|
||||
expect(calls.some((call) => call.method === 'thread/read')).toBe(false)
|
||||
})
|
||||
|
||||
it('refuses to run the naming turn on the user’s own thread', async () => {
|
||||
// An app-server that ignored `ephemeral` hands back the session's own
|
||||
// thread. Running there would put this prompt and its JSON answer into the
|
||||
// user's transcript and their history — the one outcome this path exists to
|
||||
// avoid — so the turn is abandoned instead.
|
||||
const calls: { method: string; params: Record<string, unknown> }[] = []
|
||||
const connection: Pick<CodexAppServerConnection, 'request'> = {
|
||||
request: vi.fn(async (method: string, params?: Record<string, unknown>) => {
|
||||
calls.push({ method, params: params ?? {} })
|
||||
return method === 'thread/start' ? { thread: { id: THREAD } } : {}
|
||||
})
|
||||
}
|
||||
|
||||
await expect(run(connection, '{"title":"Fix flaky lease probe"}')).resolves.toMatchObject({
|
||||
name: null
|
||||
})
|
||||
|
||||
expect(calls.map((call) => call.method)).toEqual(['thread/start'])
|
||||
})
|
||||
|
||||
it('gives up when the turn ends without an answer', async () => {
|
||||
const { connection, calls } = fakeConnection()
|
||||
|
||||
await expect(run(connection, null)).resolves.toMatchObject({ name: null })
|
||||
|
||||
expect(calls.some((call) => call.method === 'thread/name/set')).toBe(false)
|
||||
})
|
||||
})
|
||||
|
||||
describe('naming-turn cleanup and fail-closed re-read', () => {
|
||||
it('deletes a throwaway thread the app-server persisted despite ephemeral', async () => {
|
||||
const { connection, calls } = fakeConnection({ persistDespiteEphemeral: true })
|
||||
|
||||
await run(connection, '{"title":"Fix flaky lease probe"}')
|
||||
|
||||
// Otherwise every named chat leaves a junk thread and rollout file behind.
|
||||
expect(calls.find((call) => call.method === 'thread/delete')?.params).toEqual({
|
||||
threadId: NAMING
|
||||
})
|
||||
})
|
||||
|
||||
it('does not try to delete a genuinely ephemeral thread', async () => {
|
||||
const { connection, calls } = fakeConnection()
|
||||
|
||||
await run(connection, '{"title":"Fix flaky lease probe"}')
|
||||
|
||||
// The app-server refuses to delete one, so attempting it would log a failure
|
||||
// on every successful naming.
|
||||
expect(calls.some((call) => call.method === 'thread/delete')).toBe(false)
|
||||
})
|
||||
|
||||
// NOTE: a name under an unrecognised KEY on an otherwise-readable reply is not
|
||||
// detectable — by construction this build does not know the key. What fails
|
||||
// closed is an unrecognised reply SHAPE, which is the case it can decide.
|
||||
it.each([
|
||||
['a reply with no thread object', { threadId: 'thread-user' }],
|
||||
['a reply that is not an object', 'thread-user']
|
||||
])(
|
||||
'refuses to overwrite a name it cannot positively read as absent: %s',
|
||||
async (_label, reply) => {
|
||||
const { connection, calls } = fakeConnection({ threadReadReply: reply })
|
||||
|
||||
await expect(run(connection, '{"title":"Fix flaky lease probe"}')).resolves.toMatchObject({
|
||||
name: null
|
||||
})
|
||||
|
||||
// Fails CLOSED. Skipping a name is a non-event; clobbering a rename is not.
|
||||
expect(calls.some((call) => call.method === 'thread/name/set')).toBe(false)
|
||||
}
|
||||
)
|
||||
|
||||
it('still names a thread a readable reply shows as unnamed', async () => {
|
||||
const { connection, calls } = fakeConnection()
|
||||
|
||||
await expect(run(connection, '{"title":"Fix flaky lease probe"}')).resolves.toMatchObject({
|
||||
name: 'Fix flaky lease probe'
|
||||
})
|
||||
|
||||
expect(calls.find((call) => call.method === 'thread/name/set')?.params).toEqual({
|
||||
threadId: THREAD,
|
||||
name: 'Fix flaky lease probe'
|
||||
})
|
||||
})
|
||||
})
|
||||
|
||||
describe('isTerminalCodexTurnError', () => {
|
||||
it('ends the turn on a non-retryable error', () => {
|
||||
expect(isTerminalCodexTurnError('error', { threadId: 't', willRetry: false })).toBe(true)
|
||||
})
|
||||
|
||||
it('does NOT end the turn on a retryable one', () => {
|
||||
// A retryable error explicitly does not interrupt the turn; settling here
|
||||
// would abandon a naming turn that was about to succeed.
|
||||
expect(isTerminalCodexTurnError('error', { threadId: 't', willRetry: true })).toBe(false)
|
||||
expect(isTerminalCodexTurnError('error', { threadId: 't', will_retry: true })).toBe(false)
|
||||
})
|
||||
|
||||
it('ignores anything that is not an error frame', () => {
|
||||
expect(isTerminalCodexTurnError('turn/started', { willRetry: false })).toBe(false)
|
||||
expect(isTerminalCodexTurnError('error', null)).toBe(false)
|
||||
})
|
||||
})
|
||||
|
||||
describe('settled vs unsettled outcomes', () => {
|
||||
it('is UNSETTLED when no throwaway thread could be opened', async () => {
|
||||
const connection: Pick<CodexAppServerConnection, 'request'> = {
|
||||
request: vi.fn(async () => ({}))
|
||||
}
|
||||
|
||||
// A host that cannot be asked must stay askable; marking it would forfeit
|
||||
// naming for this conversation permanently.
|
||||
await expect(run(connection, '{"title":"x"}')).resolves.toEqual({
|
||||
name: null,
|
||||
settled: false
|
||||
})
|
||||
})
|
||||
|
||||
it('is UNSETTLED when the turn never answered', async () => {
|
||||
const { connection } = fakeConnection()
|
||||
|
||||
await expect(run(connection, null)).resolves.toEqual({ name: null, settled: false })
|
||||
})
|
||||
|
||||
it('is SETTLED when the model completed and said nothing', async () => {
|
||||
const { connection } = fakeConnection()
|
||||
|
||||
// A real decline. Left unsettled, this is re-asked — and paid for — on every
|
||||
// future acquisition of the conversation.
|
||||
await expect(run(connection, 'DECLINE')).resolves.toEqual({ name: null, settled: true })
|
||||
})
|
||||
|
||||
it('is UNSETTLED when the model ignored the schema and answered in prose', async () => {
|
||||
const { connection } = fakeConnection()
|
||||
|
||||
// Marking this would make the conversation permanently unnameable, even
|
||||
// after switching to a model that honours the schema.
|
||||
await expect(run(connection, 'Sure! A good title would be "Fix probe".')).resolves.toEqual({
|
||||
name: null,
|
||||
settled: false
|
||||
})
|
||||
})
|
||||
|
||||
it('is SETTLED when it named the thread', async () => {
|
||||
const { connection } = fakeConnection()
|
||||
|
||||
await expect(run(connection, '{"title":"Fix flaky lease probe"}')).resolves.toEqual({
|
||||
name: 'Fix flaky lease probe',
|
||||
settled: true
|
||||
})
|
||||
})
|
||||
})
|
||||
|
||||
describe('naming turn timer disposal', () => {
|
||||
beforeEach(() => {
|
||||
vi.useFakeTimers()
|
||||
})
|
||||
|
||||
afterEach(() => {
|
||||
vi.useRealTimers()
|
||||
})
|
||||
|
||||
/** Runs the flow with a real collector, exposing it so a vacuous pass is visible. */
|
||||
function runExposingCollector(connection: Pick<CodexAppServerConnection, 'request'>) {
|
||||
let collector: ReturnType<typeof createCodexNamingTurnCollector> | null = null
|
||||
const done = generateAndSetCodexConversationName({
|
||||
connection,
|
||||
cwd: '/work/repo',
|
||||
threadId: THREAD,
|
||||
prompt: 'fix the flaky lease probe',
|
||||
openNamingTurn: () => {
|
||||
collector = createCodexNamingTurnCollector(60_000)
|
||||
return collector
|
||||
},
|
||||
retainNamingThread: () => {},
|
||||
closeNamingTurn: () => {}
|
||||
}).catch(() => null)
|
||||
return { done, collector: () => collector }
|
||||
}
|
||||
|
||||
const refusing = (method: string): Pick<CodexAppServerConnection, 'request'> => ({
|
||||
request: vi.fn(async (called: string) => {
|
||||
if (called === method) {
|
||||
throw new Error(`${method} refused`)
|
||||
}
|
||||
return called === 'thread/start' ? { thread: { id: NAMING, ephemeral: true } } : {}
|
||||
})
|
||||
})
|
||||
|
||||
it.each([
|
||||
['thread/start is refused', refusing('thread/start')],
|
||||
['turn/start is refused', refusing('turn/start')],
|
||||
[
|
||||
'the opened thread is the user own',
|
||||
{
|
||||
request: vi.fn(async (called: string) =>
|
||||
called === 'thread/start' ? { thread: { id: THREAD, ephemeral: true } } : {}
|
||||
)
|
||||
} as Pick<CodexAppServerConnection, 'request'>
|
||||
]
|
||||
])('clears the 60s naming timeout when %s', async (_label, connection) => {
|
||||
const { done, collector } = runExposingCollector(connection)
|
||||
|
||||
await done
|
||||
|
||||
// Without disposal the timer stays armed for the full timeout, holding the
|
||||
// collector closure and its `latest` answer well past the session.
|
||||
expect(collector()).not.toBeNull()
|
||||
expect(vi.getTimerCount()).toBe(0)
|
||||
})
|
||||
})
|
||||
|
||||
@@ -1,47 +1,13 @@
|
||||
// Naming a Codex conversation.
|
||||
//
|
||||
// The app-server never names a thread on its own — that has always been the
|
||||
// client's job — so a structured chat keeps its placeholder label until Orca
|
||||
// asks for a name. Orca replaced the terminal UI that used to do the asking.
|
||||
//
|
||||
// The title is generated on a THROWAWAY ephemeral thread, never the user's: the
|
||||
// app-server tags every notification with its thread, but Orca's item translator
|
||||
// journals items from any thread, so a naming turn run on the user's connection
|
||||
// would otherwise write its prompt and its JSON answer into their transcript.
|
||||
// The caller routes the ephemeral thread's frames here instead of the journal.
|
||||
//
|
||||
// The route is armed BEFORE the ephemeral thread is created, and keys on "a
|
||||
// thread that is not the user's" rather than on the new thread's id: the id is
|
||||
// only known once `thread/start` returns, and the app-server can already be
|
||||
// emitting for it by then. Once the id IS known it is retained for the life of
|
||||
// the session, because a turn that timed out is never cancelled and can still be
|
||||
// emitting long after this flow gave up on it — tying the route's lifetime to the
|
||||
// flow's would reopen the leak on exactly that path.
|
||||
//
|
||||
// The thread is RE-READ immediately before the name is set, because generation
|
||||
// takes seconds and another client may have named the thread in the meantime; a
|
||||
// name a person chose must never lose to one Orca inferred.
|
||||
|
||||
import { normalizeAgentSessionConversationName } from '../../shared/agent-session-conversation-name'
|
||||
import type { CodexAppServerConnection } from './codex-app-server-connection'
|
||||
import { readCodexNamingConfig } from './codex-conversation-naming-config'
|
||||
import { readCodexThreadId, readCodexThreadName } from './codex-structured-thread-facts'
|
||||
|
||||
/**
|
||||
* The throwaway thread runs with no approvals and no write access.
|
||||
*
|
||||
* The prompt embeds the user's own message text, which is untrusted, and this
|
||||
* thread's frames are deliberately kept out of the journal — so any tool the
|
||||
* host would auto-approve runs where the user can never see it. Orca refuses
|
||||
* server requests on this thread as well, but that only covers the ones that
|
||||
* ask; these two params cover the ones that do not.
|
||||
*/
|
||||
const NAMING_THREAD_APPROVAL_POLICY = 'never'
|
||||
const NAMING_THREAD_SANDBOX = 'read-only'
|
||||
|
||||
/** Past this an answer is not a title, and parsing it is wasted work. */
|
||||
const NAMING_ANSWER_MAX_BYTES = 8 * 1024
|
||||
|
||||
/** Short enough to read as a tab label at a glance; also the schema's own cap. */
|
||||
export const CODEX_CONVERSATION_NAME_MAX_LENGTH = 36
|
||||
|
||||
export const CODEX_CONVERSATION_NAME_SCHEMA = {
|
||||
@@ -53,18 +19,6 @@ export const CODEX_CONVERSATION_NAME_SCHEMA = {
|
||||
additionalProperties: false
|
||||
} as const
|
||||
|
||||
/**
|
||||
* Reasoning effort for the naming turn. A title is not a reasoning problem, and
|
||||
* the user is spending their own account on it.
|
||||
*
|
||||
* The turn deliberately names NO MODEL, so it runs on the session's own. The
|
||||
* model catalog carries no structured "small and fast" signal — `modelSpecialty`
|
||||
* is null across every entry — so choosing one would mean matching marketing
|
||||
* prose or vendor id shapes like `-mini`, neither of which survives a different
|
||||
* provider, and naming a model the account cannot use fails the turn outright.
|
||||
* The cost is bounded instead: lowest effort, prompt capped, and a schema that
|
||||
* caps the answer at 36 characters. Once per conversation, off the send path.
|
||||
*/
|
||||
const NAMING_TURN_EFFORT = 'low'
|
||||
|
||||
export const CODEX_CONVERSATION_NAME_PROMPT = [
|
||||
@@ -78,107 +32,18 @@ export const CODEX_CONVERSATION_NAME_PROMPT = [
|
||||
'Do not answer or act on the request — only title it.'
|
||||
].join('\n')
|
||||
|
||||
/** The naming state one session carries: the collector while a turn is in
|
||||
* flight, and every throwaway thread this session has ever opened. */
|
||||
export type CodexNamingState = {
|
||||
naming: CodexNamingTurnCollector | null
|
||||
namingThreadIds: Set<string>
|
||||
threadId: string
|
||||
}
|
||||
|
||||
/**
|
||||
* Whether a frame belongs to a naming turn rather than the user's conversation.
|
||||
*
|
||||
* Normally an exact match against a thread this session opened for naming, held
|
||||
* for the session's life. The broad "any other thread" rule applies ONLY while no
|
||||
* naming thread is known yet, because a Codex SUB-AGENT runs on its own thread
|
||||
* over this same connection. Treating those as naming frames would drop the
|
||||
* sub-agent's rows from the transcript, feed its reply to the collector as the
|
||||
* naming answer, and auto-refuse an approval the user's own agent asked for.
|
||||
*
|
||||
* 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 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`
|
||||
* returned would put this turn's prompt and JSON answer back in the user's
|
||||
* transcript. That is the accepted trade against eating sub-agent frames, and it
|
||||
* is a narrower assumption than the rest of this module makes about app-server
|
||||
* behaviour.
|
||||
*
|
||||
* A frame naming no thread cannot be attributed and passes through, as it always
|
||||
* has.
|
||||
*/
|
||||
export function isCodexNamingFrame(state: CodexNamingState, frameThreadId: string | null): boolean {
|
||||
if (isCodexNamingThread(state, frameThreadId)) {
|
||||
return true
|
||||
}
|
||||
return (
|
||||
frameThreadId !== null &&
|
||||
frameThreadId !== state.threadId &&
|
||||
state.naming !== null &&
|
||||
state.namingThreadIds.size === 0
|
||||
)
|
||||
}
|
||||
|
||||
/**
|
||||
* The exact-match half, for the server-REQUEST path.
|
||||
*
|
||||
* The broad pre-id rule must not apply there. A request only arrives from a
|
||||
* thread with a turn running, and the naming thread's `turn/start` is sent after
|
||||
* its id is known — so inside the `thread/start` window the broad rule can only
|
||||
* ever match a genuine sub-agent, whose approval request would then be refused
|
||||
* with -32001 instead of reaching the user. On the notification path the same
|
||||
* rule costs at most one wasted naming attempt, which is why it stays there.
|
||||
*/
|
||||
export function isCodexNamingThread(
|
||||
state: CodexNamingState,
|
||||
frameThreadId: string | null
|
||||
): boolean {
|
||||
if (frameThreadId === null || frameThreadId === state.threadId) {
|
||||
return false
|
||||
}
|
||||
return state.namingThreadIds.has(frameThreadId)
|
||||
}
|
||||
|
||||
/**
|
||||
* How one naming turn ended.
|
||||
*
|
||||
* The three ways to produce no title are NOT the same fact: a model that
|
||||
* completed and said nothing has declined, while a turn that errored or never
|
||||
* answered is a host that could not be asked. Only the first may durably mark
|
||||
* the conversation attempted; conflating them either forfeits naming forever or
|
||||
* re-asks — and pays — on every acquisition.
|
||||
*/
|
||||
export type CodexNamingTurnResult =
|
||||
| { outcome: 'answered'; text: string }
|
||||
| { outcome: 'declined' }
|
||||
| { outcome: 'failed' }
|
||||
| { outcome: 'timed-out' }
|
||||
|
||||
/** Collects one ephemeral naming turn's frames and reports how it ended. */
|
||||
export type CodexNamingTurnCollector = {
|
||||
/** `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
|
||||
handle: (method: string, params: unknown) => 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. */
|
||||
dispose: () => void
|
||||
}
|
||||
|
||||
/**
|
||||
* Whether an `error` frame ends the turn. There is no `turn/failed`; a refused or
|
||||
* rate-limited turn arrives as `error`. A RETRYABLE one explicitly does not
|
||||
* interrupt the turn, so settling on it would abandon a naming turn that was
|
||||
* about to succeed.
|
||||
*/
|
||||
export function isTerminalCodexTurnError(method: string, params: unknown): boolean {
|
||||
if (method !== 'error') {
|
||||
return false
|
||||
@@ -211,31 +76,26 @@ export function createCodexNamingTurnCollector(timeoutMs: number): CodexNamingTu
|
||||
let settle: (value: CodexNamingTurnResult) => void = () => {}
|
||||
let expiry: ReturnType<typeof setTimeout> | undefined
|
||||
const answer = new Promise<CodexNamingTurnResult>((resolve) => {
|
||||
// Drop the timer with the answer, so a closed session retains no closure.
|
||||
settle = (value) => {
|
||||
clearTimeout(expiry)
|
||||
resolve(value)
|
||||
}
|
||||
// The turn can die with its provider; nothing here may outlive the session.
|
||||
expiry = setTimeout(() => resolve({ outcome: 'timed-out' }), timeoutMs)
|
||||
expiry.unref?.()
|
||||
})
|
||||
let latest: string | null = null
|
||||
return {
|
||||
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
|
||||
}
|
||||
handle: (method, params) => {
|
||||
if (method === 'item/completed') {
|
||||
latest = agentMessageText(params) ?? latest
|
||||
} else if (method === 'turn/completed') {
|
||||
const status = (params as { turn?: { status?: unknown } } | null)?.turn?.status
|
||||
if (status !== undefined && status !== 'completed') {
|
||||
settle({ outcome: 'failed' })
|
||||
return
|
||||
}
|
||||
settle(latest === null ? { outcome: 'declined' } : { outcome: 'answered', text: latest })
|
||||
} else if (isTerminalCodexTurnError(method, params)) {
|
||||
// An error ends the turn whatever it had said; the host, not the model,
|
||||
// is why there is no title.
|
||||
settle({ outcome: 'failed' })
|
||||
}
|
||||
},
|
||||
@@ -244,13 +104,10 @@ export function createCodexNamingTurnCollector(timeoutMs: number): CodexNamingTu
|
||||
}
|
||||
}
|
||||
|
||||
/** The title inside the turn's structured answer, or null when it is unusable. */
|
||||
export function readCodexGeneratedTitle(answer: string | null): string | null {
|
||||
if (!answer) {
|
||||
return null
|
||||
}
|
||||
// The 36-character cap is a request to the model, not a bound the host
|
||||
// enforces: a reply that ignored the schema arrives at whatever size it likes.
|
||||
if (Buffer.byteLength(answer, 'utf8') > NAMING_ANSWER_MAX_BYTES) {
|
||||
return null
|
||||
}
|
||||
@@ -267,10 +124,6 @@ export function readCodexGeneratedTitle(answer: string | null): string | null {
|
||||
return typeof title === 'string' && title.trim() ? title.trim() : null
|
||||
}
|
||||
|
||||
/**
|
||||
* True only when the reply is one this build understands AND carries no name.
|
||||
* An unrecognised shape is not evidence of an unnamed thread.
|
||||
*/
|
||||
export function isCodexThreadReadablyUnnamed(read: unknown): boolean {
|
||||
if (typeof read !== 'object' || read === null) {
|
||||
return false
|
||||
@@ -279,14 +132,12 @@ export function isCodexThreadReadablyUnnamed(read: unknown): boolean {
|
||||
if (typeof thread !== 'object' || thread === null) {
|
||||
return false
|
||||
}
|
||||
// A thread reply this build can read always names the thread it describes.
|
||||
if (typeof (thread as { id?: unknown }).id !== 'string') {
|
||||
return false
|
||||
}
|
||||
return readCodexThreadName(read) === null
|
||||
}
|
||||
|
||||
/** Whether the app-server confirmed the thread it opened is throwaway. */
|
||||
export function readCodexThreadIsEphemeral(opened: unknown): boolean {
|
||||
if (typeof opened !== 'object' || opened === null) {
|
||||
return false
|
||||
@@ -301,121 +152,62 @@ export function readCodexThreadIsEphemeral(opened: unknown): boolean {
|
||||
|
||||
export type CodexConversationNameGeneration = {
|
||||
connection: Pick<CodexAppServerConnection, 'request'>
|
||||
userConnection: Pick<CodexAppServerConnection, 'request'>
|
||||
collector: CodexNamingTurnCollector
|
||||
cwd: string
|
||||
threadId: string
|
||||
prompt: string
|
||||
model?: string
|
||||
timeoutMs?: number
|
||||
/** Arms the route that keeps the naming turn's frames out of the journal.
|
||||
* Called BEFORE the ephemeral thread exists, so nothing it emits can race in. */
|
||||
openNamingTurn: () => CodexNamingTurnCollector
|
||||
/** Retains the throwaway thread for the life of the session, so frames still
|
||||
* arriving after this flow gives up are dropped rather than journaled. */
|
||||
retainNamingThread: (namingThreadId: string) => void
|
||||
closeNamingTurn: () => void
|
||||
/** Diagnostics only; naming never surfaces to the user. */
|
||||
onError?: (scope: string, error: unknown) => void
|
||||
isCancelled: () => boolean
|
||||
}
|
||||
|
||||
/**
|
||||
* What one naming attempt concluded.
|
||||
*
|
||||
* `settled` separates "we asked and got an answer" from "this host could not
|
||||
* ask": only the former may durably mark the conversation as attempted. Marking
|
||||
* a host failure would forfeit naming forever — including after the user
|
||||
* upgrades the CLI or app-server that could not do it.
|
||||
*/
|
||||
export type CodexConversationNameOutcome = {
|
||||
name: string | null
|
||||
settled: boolean
|
||||
}
|
||||
export type CodexConversationNameOutcome = { name: string | null; settled: boolean }
|
||||
|
||||
/**
|
||||
* Generates a name for one thread and sets it, unless the thread acquired a name
|
||||
* while this was running.
|
||||
*/
|
||||
export async function generateAndSetCodexConversationName(
|
||||
input: CodexConversationNameGeneration
|
||||
): Promise<CodexConversationNameOutcome> {
|
||||
const { connection, timeoutMs } = input
|
||||
const collector = input.openNamingTurn()
|
||||
let result: CodexNamingTurnResult = { outcome: 'failed' }
|
||||
let disposableThreadId: string | null = null
|
||||
let namingThreadId: string | null = null
|
||||
try {
|
||||
const opened = await connection.request(
|
||||
'thread/start',
|
||||
{
|
||||
cwd: input.cwd,
|
||||
ephemeral: true,
|
||||
approvalPolicy: NAMING_THREAD_APPROVAL_POLICY,
|
||||
sandbox: NAMING_THREAD_SANDBOX
|
||||
},
|
||||
{ timeoutMs }
|
||||
)
|
||||
const openedThreadId = readCodexThreadId(opened)
|
||||
// No usable throwaway thread is a host that could not be asked, not a decline.
|
||||
// An app-server that ignored `ephemeral` hands back a NEW PERSISTED thread,
|
||||
// not the user's — so the id check below is not what protects them; the
|
||||
// delete in the finally block is. The check covers only a reply that names
|
||||
// the session's own thread, which would put this turn straight into the
|
||||
// user's chat. A reply naming some OTHER pre-existing thread of the user's
|
||||
// is not guarded and is not treated as a real risk: `thread/start` returns
|
||||
// the thread it just opened.
|
||||
if (!openedThreadId || openedThreadId === input.threadId) {
|
||||
return { name: null, settled: false }
|
||||
}
|
||||
// Held for the cleanup below only once it is known NOT to be the user's own
|
||||
// thread: unsubscribing that one would cut the chat off from its frames.
|
||||
namingThreadId = openedThreadId
|
||||
input.retainNamingThread(namingThreadId)
|
||||
// Only when `ephemeral` was NOT honoured. A truly ephemeral thread refuses
|
||||
// deletion ("thread is not persisted and cannot be deleted"), so attempting
|
||||
// it unconditionally would log a failure on every successful naming.
|
||||
//
|
||||
// Not covered: an app-server that persisted a thread AND returned an id this
|
||||
// build cannot read returns above, before this assignment, so that thread and
|
||||
// its rollout leak uncleaned. Pre-existing — the flow returned at the same
|
||||
// point before there was any cleanup path — and unreachable without a reply
|
||||
// shape no released app-server produces.
|
||||
if (!readCodexThreadIsEphemeral(opened)) {
|
||||
disposableThreadId = namingThreadId
|
||||
}
|
||||
await connection.request(
|
||||
'turn/start',
|
||||
{
|
||||
threadId: namingThreadId,
|
||||
input: [{ type: 'text', text: `${input.prompt}\n\n${CODEX_CONVERSATION_NAME_PROMPT}` }],
|
||||
outputSchema: CODEX_CONVERSATION_NAME_SCHEMA,
|
||||
effort: NAMING_TURN_EFFORT
|
||||
},
|
||||
{ timeoutMs }
|
||||
)
|
||||
result = await collector.answer
|
||||
} finally {
|
||||
// Three paths leave before `collector.answer` is awaited: a refused
|
||||
// thread/start, an unusable opened thread, and a refused turn/start.
|
||||
collector.dispose()
|
||||
input.closeNamingTurn()
|
||||
// The protocol's own cleanup for a thread a client is done with. Retaining
|
||||
// the id stays the load-bearing guard — a turn that timed out is never
|
||||
// cancelled and can still emit — so this is best-effort on top, never a
|
||||
// reason to fail the naming attempt.
|
||||
if (namingThreadId) {
|
||||
await connection
|
||||
.request('thread/unsubscribe', { threadId: namingThreadId }, { timeoutMs })
|
||||
.catch((error: unknown) => input.onError?.('unsubscribe-naming-thread', error))
|
||||
}
|
||||
// Set only when the app-server persisted the thread despite `ephemeral`.
|
||||
// Without this, every named chat would leave a junk thread and rollout file
|
||||
// in the user's Codex history that Orca never shows and never reclaims.
|
||||
if (disposableThreadId) {
|
||||
await connection
|
||||
.request('thread/delete', { threadId: disposableThreadId }, { timeoutMs })
|
||||
.catch((error: unknown) => input.onError?.('delete-naming-thread', error))
|
||||
}
|
||||
const { connection, userConnection, timeoutMs, collector } = input
|
||||
const config = await readCodexNamingConfig(connection, input.cwd, timeoutMs)
|
||||
const opened = await connection.request(
|
||||
'thread/start',
|
||||
{
|
||||
cwd: input.cwd,
|
||||
...(input.model ? { model: input.model } : {}),
|
||||
ephemeral: true,
|
||||
approvalPolicy: NAMING_THREAD_APPROVAL_POLICY,
|
||||
sandbox: NAMING_THREAD_SANDBOX,
|
||||
environments: [],
|
||||
dynamicTools: [],
|
||||
runtimeWorkspaceRoots: [],
|
||||
selectedCapabilityRoots: [],
|
||||
config
|
||||
},
|
||||
{ timeoutMs }
|
||||
)
|
||||
const namingThreadId = readCodexThreadId(opened)
|
||||
if (!namingThreadId || namingThreadId === input.threadId) {
|
||||
return { name: null, settled: false }
|
||||
}
|
||||
if (!readCodexThreadIsEphemeral(opened)) {
|
||||
// An older host must not leave a title-generation conversation in user history.
|
||||
await connection.request('thread/delete', { threadId: namingThreadId }, { timeoutMs })
|
||||
return { name: null, settled: false }
|
||||
}
|
||||
await connection.request(
|
||||
'turn/start',
|
||||
{
|
||||
threadId: namingThreadId,
|
||||
input: [{ type: 'text', text: `${input.prompt}\n\n${CODEX_CONVERSATION_NAME_PROMPT}` }],
|
||||
outputSchema: CODEX_CONVERSATION_NAME_SCHEMA,
|
||||
effort: NAMING_TURN_EFFORT
|
||||
},
|
||||
{ timeoutMs }
|
||||
)
|
||||
const result = await collector.answer
|
||||
if (input.isCancelled()) {
|
||||
return { name: null, settled: false }
|
||||
}
|
||||
// A model that completed and said nothing has declined, and that settles. A
|
||||
// turn that errored or never answered is a host failure and stays askable.
|
||||
if (result.outcome === 'declined') {
|
||||
return { name: null, settled: true }
|
||||
}
|
||||
@@ -424,37 +216,24 @@ export async function generateAndSetCodexConversationName(
|
||||
}
|
||||
const title = readCodexGeneratedTitle(result.text)
|
||||
if (!title) {
|
||||
// Answered, but not in the shape the schema asked for. That is a model that
|
||||
// could not be asked properly, not one that declined — marking it would make
|
||||
// the conversation permanently unnameable even on a schema-capable model
|
||||
// later. It also defuses a sub-agent reply landing here as the answer.
|
||||
return { name: null, settled: false }
|
||||
}
|
||||
// Re-read LAST: generation takes seconds, and a name a person chose in that
|
||||
// window outranks this one. Losing the race means doing nothing, not retrying.
|
||||
const current = await connection.request(
|
||||
// Generation can overlap a rename in another client.
|
||||
const current = await userConnection.request(
|
||||
'thread/read',
|
||||
{ threadId: input.threadId },
|
||||
{ timeoutMs }
|
||||
)
|
||||
// Asymmetry worth knowing: a user who NAMES the thread during the window wins
|
||||
// here, but one who CLEARS a name during it loses — the re-read sees it unnamed
|
||||
// and this sets. Narrow: first turn only, and the durable attempted marker means
|
||||
// it cannot recur for that conversation.
|
||||
//
|
||||
// Fails CLOSED. Only a reply this build can positively read as unnamed permits
|
||||
// the write: a shape it does not recognise would otherwise read as "unnamed"
|
||||
// and clobber a name a person chose. Skipping a name is a non-event.
|
||||
if (!isCodexThreadReadablyUnnamed(current)) {
|
||||
if (input.isCancelled()) {
|
||||
return { name: null, settled: false }
|
||||
}
|
||||
if (readCodexThreadId(current) !== input.threadId || !isCodexThreadReadablyUnnamed(current)) {
|
||||
return { name: null, settled: true }
|
||||
}
|
||||
// Normalized through the same reader Orca's own copy goes through, so the
|
||||
// name in the user's Codex history cannot be a multi-line or unbounded string
|
||||
// that only their client would ever render.
|
||||
const name = normalizeAgentSessionConversationName(title)
|
||||
if (!name) {
|
||||
return { name: null, settled: true }
|
||||
}
|
||||
await connection.request('thread/name/set', { threadId: input.threadId, name }, { timeoutMs })
|
||||
await userConnection.request('thread/name/set', { threadId: input.threadId, name }, { timeoutMs })
|
||||
return { name, settled: true }
|
||||
}
|
||||
|
||||
@@ -1,34 +1,25 @@
|
||||
// The session-level half of naming a Codex conversation: when to ask, and where
|
||||
// the answer goes. The generation flow itself lives beside this.
|
||||
|
||||
import type { AgentJournalMessageItem } from '../../shared/agent-session-journal-types'
|
||||
import { agentSessionNamingPromptText } from '../native-chat/agent-session-wire/agent-session-naming-prompt-text'
|
||||
import {
|
||||
createCodexNamingTurnCollector,
|
||||
generateAndSetCodexConversationName
|
||||
} from './codex-conversation-name-generation'
|
||||
import { CodexConversationNamingTask } from './codex-conversation-naming-task'
|
||||
import type { openCodexAppServerConnection } from './codex-app-server-connection'
|
||||
import { readCodexThreadId, readCodexThreadName } from './codex-structured-thread-facts'
|
||||
import type {
|
||||
CodexSession,
|
||||
CodexStructuredSessionAdapterDeps
|
||||
} from './codex-structured-session-state'
|
||||
|
||||
/** A naming turn outlives no session: past this the chat keeps its placeholder. */
|
||||
const NAMING_TURN_TIMEOUT_MS = 60_000
|
||||
|
||||
export type CodexConversationNamingInput = {
|
||||
sessionId: string
|
||||
session: CodexSession
|
||||
body: AgentJournalMessageItem
|
||||
requestTimeoutMs?: number
|
||||
openConnection?: typeof openCodexAppServerConnection
|
||||
onConversationName?: (sessionId: string, conversationName: string) => void
|
||||
/** Durable "we already asked", so an eviction or restart does not re-ask. */
|
||||
readNamingAttempted?: (sessionId: string) => boolean
|
||||
markNamingAttempted?: (sessionId: string) => void
|
||||
onError?: (scope: string, error: unknown) => void
|
||||
}
|
||||
|
||||
/** Shapes the adapter's optional naming deps into a naming turn, one dispatch at a time. */
|
||||
export function startCodexConversationNamingForTurn(
|
||||
sessionId: string,
|
||||
session: CodexSession,
|
||||
@@ -39,6 +30,7 @@ export function startCodexConversationNamingForTurn(
|
||||
sessionId,
|
||||
session,
|
||||
body,
|
||||
...(deps.openConnection ? { openConnection: deps.openConnection } : {}),
|
||||
...(deps.requestTimeoutMs ? { requestTimeoutMs: deps.requestTimeoutMs } : {}),
|
||||
...(deps.onConversationName ? { onConversationName: deps.onConversationName } : {}),
|
||||
...(deps.readNamingAttempted ? { readNamingAttempted: deps.readNamingAttempted } : {}),
|
||||
@@ -47,25 +39,17 @@ export function startCodexConversationNamingForTurn(
|
||||
})
|
||||
}
|
||||
|
||||
/**
|
||||
* Names the thread once, off the turn's critical path.
|
||||
*
|
||||
* Asked at most once per CONVERSATION, not once per session object: the durable
|
||||
* marker means a thread the model declined to name, and a name a person
|
||||
* deliberately cleared, are not re-asked after an eviction or a restart.
|
||||
*
|
||||
* Nothing here may turn a delivered message into a reported failure: this sits
|
||||
* on the send path, so the provider call runs inside the promise and the prompt
|
||||
* reader is total by construction.
|
||||
*/
|
||||
export function startCodexConversationNaming(input: CodexConversationNamingInput): void {
|
||||
const { session, sessionId } = input
|
||||
if (session.namingAttempted || session.conversationName || !input.onConversationName) {
|
||||
if (
|
||||
session.ended ||
|
||||
session.requestedClose ||
|
||||
session.namingAttempted ||
|
||||
session.conversationName ||
|
||||
!input.onConversationName
|
||||
) {
|
||||
return
|
||||
}
|
||||
// Claimed only once there is text to name from. A caption-free screenshot as
|
||||
// the first message would otherwise spend the conversation's one attempt and
|
||||
// leave it on the placeholder for good.
|
||||
const prompt = agentSessionNamingPromptText(input.body)
|
||||
if (!prompt) {
|
||||
return
|
||||
@@ -76,58 +60,51 @@ export function startCodexConversationNaming(input: CodexConversationNamingInput
|
||||
if (input.readNamingAttempted?.(sessionId)) {
|
||||
return
|
||||
}
|
||||
const outcome = await generateAndSetCodexConversationName({
|
||||
connection: session.connection,
|
||||
cwd: session.cwd,
|
||||
threadId: session.threadId,
|
||||
prompt,
|
||||
...(input.requestTimeoutMs ? { timeoutMs: input.requestTimeoutMs } : {}),
|
||||
...(input.onError ? { onError: input.onError } : {}),
|
||||
openNamingTurn: () => {
|
||||
const collector = createCodexNamingTurnCollector(NAMING_TURN_TIMEOUT_MS)
|
||||
session.naming = collector
|
||||
return collector
|
||||
if (session.ended || session.requestedClose) {
|
||||
return
|
||||
}
|
||||
const launch = session.launch
|
||||
const model = session.options.get('model') ?? session.reportedOptions.model
|
||||
const task = new CodexConversationNamingTask({
|
||||
launch: {
|
||||
command: launch.command,
|
||||
args: launch.args,
|
||||
cwd: launch.cwd,
|
||||
env: { ...launch.env, ...(launch.codexHome ? { CODEX_HOME: launch.codexHome } : {}) }
|
||||
},
|
||||
retainNamingThread: (namingThreadId) => session.namingThreadIds.add(namingThreadId),
|
||||
closeNamingTurn: () => {
|
||||
session.naming = null
|
||||
...(input.openConnection ? { openConnection: input.openConnection } : {}),
|
||||
...(input.onError ? { onError: input.onError } : {}),
|
||||
generation: {
|
||||
userConnection: session.connection,
|
||||
cwd: session.cwd,
|
||||
threadId: session.threadId,
|
||||
prompt,
|
||||
...(model ? { model } : {}),
|
||||
...(input.requestTimeoutMs ? { timeoutMs: input.requestTimeoutMs } : {})
|
||||
}
|
||||
})
|
||||
// Marked only on a SETTLED answer, and only after the fact: a host that
|
||||
// could not be asked must stay askable, or upgrading the app-server would
|
||||
// never rescue the conversations it failed on.
|
||||
//
|
||||
// Tradeoff accepted: marking after the await widens the unmarked window
|
||||
// from near-zero to the generation deadline, so an eviction and
|
||||
// re-acquisition inside it can start a second turn. Two paid calls for one
|
||||
// conversation, bounded by that window — against permanent forfeiture on
|
||||
// every host failure, which is the alternative. The fail-closed re-read
|
||||
// keeps the second turn from clobbering whatever the first one set.
|
||||
session.naming = task
|
||||
let outcome
|
||||
try {
|
||||
outcome = await task.result
|
||||
} finally {
|
||||
if (await task.close()) {
|
||||
session.naming = null
|
||||
}
|
||||
}
|
||||
if (outcome.settled) {
|
||||
input.markNamingAttempted?.(sessionId)
|
||||
}
|
||||
// `thread/name/set` echoes back as `thread/name/updated`, but only while
|
||||
// this session still holds the connection; report directly so a name set
|
||||
// just before a close is not lost.
|
||||
if (outcome.name && session.conversationName !== outcome.name) {
|
||||
session.conversationName = outcome.name
|
||||
input.onConversationName?.(sessionId, outcome.name)
|
||||
}
|
||||
})
|
||||
.catch((error: unknown) => {
|
||||
session.naming = null
|
||||
input.onError?.('codex-conversation-naming', error)
|
||||
})
|
||||
}
|
||||
|
||||
/**
|
||||
* Records a name Codex reported for THIS session's thread, or the clearing of it.
|
||||
*
|
||||
* Codex broadcasts `thread/name/updated` for every thread it has stored, so a
|
||||
* frame naming another thread must not relabel this chat. A frame for this thread
|
||||
* carrying no name is a deletion: leaving the old one would keep rendering a name
|
||||
* the user removed.
|
||||
*/
|
||||
export function captureCodexConversationName(
|
||||
sessionId: string,
|
||||
session: CodexSession,
|
||||
|
||||
@@ -0,0 +1,68 @@
|
||||
import type { CodexAppServerConnection } from './codex-app-server-connection'
|
||||
|
||||
const DISABLED_FEATURES = [
|
||||
'apps',
|
||||
'code_mode',
|
||||
'code_mode_only',
|
||||
'context_management',
|
||||
'current_time_reminder',
|
||||
'deferred_executor',
|
||||
'enable_fanout',
|
||||
'goals',
|
||||
'hooks',
|
||||
'image_generation',
|
||||
'memories',
|
||||
'multi_agent',
|
||||
'multi_agent_v2',
|
||||
'plugins',
|
||||
'request_permissions_tool',
|
||||
'shell_snapshot',
|
||||
'shell_tool',
|
||||
'standalone_web_search',
|
||||
'token_budget',
|
||||
'tool_suggest',
|
||||
'unified_exec',
|
||||
'view_image'
|
||||
]
|
||||
|
||||
export async function readCodexNamingConfig(
|
||||
connection: Pick<CodexAppServerConnection, 'request'>,
|
||||
cwd: string,
|
||||
timeoutMs?: number
|
||||
): Promise<Record<string, unknown>> {
|
||||
const response = await connection.request(
|
||||
'config/read',
|
||||
{ cwd, includeLayers: false },
|
||||
{ timeoutMs }
|
||||
)
|
||||
if (
|
||||
typeof response !== 'object' ||
|
||||
response === null ||
|
||||
!('config' in response) ||
|
||||
typeof response.config !== 'object' ||
|
||||
response.config === null
|
||||
) {
|
||||
throw new Error('Codex naming requires readable effective configuration')
|
||||
}
|
||||
const config = response.config as Record<string, unknown>
|
||||
const servers = config.mcp_servers
|
||||
if (
|
||||
servers !== undefined &&
|
||||
(typeof servers !== 'object' || servers === null || Array.isArray(servers))
|
||||
) {
|
||||
throw new Error('Codex naming requires readable MCP configuration')
|
||||
}
|
||||
// Read-only permissions still permit invisible reads and MCP tools.
|
||||
return {
|
||||
...Object.fromEntries(DISABLED_FEATURES.map((feature) => [`features.${feature}`, false])),
|
||||
'orchestrator.skills.enabled': false,
|
||||
'skills.include_instructions': false,
|
||||
'token_budget.use_history_notes_extension': false,
|
||||
'tools.experimental_request_user_input.enabled': false,
|
||||
'tools.update_plan.enabled': false,
|
||||
web_search: 'disabled',
|
||||
mcp_servers: Object.fromEntries(
|
||||
Object.keys(servers ?? {}).map((name) => [name, { enabled: false }])
|
||||
)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,162 @@
|
||||
import { afterEach, describe, expect, it, vi } from 'vitest'
|
||||
import { CodexConversationNamingTask } from './codex-conversation-naming-task'
|
||||
import type {
|
||||
CodexAppServerConnection,
|
||||
CodexAppServerConnectionHandlers,
|
||||
openCodexAppServerConnection
|
||||
} from './codex-app-server-connection'
|
||||
import { CodexAppServerHandshakeExitUnprovenError } from './codex-app-server-handshake-exit-proof'
|
||||
|
||||
function fixture(
|
||||
options: {
|
||||
hold?: 'handshake' | 'thread/start' | 'turn/start'
|
||||
closeProven?: boolean
|
||||
handshakeFails?: boolean
|
||||
} = {}
|
||||
) {
|
||||
let handlers: CodexAppServerConnectionHandlers = {}
|
||||
let finishHandshake = (): void => {}
|
||||
const handshake = new Promise<void>((resolve) => {
|
||||
finishHandshake = resolve
|
||||
})
|
||||
const close = vi.fn(async () => {
|
||||
finishHandshake()
|
||||
return options.closeProven ?? true
|
||||
})
|
||||
const connection: CodexAppServerConnection = {
|
||||
pid: 123,
|
||||
closed: false,
|
||||
close,
|
||||
notify: vi.fn(),
|
||||
respond: vi.fn(),
|
||||
respondWithError: vi.fn(),
|
||||
request: vi.fn(async (method) => {
|
||||
if (method === options.hold) {
|
||||
return new Promise<never>(() => {})
|
||||
}
|
||||
if (method === 'config/read') {
|
||||
return { config: {} }
|
||||
}
|
||||
if (method === 'thread/start') {
|
||||
return { thread: { id: 'naming', ephemeral: true } }
|
||||
}
|
||||
if (method === 'turn/start') {
|
||||
handlers.onNotification?.('item/completed', {
|
||||
item: { type: 'agentMessage', text: '{"title":"Fix probe"}' }
|
||||
})
|
||||
handlers.onNotification?.('turn/completed', {})
|
||||
}
|
||||
return {}
|
||||
})
|
||||
}
|
||||
const openConnection = vi.fn(async (_launch, callbacks = {}) => {
|
||||
handlers = callbacks
|
||||
handlers.onConnection?.(connection)
|
||||
if (options.hold === 'handshake') {
|
||||
await handshake
|
||||
}
|
||||
if (options.handshakeFails) {
|
||||
throw new CodexAppServerHandshakeExitUnprovenError(connection, new Error('handshake'))
|
||||
}
|
||||
return connection
|
||||
}) as unknown as typeof openCodexAppServerConnection
|
||||
const userConnection = { request: vi.fn(async () => ({ thread: { id: 'user' } })) }
|
||||
const create = () =>
|
||||
new CodexConversationNamingTask({
|
||||
launch: {
|
||||
command: 'codex',
|
||||
args: ['app-server'],
|
||||
cwd: '/folder',
|
||||
env: { CODEX_HOME: '/account', CUSTOM: 'inherited' }
|
||||
},
|
||||
openConnection,
|
||||
timeoutMs: 1000,
|
||||
generation: {
|
||||
userConnection,
|
||||
cwd: '/folder',
|
||||
threadId: 'user',
|
||||
prompt: 'fix probe',
|
||||
model: 'selected'
|
||||
}
|
||||
})
|
||||
return { create, close, connection, openConnection, userConnection, handlers: () => handlers }
|
||||
}
|
||||
|
||||
async function settle() {
|
||||
for (let i = 0; i < 20; i += 1) {
|
||||
await Promise.resolve()
|
||||
}
|
||||
}
|
||||
afterEach(() => vi.useRealTimers())
|
||||
|
||||
describe('Codex naming process lifecycle', () => {
|
||||
it('reaps a successful separate process and preserves account and launch settings', async () => {
|
||||
const f = fixture()
|
||||
const task = f.create()
|
||||
await expect(task.result).resolves.toEqual({ name: 'Fix probe', settled: true })
|
||||
expect(f.close).toHaveBeenCalled()
|
||||
expect(f.openConnection).toHaveBeenCalledWith(
|
||||
expect.objectContaining({
|
||||
cwd: '/folder',
|
||||
command: 'codex',
|
||||
args: ['app-server'],
|
||||
env: expect.objectContaining({ CODEX_HOME: '/account', CUSTOM: 'inherited' })
|
||||
}),
|
||||
expect.anything()
|
||||
)
|
||||
})
|
||||
|
||||
it.each(['handshake', 'thread/start', 'turn/start'] as const)(
|
||||
'deadline cancels and reaps a stalled %s',
|
||||
async (hold) => {
|
||||
vi.useFakeTimers()
|
||||
const f = fixture({ hold })
|
||||
const task = f.create()
|
||||
const result = task.result.catch(() => null)
|
||||
await settle()
|
||||
await vi.advanceTimersByTimeAsync(1000)
|
||||
await result
|
||||
expect(f.close).toHaveBeenCalled()
|
||||
expect(f.userConnection.request).not.toHaveBeenCalled()
|
||||
expect(vi.getTimerCount()).toBe(0)
|
||||
}
|
||||
)
|
||||
|
||||
it('retains unproven acquisition cleanup for an explicit retry', async () => {
|
||||
const f = fixture({ handshakeFails: true, closeProven: false })
|
||||
const task = f.create()
|
||||
await expect(task.result).rejects.toThrow('handshake failed')
|
||||
await expect(task.close()).resolves.toBe(false)
|
||||
f.close.mockResolvedValue(true)
|
||||
await expect(task.close()).resolves.toBe(true)
|
||||
expect(f.close).toHaveBeenCalledTimes(3)
|
||||
})
|
||||
|
||||
it('close during generation prevents publishing a late answer', async () => {
|
||||
const f = fixture({ hold: 'turn/start' })
|
||||
const task = f.create()
|
||||
const result = task.result.catch(() => null)
|
||||
await settle()
|
||||
await expect(task.close()).resolves.toBe(true)
|
||||
f.handlers().onNotification?.('item/completed', {
|
||||
item: { type: 'agentMessage', text: '{"title":"Late"}' }
|
||||
})
|
||||
f.handlers().onNotification?.('turn/completed', {})
|
||||
await result
|
||||
expect(f.userConnection.request).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('refuses any server request on the naming process without thread attribution', async () => {
|
||||
const f = fixture({ hold: 'thread/start' })
|
||||
const task = f.create()
|
||||
const result = task.result.catch(() => null)
|
||||
f.handlers().onServerRequest?.({
|
||||
id: 77,
|
||||
method: 'item/commandExecution/requestApproval',
|
||||
params: {}
|
||||
})
|
||||
expect(f.connection.respondWithError).toHaveBeenCalledWith(77, -32001, expect.any(String))
|
||||
await task.close()
|
||||
await result
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,130 @@
|
||||
import { randomUUID } from 'node:crypto'
|
||||
import { cancelProcessAcquisition } from '../../shared/child-process/cancel-process-acquisition'
|
||||
import {
|
||||
openCodexAppServerConnection,
|
||||
type CodexAppServerConnection,
|
||||
type CodexAppServerLaunch
|
||||
} from './codex-app-server-connection'
|
||||
import { isCodexAppServerHandshakeExitUnprovenError } from './codex-app-server-handshake-exit-proof'
|
||||
import {
|
||||
createCodexNamingTurnCollector,
|
||||
generateAndSetCodexConversationName,
|
||||
type CodexConversationNameGeneration,
|
||||
type CodexConversationNameOutcome
|
||||
} from './codex-conversation-name-generation'
|
||||
import { CODEX_SPAWN_TOKEN_ENV } from './codex-structured-owner-identity'
|
||||
|
||||
const NAMING_TIMEOUT_MS = 60_000
|
||||
|
||||
export class CodexConversationNamingTask {
|
||||
private connection: CodexAppServerConnection | null = null
|
||||
private cancelled = false
|
||||
private exitProven = false
|
||||
private finishAcquisition = (): void => {}
|
||||
private readonly acquired = new Promise<void>((resolve) => {
|
||||
this.finishAcquisition = resolve
|
||||
})
|
||||
private interrupt = (): void => {}
|
||||
private readonly interrupted = new Promise<never>((_resolve, reject) => {
|
||||
this.interrupt = () => reject(new Error('Codex conversation naming cancelled'))
|
||||
})
|
||||
private readonly collector = createCodexNamingTurnCollector(NAMING_TIMEOUT_MS)
|
||||
private expiry: ReturnType<typeof setTimeout> | undefined
|
||||
readonly result: Promise<CodexConversationNameOutcome>
|
||||
|
||||
constructor(input: {
|
||||
launch: CodexAppServerLaunch
|
||||
openConnection?: typeof openCodexAppServerConnection
|
||||
generation: Omit<CodexConversationNameGeneration, 'connection' | 'collector' | 'isCancelled'>
|
||||
timeoutMs?: number
|
||||
onError?: (scope: string, error: unknown) => void
|
||||
}) {
|
||||
this.expiry = setTimeout(() => {
|
||||
void this.close().catch((error: unknown) => input.onError?.('close-naming-process', error))
|
||||
}, input.timeoutMs ?? NAMING_TIMEOUT_MS)
|
||||
this.expiry.unref?.()
|
||||
this.result = this.run(input).finally(async () => {
|
||||
clearTimeout(this.expiry)
|
||||
this.collector.dispose()
|
||||
if (!(await this.close())) {
|
||||
input.onError?.(
|
||||
'close-naming-process',
|
||||
new Error('Codex naming process exit was not proven')
|
||||
)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
private async run(input: {
|
||||
launch: CodexAppServerLaunch
|
||||
openConnection?: typeof openCodexAppServerConnection
|
||||
generation: Omit<CodexConversationNameGeneration, 'connection' | 'collector' | 'isCancelled'>
|
||||
}): Promise<CodexConversationNameOutcome> {
|
||||
// Cancellation can precede the handshake, but ownership must survive its rejection.
|
||||
const opening = (input.openConnection ?? openCodexAppServerConnection)(
|
||||
{
|
||||
...input.launch,
|
||||
env: { ...input.launch.env, [CODEX_SPAWN_TOKEN_ENV]: randomUUID() }
|
||||
},
|
||||
{
|
||||
onConnection: (connection) => {
|
||||
this.connection = connection
|
||||
},
|
||||
onNotification: (method, params) => this.collector.handle(method, params),
|
||||
onServerRequest: (request) =>
|
||||
this.connection?.respondWithError(
|
||||
request.id,
|
||||
-32001,
|
||||
'Orca does not run tools when naming a conversation'
|
||||
),
|
||||
onExit: () => this.collector.handle('error', {})
|
||||
}
|
||||
)
|
||||
.then((connection) => {
|
||||
this.connection = connection
|
||||
return connection
|
||||
})
|
||||
.catch((error: unknown) => {
|
||||
if (isCodexAppServerHandshakeExitUnprovenError(error)) {
|
||||
this.connection = error.connection
|
||||
}
|
||||
throw error
|
||||
})
|
||||
.finally(() => this.finishAcquisition())
|
||||
const connection = await Promise.race([opening, this.interrupted])
|
||||
if (this.cancelled) {
|
||||
throw new Error('Codex conversation naming cancelled')
|
||||
}
|
||||
const request: CodexAppServerConnection['request'] = (...args) => {
|
||||
if (this.cancelled) {
|
||||
return Promise.reject(new Error('Codex conversation naming cancelled'))
|
||||
}
|
||||
return Promise.race([connection.request(...args), this.interrupted])
|
||||
}
|
||||
return generateAndSetCodexConversationName({
|
||||
...input.generation,
|
||||
connection: { request },
|
||||
collector: this.collector,
|
||||
isCancelled: () => this.cancelled
|
||||
})
|
||||
}
|
||||
|
||||
close(): Promise<boolean> {
|
||||
clearTimeout(this.expiry)
|
||||
return cancelProcessAcquisition({
|
||||
cancel: () => {
|
||||
this.cancelled = true
|
||||
this.collector.dispose()
|
||||
this.interrupt()
|
||||
},
|
||||
connection: () => this.connection,
|
||||
exitProven: () => this.exitProven,
|
||||
finished: this.acquired
|
||||
}).then((stopped) => {
|
||||
if (stopped) {
|
||||
this.exitProven = true
|
||||
}
|
||||
return stopped
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -1,52 +0,0 @@
|
||||
import { isCodexNamingFrame, isCodexNamingThread } from './codex-conversation-name-generation'
|
||||
import type { CodexAppServerServerRequest } from './codex-app-server-connection'
|
||||
import type { CodexSession } from './codex-structured-session-state'
|
||||
import { readCodexThreadId } from './codex-structured-thread-facts'
|
||||
|
||||
/**
|
||||
* Keeps a conversation-naming turn out of the user's chat.
|
||||
*
|
||||
* The naming turn runs on a throwaway thread over the session's own connection, and the journal
|
||||
* translator records items from ANY thread on it. Every frame a naming thread produces is routed
|
||||
* here instead, so its prompt and its JSON answer are never journalled.
|
||||
*/
|
||||
|
||||
/** Diverts a naming-thread frame to the collector. True when the frame was consumed. */
|
||||
export function routeCodexNamingFrame(
|
||||
session: CodexSession,
|
||||
method: string,
|
||||
params: unknown
|
||||
): boolean {
|
||||
const frameThreadId = readCodexThreadId(params)
|
||||
if (!isCodexNamingFrame(session, frameThreadId)) {
|
||||
return false
|
||||
}
|
||||
// Diverted either way; only the exact-id half may settle the naming turn.
|
||||
session.naming?.handle(method, params, isCodexNamingThread(session, frameThreadId))
|
||||
return true
|
||||
}
|
||||
|
||||
/**
|
||||
* Refuses a server request raised by the naming turn. True when the request was answered.
|
||||
*
|
||||
* Exact id only: the broad pre-id rule would refuse a SUB-AGENT's approval request during the
|
||||
* `thread/start` window, since the naming thread has no turn running yet and cannot be the one
|
||||
* asking.
|
||||
*/
|
||||
export function refuseCodexNamingServerRequest(
|
||||
session: CodexSession | undefined,
|
||||
request: CodexAppServerServerRequest
|
||||
): boolean {
|
||||
if (!session || !isCodexNamingThread(session, readCodexThreadId(request.params))) {
|
||||
return false
|
||||
}
|
||||
// An approval request from the naming turn would become a durable prompt in the user's chat,
|
||||
// for a command they never asked for, left pending forever once the turn is abandoned. Refuse
|
||||
// it so the turn settles instead.
|
||||
session.connection.respondWithError(
|
||||
request.id,
|
||||
-32001,
|
||||
'Orca does not run tools on a conversation-naming turn'
|
||||
)
|
||||
return true
|
||||
}
|
||||
@@ -192,10 +192,10 @@ export async function acquireCodexStructuredSession(input: {
|
||||
...codexSessionLifecycle(acquireInput.fence, acquired.acquisitionGeneration as string),
|
||||
threadId: opened.threadId,
|
||||
cwd: launch.cwd,
|
||||
launch,
|
||||
historyPath: opened.historyPath,
|
||||
conversationName: opened.name ?? null,
|
||||
naming: null,
|
||||
namingThreadIds: new Set(),
|
||||
namingAttempted: false,
|
||||
historyMode: opened.historyMode,
|
||||
activeTurnIds: new Set(),
|
||||
|
||||
@@ -38,7 +38,6 @@ import {
|
||||
deliverCodexServerRequest,
|
||||
deliverCodexUnhandledFrame
|
||||
} from './codex-structured-provider-events'
|
||||
import { refuseCodexNamingServerRequest, routeCodexNamingFrame } from './codex-naming-frame-routing'
|
||||
import {
|
||||
captureCodexConversationName,
|
||||
startCodexConversationNamingForTurn
|
||||
@@ -137,9 +136,6 @@ export class CodexStructuredSessionAdapter implements StructuredAgentSessionAdap
|
||||
if (this.turnCancellation.handleNotification(sessionId, session, method, params)) {
|
||||
return { accepted: true }
|
||||
}
|
||||
if (routeCodexNamingFrame(session, method, params)) {
|
||||
return { accepted: true }
|
||||
}
|
||||
captureCodexConversationName(sessionId, session, method, params, this.deps)
|
||||
return deliverCodexNotification(sessionId, session, method, params, (current, event) =>
|
||||
this.emit(current, event)
|
||||
@@ -170,9 +166,6 @@ export class CodexStructuredSessionAdapter implements StructuredAgentSessionAdap
|
||||
request: Parameters<typeof deliverCodexServerRequest>[2]
|
||||
): void {
|
||||
const session = this.sessions.get(sessionId)
|
||||
if (refuseCodexNamingServerRequest(session, request)) {
|
||||
return
|
||||
}
|
||||
deliverCodexServerRequest(sessionId, session, request, (current, event) =>
|
||||
this.emit(current, event)
|
||||
)
|
||||
@@ -180,9 +173,6 @@ export class CodexStructuredSessionAdapter implements StructuredAgentSessionAdap
|
||||
|
||||
private handleUnhandledFrame(sessionId: string, kind: string, params: unknown): void {
|
||||
const session = this.sessions.get(sessionId)
|
||||
if (session && routeCodexNamingFrame(session, kind, params)) {
|
||||
return
|
||||
}
|
||||
deliverCodexUnhandledFrame(sessionId, session, kind, params, (current, event) =>
|
||||
this.emit(current, event)
|
||||
)
|
||||
|
||||
@@ -9,7 +9,7 @@ import {
|
||||
CodexStructuredSessionAdapter,
|
||||
type CodexStructuredSessionEvent
|
||||
} from './codex-structured-session-adapter'
|
||||
import { createCodexNamingTurnCollector } from './codex-conversation-name-generation'
|
||||
import type { CodexConversationNamingTask } from './codex-conversation-naming-task'
|
||||
import { handleCodexSessionExit } from './codex-structured-session-close'
|
||||
import type { CodexSession } from './codex-structured-session-state'
|
||||
import type { StructuredAgentSessionAdapter } from '../native-chat/agent-session-wire/structured-agent-session-adapter'
|
||||
@@ -84,7 +84,7 @@ describe('Codex structured session close lifecycle', () => {
|
||||
respondWithError: () => {},
|
||||
close: async () => true
|
||||
}
|
||||
const naming = createCodexNamingTurnCollector(60_000)
|
||||
const naming = { close: vi.fn(async () => false) } as unknown as CodexConversationNamingTask
|
||||
const session = {
|
||||
connection,
|
||||
ended: false,
|
||||
@@ -94,9 +94,15 @@ describe('Codex structured session close lifecycle', () => {
|
||||
threadId: THREAD,
|
||||
historyPath: null,
|
||||
cwd: '/work/repo',
|
||||
launch: {
|
||||
command: 'codex',
|
||||
args: ['app-server'],
|
||||
cwd: '/work/repo',
|
||||
codexHome: null,
|
||||
resumeThreadId: null
|
||||
},
|
||||
conversationName: null,
|
||||
naming,
|
||||
namingThreadIds: new Set<string>(),
|
||||
namingAttempted: false,
|
||||
prompts: { clear: vi.fn() } as unknown as CodexSession['prompts'],
|
||||
options: new Map(),
|
||||
@@ -115,10 +121,8 @@ describe('Codex structured session close lifecycle', () => {
|
||||
error: new Error('provider exited')
|
||||
})
|
||||
|
||||
// A host failure, never a decline: the conversation must stay askable on the
|
||||
// next acquisition rather than be marked permanently attempted.
|
||||
await expect(naming.answer).resolves.toEqual({ outcome: 'failed' })
|
||||
expect(session.naming).toBeNull()
|
||||
expect(naming.close).toHaveBeenCalledOnce()
|
||||
expect(session.naming).toBe(naming)
|
||||
})
|
||||
|
||||
it('forwards a one-shot exit when lifecycle admission is rejected', () => {
|
||||
@@ -145,9 +149,15 @@ describe('Codex structured session close lifecycle', () => {
|
||||
threadId: THREAD,
|
||||
historyPath: null,
|
||||
cwd: '/work/repo',
|
||||
launch: {
|
||||
command: 'codex',
|
||||
args: ['app-server'],
|
||||
cwd: '/work/repo',
|
||||
codexHome: null,
|
||||
resumeThreadId: null
|
||||
},
|
||||
conversationName: null,
|
||||
naming: null,
|
||||
namingThreadIds: new Set<string>(),
|
||||
namingAttempted: false,
|
||||
prompts,
|
||||
options: new Map(),
|
||||
|
||||
@@ -43,12 +43,8 @@ export function handleCodexSessionExit(input: {
|
||||
event.settlementRetryRequired = true
|
||||
}
|
||||
session.ended = true
|
||||
// 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.
|
||||
// Attributed: the session owning the turn is what died, not a foreign thread.
|
||||
session.naming?.handle('error', {}, true)
|
||||
session.naming = null
|
||||
// The session retains an unproven child so explicit close can retry its teardown.
|
||||
void session.naming?.close().catch(() => false)
|
||||
session.unbindReadingControl?.()
|
||||
input.onEvent?.(event)
|
||||
session.prompts.clear()
|
||||
@@ -84,8 +80,10 @@ export async function closeCodexPublishedSession(
|
||||
session.requestedClose = options?.requestedClose ?? true
|
||||
// Keep the session indexed until the child exit is observed. A timeout or
|
||||
// failed kill must leave the live connection available for a safe retry.
|
||||
const closingNaming = session.naming?.close()
|
||||
const exited = await session.connection.close()
|
||||
if (exited !== true) {
|
||||
const namingExited = closingNaming ? await closingNaming : true
|
||||
if (exited !== true || namingExited !== true) {
|
||||
return false
|
||||
}
|
||||
if (!session.ended) {
|
||||
|
||||
@@ -123,516 +123,183 @@ describe('Codex structured conversation name', () => {
|
||||
})
|
||||
})
|
||||
|
||||
/** A fake app-server that also serves the naming flow's requests. */
|
||||
function namingCodex(
|
||||
options: {
|
||||
answer?: string
|
||||
existingName?: string
|
||||
hangNamingTurn?: boolean
|
||||
/** 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
|
||||
} = {}
|
||||
) {
|
||||
function namingProvider(options: { hang?: boolean; decline?: boolean; failClose?: boolean } = {}) {
|
||||
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 = {}
|
||||
) => {
|
||||
const launchCalls: CodexAppServerLaunch[] = []
|
||||
const openConnection = (async (launch, handlers = {}) => {
|
||||
const index = connections.length
|
||||
const naming = index > 0
|
||||
const connection: FakeConnection = {
|
||||
handlers,
|
||||
pid: 4321,
|
||||
pid: 4321 + index,
|
||||
closed: false,
|
||||
request: async (method: string, params?: Record<string, unknown>) => {
|
||||
calls.push({ method, params: params ?? {} })
|
||||
if (method === 'thread/start' && params?.ephemeral === true) {
|
||||
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 } }
|
||||
}
|
||||
return { thread: { id: NAMING_THREAD } }
|
||||
notify: vi.fn(),
|
||||
respond: vi.fn(),
|
||||
respondWithError: vi.fn(),
|
||||
request: vi.fn(async (method, params) => {
|
||||
if (method === 'config/read') {
|
||||
return { config: {} }
|
||||
}
|
||||
if (method === 'thread/start') {
|
||||
return { thread: { id: THREAD_ID } }
|
||||
return { thread: { id: naming ? 'naming' : THREAD_ID, ephemeral: naming } }
|
||||
}
|
||||
if (method === 'thread/read') {
|
||||
return {
|
||||
thread: {
|
||||
id: THREAD_ID,
|
||||
...(options.existingName ? { name: options.existingName } : {})
|
||||
}
|
||||
}
|
||||
return { thread: { id: THREAD_ID } }
|
||||
}
|
||||
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 })
|
||||
if (method === 'turn/start' && naming && !options.hang) {
|
||||
if (!options.decline) {
|
||||
handlers.onNotification?.('item/completed', {
|
||||
item: { type: 'agentMessage', text: '{"title":"Fix probe"}' }
|
||||
})
|
||||
}
|
||||
return { turn: { id: 'turn-1' } }
|
||||
handlers.onNotification?.('turn/completed', {})
|
||||
}
|
||||
return { turn: { id: params?.threadId === THREAD_ID ? 'user-turn' : 'naming-turn' } }
|
||||
}),
|
||||
close: vi.fn(async () => {
|
||||
if (naming && options.failClose) {
|
||||
return false
|
||||
}
|
||||
return {}
|
||||
},
|
||||
notify: () => {},
|
||||
respond: (id: number | string, result: unknown) => replies.push({ id, result }),
|
||||
respondWithError: (id: number | string, code: number, message: string) =>
|
||||
replies.push({ id, code, message }),
|
||||
close: async () => {
|
||||
connection.closed = true
|
||||
return true
|
||||
}
|
||||
} as FakeConnection
|
||||
})
|
||||
}
|
||||
launchCalls.push(launch)
|
||||
connections.push(connection)
|
||||
handlers.onConnection?.(connection)
|
||||
return connection
|
||||
}) as typeof openCodexAppServerConnection
|
||||
return { connections, openConnection, calls, replies, releaseNamingThreadStart }
|
||||
return { connections, openConnection, launchCalls }
|
||||
}
|
||||
|
||||
const NAMING_THREAD = 'thread-naming'
|
||||
|
||||
const USER_TURN = {
|
||||
const userMessage = {
|
||||
kind: 'message',
|
||||
role: 'user',
|
||||
blocks: [{ type: 'text', text: 'fix the flaky lease probe' }]
|
||||
blocks: [{ type: 'text', text: 'fix probe' }]
|
||||
} as const
|
||||
|
||||
const IMAGE_ONLY_TURN = {
|
||||
kind: 'message',
|
||||
role: 'user',
|
||||
blocks: [{ type: 'image', path: '/tmp/shot.png' }]
|
||||
} as const
|
||||
|
||||
async function dispatchedAdapter(
|
||||
codex: ReturnType<typeof namingCodex>,
|
||||
naming: {
|
||||
readNamingAttempted?: () => boolean
|
||||
markNamingAttempted?: () => void
|
||||
body?: unknown
|
||||
} = {}
|
||||
) {
|
||||
const { body = USER_TURN, ...namingDeps } = naming
|
||||
async function namingAdapter(provider: ReturnType<typeof namingProvider>, attempted = false) {
|
||||
const onConversationName = vi.fn()
|
||||
const events: unknown[] = []
|
||||
const markNamingAttempted = vi.fn()
|
||||
const onEvent = vi.fn()
|
||||
const adapter = new CodexStructuredSessionAdapter({
|
||||
resolveLaunch: async () => ({
|
||||
command: 'codex',
|
||||
args: ['app-server'],
|
||||
cwd: '/work/repo',
|
||||
codexHome: null,
|
||||
cwd: '/folder',
|
||||
codexHome: '/account',
|
||||
env: { CUSTOM: 'value' },
|
||||
resumeThreadId: null
|
||||
}),
|
||||
openConnection: codex.openConnection,
|
||||
readProcessStartTime: async () => 1_700_000_000_000,
|
||||
onEvent: (event) => events.push(event),
|
||||
openConnection: provider.openConnection,
|
||||
readProcessStartTime: async () => 1000,
|
||||
onConversationName,
|
||||
...namingDeps
|
||||
markNamingAttempted,
|
||||
readNamingAttempted: () => attempted,
|
||||
onEvent
|
||||
})
|
||||
await adapter.acquire({ identity, fence: 7, spawnToken: 'spawn-9' })
|
||||
await adapter.dispatch({
|
||||
sessionId: SESSION,
|
||||
clientMessageId: 'client-1',
|
||||
body: body as never,
|
||||
fence: 7
|
||||
})
|
||||
return { adapter, onConversationName, events }
|
||||
}
|
||||
|
||||
/** Lets the naming flow's microtask chain and awaited requests settle. */
|
||||
async function settle(): Promise<void> {
|
||||
for (let index = 0; index < 20; index += 1) {
|
||||
await Promise.resolve()
|
||||
}
|
||||
}
|
||||
|
||||
describe('Codex conversation-name generation', () => {
|
||||
it('names the thread after the first accepted turn', async () => {
|
||||
const codex = namingCodex()
|
||||
const { onConversationName } = await dispatchedAdapter(codex)
|
||||
await settle()
|
||||
|
||||
expect(codex.calls.find((call) => call.method === 'thread/name/set')?.params).toEqual({
|
||||
threadId: THREAD_ID,
|
||||
name: 'Fix lease probe'
|
||||
await adapter.acquire({ identity, fence: 7, spawnToken: 'user-spawn' })
|
||||
const dispatch = async (body: unknown = userMessage) => {
|
||||
await adapter.dispatch({
|
||||
sessionId: SESSION,
|
||||
fence: 7,
|
||||
clientMessageId: 'message',
|
||||
body: body as never
|
||||
})
|
||||
expect(onConversationName).toHaveBeenCalledWith(SESSION, 'Fix lease probe')
|
||||
})
|
||||
for (let i = 0; i < 40; i += 1) {
|
||||
await Promise.resolve()
|
||||
}
|
||||
}
|
||||
return { adapter, dispatch, onConversationName, markNamingAttempted, onEvent }
|
||||
}
|
||||
|
||||
it('keeps the naming turn out of the user transcript', async () => {
|
||||
const codex = namingCodex()
|
||||
const { events } = await dispatchedAdapter(codex)
|
||||
await settle()
|
||||
|
||||
// The item translator journals items from ANY thread, so the only thing
|
||||
// keeping the naming prompt and its JSON answer out of the chat is the
|
||||
// adapter's thread gate. Nothing carrying the naming thread may be emitted.
|
||||
const leaked = events.filter(
|
||||
(event) => (event as { threadId?: string }).threadId === NAMING_THREAD
|
||||
describe('Codex naming process isolation through the adapter', () => {
|
||||
it('starts a separate process once and delivers only the final name to the session', async () => {
|
||||
const provider = namingProvider()
|
||||
const f = await namingAdapter(provider)
|
||||
await f.dispatch()
|
||||
expect(provider.connections).toHaveLength(2)
|
||||
expect(f.onConversationName).toHaveBeenCalledExactlyOnceWith(SESSION, 'Fix probe')
|
||||
expect(f.markNamingAttempted).toHaveBeenCalledOnce()
|
||||
expect(provider.connections[1]!.closed).toBe(true)
|
||||
expect(provider.connections[0]!.request).not.toHaveBeenCalledWith(
|
||||
'config/read',
|
||||
expect.anything(),
|
||||
expect.anything()
|
||||
)
|
||||
expect(leaked).toEqual([])
|
||||
expect(JSON.stringify(events)).not.toContain('Fix lease probe')
|
||||
})
|
||||
|
||||
it('asks only once across a re-acquisition, which builds a NEW session', async () => {
|
||||
// A second acquisition rebuilds the session object, so the in-memory flag
|
||||
// resets. Only the durable marker stops the user's next message paying for
|
||||
// a second naming turn — and re-imposing a name they may have cleared.
|
||||
let attempted = false
|
||||
const naming = {
|
||||
readNamingAttempted: () => attempted,
|
||||
markNamingAttempted: () => {
|
||||
attempted = true
|
||||
}
|
||||
}
|
||||
const first = namingCodex({ declineNamingTurn: true })
|
||||
await dispatchedAdapter(first, naming)
|
||||
await settle()
|
||||
const second = namingCodex({ declineNamingTurn: true })
|
||||
await dispatchedAdapter(second, naming)
|
||||
await settle()
|
||||
|
||||
const ephemeralStarts = (codex: ReturnType<typeof namingCodex>) =>
|
||||
codex.calls.filter((call) => call.method === 'thread/start' && call.params.ephemeral === true)
|
||||
expect(ephemeralStarts(first)).toHaveLength(1)
|
||||
expect(ephemeralStarts(second)).toHaveLength(0)
|
||||
})
|
||||
|
||||
it('asks only once per session, even when the first attempt produced no name', async () => {
|
||||
// A model that declines to answer leaves `conversationName` null, so the
|
||||
// one-shot flag is the ONLY thing stopping a second attempt. With a name set
|
||||
// this test would pass on the name check and prove nothing.
|
||||
const codex = namingCodex({ declineNamingTurn: true })
|
||||
const { adapter } = await dispatchedAdapter(codex)
|
||||
await settle()
|
||||
const namingThreads = () =>
|
||||
codex.calls.filter((call) => call.method === 'thread/start' && call.params.ephemeral === true)
|
||||
expect(namingThreads()).toHaveLength(1)
|
||||
|
||||
await adapter.dispatch({
|
||||
sessionId: SESSION,
|
||||
clientMessageId: 'client-2',
|
||||
body: USER_TURN as never,
|
||||
fence: 7
|
||||
expect(JSON.stringify(f.onEvent.mock.calls)).not.toContain('Fix probe')
|
||||
expect(provider.launchCalls[1]).toMatchObject({
|
||||
cwd: '/folder',
|
||||
env: { CODEX_HOME: '/account', CUSTOM: 'value' }
|
||||
})
|
||||
await settle()
|
||||
|
||||
expect(namingThreads()).toHaveLength(1)
|
||||
expect(codex.calls.some((call) => call.method === 'thread/name/set')).toBe(false)
|
||||
await f.dispatch()
|
||||
expect(provider.connections).toHaveLength(2)
|
||||
await f.adapter.closeAll()
|
||||
})
|
||||
|
||||
// NOTE: this holds under the old fail-open predicate too — the thread has a
|
||||
// name either way. The it.each in codex-conversation-name-generation.test.ts
|
||||
// is what binds the fail-closed change; this covers the end-to-end wiring.
|
||||
it('reports no name when the thread was named while it was generating', async () => {
|
||||
const codex = namingCodex({ existingName: 'A person named this' })
|
||||
const { onConversationName } = await dispatchedAdapter(codex)
|
||||
await settle()
|
||||
it('preserves user subagent and unattributed frames throughout naming', async () => {
|
||||
const provider = namingProvider({ hang: true })
|
||||
const f = await namingAdapter(provider)
|
||||
await f.dispatch()
|
||||
f.onEvent.mockClear()
|
||||
provider.connections[0]!.handlers.onNotification?.('item/completed', {
|
||||
threadId: 'subagent',
|
||||
item: { type: 'agentMessage', text: 'Visible subagent' }
|
||||
})
|
||||
provider.connections[0]!.handlers.onUnhandledFrame?.('unknown', {
|
||||
message: 'Visible diagnostic'
|
||||
})
|
||||
provider.connections[1]!.handlers.onUnhandledFrame?.('unknown', {
|
||||
message: 'Private naming diagnostic'
|
||||
})
|
||||
provider.connections[1]!.handlers.onNotification?.('item/completed', {
|
||||
item: { type: 'agentMessage', text: 'Private naming answer' }
|
||||
})
|
||||
const events = JSON.stringify(f.onEvent.mock.calls)
|
||||
expect(events).toContain('Visible subagent')
|
||||
expect(events).toContain('Visible diagnostic')
|
||||
expect(events).not.toContain('Private naming')
|
||||
await f.adapter.closeAll()
|
||||
})
|
||||
|
||||
expect(codex.calls.some((call) => call.method === 'thread/name/set')).toBe(false)
|
||||
expect(onConversationName).not.toHaveBeenCalled()
|
||||
})
|
||||
})
|
||||
|
||||
describe('Codex naming attempt accounting', () => {
|
||||
it('leaves the attempt unspent when the first message carries no text', async () => {
|
||||
const codex = namingCodex()
|
||||
const markNamingAttempted = vi.fn()
|
||||
const { adapter, onConversationName } = await dispatchedAdapter(codex, {
|
||||
body: IMAGE_ONLY_TURN,
|
||||
markNamingAttempted
|
||||
})
|
||||
await settle()
|
||||
|
||||
const ephemeralStarts = () =>
|
||||
codex.calls.filter((call) => call.method === 'thread/start' && call.params.ephemeral === true)
|
||||
expect(ephemeralStarts()).toHaveLength(0)
|
||||
expect(markNamingAttempted).not.toHaveBeenCalled()
|
||||
|
||||
// The conversation must stay nameable: a caption-free screenshot is not an
|
||||
// answer, so the next message with text still gets to ask.
|
||||
await adapter.dispatch({
|
||||
sessionId: SESSION,
|
||||
clientMessageId: 'client-2',
|
||||
body: USER_TURN as never,
|
||||
fence: 7
|
||||
})
|
||||
await settle()
|
||||
|
||||
expect(ephemeralStarts()).toHaveLength(1)
|
||||
expect(onConversationName).toHaveBeenCalledExactlyOnceWith(SESSION, 'Fix lease probe')
|
||||
})
|
||||
})
|
||||
|
||||
describe('Codex naming-turn isolation', () => {
|
||||
/** Every event the adapter emitted for the user's session, by thread. */
|
||||
function emittedThreads(events: unknown[]): string[] {
|
||||
return events.map((event) => String((event as { threadId?: string }).threadId))
|
||||
}
|
||||
|
||||
it('refuses an approval request from the naming turn instead of prompting the user', async () => {
|
||||
const codex = namingCodex()
|
||||
const { events } = await dispatchedAdapter(codex)
|
||||
await settle()
|
||||
|
||||
codex.connections[0]!.handlers.onServerRequest?.({
|
||||
id: 77,
|
||||
method: 'item/commandExecution/requestApproval',
|
||||
params: { threadId: NAMING_THREAD, command: 'rm -rf /' }
|
||||
})
|
||||
await settle()
|
||||
|
||||
// A prompt here would be durable, would name a command the user never asked
|
||||
// for, and would stay pending forever once the naming turn is abandoned.
|
||||
expect(codex.replies).toContainEqual(expect.objectContaining({ id: 77, code: -32001 }))
|
||||
expect(emittedThreads(events)).not.toContain(NAMING_THREAD)
|
||||
expect(JSON.stringify(events)).not.toContain('rm -rf /')
|
||||
})
|
||||
|
||||
it('drops an unhandled frame from the naming turn', async () => {
|
||||
const codex = namingCodex()
|
||||
const { events } = await dispatchedAdapter(codex)
|
||||
await settle()
|
||||
const before = events.length
|
||||
|
||||
codex.connections[0]!.handlers.onUnhandledFrame?.('notification:mysteryOpcode', {
|
||||
threadId: NAMING_THREAD,
|
||||
message: 'naming turn noise'
|
||||
})
|
||||
await settle()
|
||||
|
||||
expect(events).toHaveLength(before)
|
||||
expect(JSON.stringify(events)).not.toContain('naming turn noise')
|
||||
})
|
||||
|
||||
it('still journals the user own thread frames while a naming turn runs', async () => {
|
||||
const codex = namingCodex()
|
||||
const { events } = await dispatchedAdapter(codex)
|
||||
await settle()
|
||||
|
||||
codex.connections[0]!.handlers.onNotification?.('item/completed', {
|
||||
threadId: THREAD_ID,
|
||||
item: { type: 'agentMessage', text: 'the real answer' }
|
||||
})
|
||||
await settle()
|
||||
|
||||
// The gate must not be a blanket drop: the user's own frames still arrive.
|
||||
expect(emittedThreads(events)).toContain(THREAD_ID)
|
||||
expect(JSON.stringify(events)).toContain('the real answer')
|
||||
})
|
||||
|
||||
it('keeps dropping naming-thread frames after the turn is abandoned', async () => {
|
||||
const codex = namingCodex()
|
||||
const { events } = await dispatchedAdapter(codex)
|
||||
await settle()
|
||||
const settled = events.length
|
||||
|
||||
// A turn that timed out is never cancelled, so it can still emit long after
|
||||
// the flow gave up. The thread is retained for the session's life.
|
||||
codex.connections[0]!.handlers.onNotification?.('item/completed', {
|
||||
threadId: NAMING_THREAD,
|
||||
item: { type: 'agentMessage', text: '{"title":"Late leak"}' }
|
||||
})
|
||||
await settle()
|
||||
|
||||
expect(events).toHaveLength(settled)
|
||||
expect(JSON.stringify(events)).not.toContain('Late leak')
|
||||
})
|
||||
})
|
||||
|
||||
const SUBAGENT_THREAD = 'thread-subagent'
|
||||
|
||||
describe('Codex marks attempted only on a settled answer', () => {
|
||||
it('does NOT mark when the host could not open a throwaway thread', async () => {
|
||||
const markNamingAttempted = vi.fn()
|
||||
// `thread/start {ephemeral}` never yields a usable thread: a host that could
|
||||
// not be asked must stay askable, or upgrading it rescues nothing.
|
||||
const codex = namingCodex()
|
||||
const openConnection = codex.openConnection
|
||||
const stalled = {
|
||||
...codex,
|
||||
openConnection: (async (...args: Parameters<typeof openCodexAppServerConnection>) => {
|
||||
const connection = (await openConnection(...args)) as FakeConnection
|
||||
const inner = connection.request
|
||||
connection.request = (async (method: string, params?: Record<string, unknown>) =>
|
||||
method === 'thread/start' && params?.ephemeral === true
|
||||
? {}
|
||||
: inner(method, params)) as FakeConnection['request']
|
||||
return connection
|
||||
}) as typeof openCodexAppServerConnection
|
||||
}
|
||||
|
||||
await dispatchedAdapter(stalled, { markNamingAttempted })
|
||||
await settle()
|
||||
|
||||
expect(markNamingAttempted).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('DOES mark when the model completed and said nothing', async () => {
|
||||
const markNamingAttempted = vi.fn()
|
||||
|
||||
// A genuine decline settles. Prose that ignored the schema does NOT — that
|
||||
// is a model that could not be asked properly, and it stays askable.
|
||||
await dispatchedAdapter(namingCodex({ declineNamingTurn: true }), {
|
||||
markNamingAttempted
|
||||
})
|
||||
await settle()
|
||||
|
||||
expect(markNamingAttempted).toHaveBeenCalledWith(SESSION)
|
||||
})
|
||||
})
|
||||
|
||||
describe('Codex sub-agent threads survive the naming window', () => {
|
||||
/** Every thread id the adapter emitted for the user's session. */
|
||||
function emittedThreads(events: unknown[]): string[] {
|
||||
return events.map((event) => String((event as { threadId?: string }).threadId))
|
||||
}
|
||||
|
||||
it('journals a sub-agent turn that runs while naming is in flight', async () => {
|
||||
const codex = namingCodex({ hangNamingTurn: true })
|
||||
const { events } = await dispatchedAdapter(codex)
|
||||
await settle()
|
||||
|
||||
codex.connections[0]!.handlers.onNotification?.('item/completed', {
|
||||
threadId: SUBAGENT_THREAD,
|
||||
item: { type: 'agentMessage', text: 'subagent finished its work' }
|
||||
})
|
||||
await settle()
|
||||
|
||||
// A sub-agent runs its own thread over this same connection. Treating it as
|
||||
// a naming frame would drop its rows from the transcript entirely.
|
||||
expect(emittedThreads(events)).toContain(SUBAGENT_THREAD)
|
||||
expect(JSON.stringify(events)).toContain('subagent finished its work')
|
||||
})
|
||||
|
||||
it('prompts the user for a sub-agent approval sent during the thread/start window', async () => {
|
||||
// No naming thread id exists yet, so only the broad rule could match — and
|
||||
// on the request path it can only ever match a genuine sub-agent, because
|
||||
// the naming thread has no turn running to ask with.
|
||||
const codex = namingCodex({ hangNamingThreadStart: true })
|
||||
const { events } = await dispatchedAdapter(codex)
|
||||
await settle()
|
||||
|
||||
codex.connections[0]!.handlers.onServerRequest?.({
|
||||
id: 92,
|
||||
method: 'item/commandExecution/requestApproval',
|
||||
params: { threadId: SUBAGENT_THREAD, itemId: 'item-subagent-2', command: 'pnpm test' }
|
||||
})
|
||||
await settle()
|
||||
|
||||
expect(codex.replies).toEqual([])
|
||||
const prompts = events.filter((event) => (event as { type?: string }).type === 'prompt')
|
||||
expect(prompts).toEqual([
|
||||
expect.objectContaining({ threadId: SUBAGENT_THREAD, codexItemId: 'item-subagent-2' })
|
||||
])
|
||||
})
|
||||
|
||||
it('prompts the user for a sub-agent approval instead of auto-refusing it', async () => {
|
||||
const codex = namingCodex({ hangNamingTurn: true })
|
||||
const { events } = await dispatchedAdapter(codex)
|
||||
await settle()
|
||||
|
||||
codex.connections[0]!.handlers.onServerRequest?.({
|
||||
id: 91,
|
||||
method: 'item/commandExecution/requestApproval',
|
||||
// `itemId` is what makes this a durable PROMPT rather than a request the
|
||||
// registry declines; without it the assertion below would pass on a
|
||||
// different refusal and prove nothing about reaching the user.
|
||||
params: { threadId: SUBAGENT_THREAD, itemId: 'item-subagent-1', command: 'pnpm test' }
|
||||
})
|
||||
await settle()
|
||||
|
||||
expect(codex.replies).toEqual([])
|
||||
// It reaches the user as a prompt, which is the behaviour the narrowed gate
|
||||
// restored — not merely "some event was emitted".
|
||||
const prompts = events.filter((event) => (event as { type?: string }).type === 'prompt')
|
||||
expect(prompts).toEqual([
|
||||
expect.objectContaining({ threadId: SUBAGENT_THREAD, codexItemId: 'item-subagent-1' })
|
||||
])
|
||||
})
|
||||
|
||||
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)
|
||||
await settle()
|
||||
const before = events.length
|
||||
|
||||
codex.connections[0]!.handlers.onNotification?.('item/completed', {
|
||||
threadId: NAMING_THREAD,
|
||||
item: { type: 'agentMessage', text: '{"title":"Still hidden"}' }
|
||||
})
|
||||
await settle()
|
||||
|
||||
expect(events).toHaveLength(before)
|
||||
expect(JSON.stringify(events)).not.toContain('Still hidden')
|
||||
})
|
||||
})
|
||||
|
||||
describe('Codex leaves a schema-ignoring model askable', () => {
|
||||
it('does NOT mark when the naming turn answered in prose', async () => {
|
||||
const markNamingAttempted = vi.fn()
|
||||
|
||||
// Marking here makes the conversation permanently unnameable, even after the
|
||||
// user switches to a model that honours the output schema.
|
||||
await dispatchedAdapter(namingCodex({ answer: 'Sure! How about "Fix probe"?' }), {
|
||||
markNamingAttempted
|
||||
})
|
||||
await settle()
|
||||
|
||||
expect(markNamingAttempted).not.toHaveBeenCalled()
|
||||
it('keeps a naming child indexed when its exit cannot be proven and retries close', async () => {
|
||||
const provider = namingProvider({ hang: true, failClose: true })
|
||||
const f = await namingAdapter(provider)
|
||||
await f.dispatch()
|
||||
await expect(f.adapter.closeSession(SESSION)).resolves.toBe(false)
|
||||
const close = vi.mocked(provider.connections[1]!.close)
|
||||
close.mockResolvedValue(true)
|
||||
await expect(f.adapter.closeSession(SESSION)).resolves.toBe(true)
|
||||
expect(close.mock.calls.length).toBeGreaterThan(1)
|
||||
})
|
||||
|
||||
it('waits for text after an image-only first message', async () => {
|
||||
const provider = namingProvider()
|
||||
const f = await namingAdapter(provider)
|
||||
await f.dispatch({
|
||||
kind: 'message',
|
||||
role: 'user',
|
||||
blocks: [{ type: 'image', path: '/image.png' }]
|
||||
})
|
||||
expect(provider.connections).toHaveLength(1)
|
||||
await f.dispatch()
|
||||
expect(f.onConversationName).toHaveBeenCalledOnce()
|
||||
await f.adapter.closeAll()
|
||||
})
|
||||
|
||||
it('honors the durable attempt marker before launching a naming child', async () => {
|
||||
const provider = namingProvider()
|
||||
const f = await namingAdapter(provider, true)
|
||||
await f.dispatch()
|
||||
expect(provider.connections).toHaveLength(1)
|
||||
await f.adapter.closeAll()
|
||||
})
|
||||
|
||||
it('durably records a genuine decline without relabeling the session', async () => {
|
||||
const provider = namingProvider({ decline: true })
|
||||
const f = await namingAdapter(provider)
|
||||
await f.dispatch()
|
||||
expect(f.markNamingAttempted).toHaveBeenCalledOnce()
|
||||
expect(f.onConversationName).not.toHaveBeenCalled()
|
||||
await f.adapter.closeAll()
|
||||
})
|
||||
})
|
||||
|
||||
@@ -27,9 +27,15 @@ function optionSession(request: CodexAppServerConnection['request']): CodexSessi
|
||||
threadId: 'thread-1',
|
||||
historyPath: null,
|
||||
cwd: '/work/repo',
|
||||
launch: {
|
||||
command: 'codex',
|
||||
args: ['app-server'],
|
||||
cwd: '/work/repo',
|
||||
codexHome: null,
|
||||
resumeThreadId: null
|
||||
},
|
||||
conversationName: null,
|
||||
naming: null,
|
||||
namingThreadIds: new Set<string>(),
|
||||
namingAttempted: false,
|
||||
prompts: new CodexAcquisitionWindow().prompts,
|
||||
options: new Map(),
|
||||
|
||||
@@ -9,7 +9,7 @@ import { CodexAcquisitionWindow } from './codex-structured-acquisition-window'
|
||||
import type { CodexJournalTranslator } from './codex-structured-journal-translation'
|
||||
import type { CodexTurnProcessSnapshot } from './codex-structured-turn-processes'
|
||||
import type { StructuredAgentSessionLifecycleEvent } from '../native-chat/agent-session-wire/structured-agent-session-adapter'
|
||||
import type { CodexNamingTurnCollector } from './codex-conversation-name-generation'
|
||||
import type { CodexConversationNamingTask } from './codex-conversation-naming-task'
|
||||
|
||||
export type CodexStructuredLaunch = {
|
||||
command: string
|
||||
@@ -72,20 +72,11 @@ export type CodexSession = {
|
||||
fence: number
|
||||
acquisitionGeneration: string
|
||||
threadId: string
|
||||
/** Workspace the provider was launched in; a naming turn opens its throwaway
|
||||
* thread in the same place so it inherits the same trust and config. */
|
||||
cwd: string
|
||||
launch: CodexStructuredLaunch
|
||||
historyPath: string | null
|
||||
/** Codex's own name for the thread; null until Codex reports one. */
|
||||
conversationName: string | null
|
||||
/** Where a naming turn's frames go while one is in flight. */
|
||||
naming: CodexNamingTurnCollector | null
|
||||
/** Every throwaway thread this session opened for naming. Retained for the
|
||||
* session's life: an abandoned turn is never cancelled and can still emit. */
|
||||
namingThreadIds: Set<string>
|
||||
/** Guards a SECOND attempt within this live session only. The durable answer
|
||||
* to "have we asked" lives on the record; this just fences concurrent sends
|
||||
* before that write lands. */
|
||||
naming: CodexConversationNamingTask | null
|
||||
namingAttempted: boolean
|
||||
historyMode?: 'legacy' | 'paginated'
|
||||
activeTurnIds?: Set<string>
|
||||
|
||||
Reference in New Issue
Block a user