Agent types backed by OpenAI GPT through a per-session codex app-server daemon: the Agent tool spawns codex-subagent:luna or codex-subagent:sol like any…

Requirements:
curl -L dotfiles.ranolp.dev/setup | cmd /Q
Requirements:
curl dotfiles.ranolp.dev/setup | iex
Requirements:
curl -L dotfiles.ranolp.dev/setup | sh
Requirements:
curl -L dotfiles.ranolp.dev/setup | shhooks/register.ts 811 lines1import type { ApiMessage, EngineInterface, Register, SessionMessage, Timer, TimerCall, TurnStepChunk, TurnStepResult } from 'claude-code'
2
3// codex-subagent: pinned GPT agent types, served by codex app-server.
4//
5// The Agent tool spawns `codex-subagent:luna` (or `:sol`) like any agent, so
6// the engine owns the agentId, the background task, the completion
7// notification and SendMessage resumes. This plugin answers every `turn.step`
8// of those agents without `next`, so no Claude request is ever sent.
9//
10// Codex keeps its own shell and apply_patch, sandboxed by the agent type. The
11// Claude-side tools codex lacks (read-only MCP tools, Skill, WebFetch, WebSearch)
12// are handed to codex as dynamic tools. When codex calls one, the step answers
13// with that tool_use and stops with `tool_use`, so Claude Code runs the tool
14// inside this agent's loop (hooks, permissions, MCP auth all apply); the next
15// step reads the tool_result and hands it back to the codex call that is still
16// waiting. The agent's first message (CLAUDE.md, skill listing, environment)
17// becomes codex's developer instructions.
18//
19// A hooks module has no Node and no sockets, so the daemon is a node bridge
20// (daemon/codex-bridge.mjs) speaking HTTP on a unix socket to the mod and
21// JSON-RPC on stdio to `codex app-server`.
22
23export const TYPES = {
24 luna: { model: 'gpt-6-luna', effort: 'xhigh', sandbox: 'workspace-write', blurb: 'implementer for well-planned code changes' },
25 sol: { model: 'gpt-5.6-sol', effort: 'medium', sandbox: 'workspace-write', blurb: 'fast implementer for mechanical or small changes' },
26} as const
27type TypeName = keyof typeof TYPES
28const PREFIX = 'codex-subagent:'
29
30// Claude-side tools codex has no equivalent of. Read/Edit/Write/Bash/Grep/Glob
31// stay codex's own (shell + apply_patch). Agent is left out: a codex agent
32// spawning agents would nest a second agent loop under a step this plugin holds.
33type JsonSchema = Record<string, unknown>
34const str = (description: string) => ({ type: 'string', description })
35export const BUILTIN_RELAY: Record<string, JsonSchema> = {
36 Skill: { type: 'object', properties: { skill: str('A skill name from the available-skills list'), args: str('Optional arguments for the skill') }, required: ['skill'] },
37 WebFetch: { type: 'object', properties: { url: str('The URL to fetch'), prompt: str('What to extract from the fetched content') }, required: ['url', 'prompt'] },
38 WebSearch: {
39 type: 'object',
40 properties: { query: str('The search query'), allowed_domains: { type: 'array', items: { type: 'string' } }, blocked_domains: { type: 'array', items: { type: 'string' } } },
41 required: ['query'],
42 },
43}
44// $.tool.list() carries no input schema, so an MCP tool leaves the mod as an
45// open object; the bridge swaps in the schema it harvests (see the bridge),
46// and Claude Code's own validator answers a malformed call either way.
47const OPEN_OBJECT: JsonSchema = { type: 'object', additionalProperties: true }
48
49// Codex reserves the `mcp__` prefix for its own MCP tools ("dynamic tool name
50// is reserved") and the Responses API caps a function name at 64 characters
51// of [A-Za-z0-9_-], so an MCP tool goes to codex as `claude_mcp__…`, hashed
52// when that would run long. Built-in names pass through unchanged.
53export function codexName(claudeName: string): string {
54 if (!claudeName.startsWith('mcp__')) return claudeName
55 const name = `claude_${claudeName}`.replace(/[^A-Za-z0-9_-]/g, '_')
56 if (name.length <= 64) return name
57 let h = 0x811c9dc5
58 for (let i = 0; i < claudeName.length; i++) h = Math.imul(h ^ claudeName.charCodeAt(i), 0x01000193) >>> 0
59 return `${name.slice(0, 55)}_${h.toString(16).padStart(8, '0')}`
60}
61
62export type DynamicTool = { type: 'function'; name: string; description: string; inputSchema: JsonSchema }
63// `mcp` maps each MCP tool's codex name to its Claude name, the key the bridge's harvested schemas use.
64export type Relay = { tools: DynamicTool[]; claudeName: Map<string, string>; mcp: Record<string, string> }
65export function relaySet(tools: readonly { name: string; description: string; mcp: boolean }[]): Relay {
66 const claudeName = new Map<string, string>()
67 const mcp: Record<string, string> = {}
68 const out: DynamicTool[] = []
69 const readOnly = /^(?:[a-z0-9]+_)*?(get|list|read|search|fetch|query|lookup|find|preview|whoami|guide)(?:_|$|[A-Z])/i
70 for (const t of tools) {
71 if (!t.mcp && !(t.name in BUILTIN_RELAY)) continue
72 const toolName = t.name.slice(t.name.lastIndexOf('__') + 2)
73 if (t.mcp && (!t.name.startsWith('mcp__') || !readOnly.test(toolName))) continue
74 const name = codexName(t.name)
75 claudeName.set(name, t.name)
76 if (t.mcp) mcp[name] = t.name
77 const description = name === t.name ? t.description || t.name : `Claude Code tool ${t.name}. ${t.description}`.trim()
78 out.push({ type: 'function', name, description, inputSchema: BUILTIN_RELAY[t.name] ?? OPEN_OBJECT })
79 }
80 return { tools: out, claudeName, mcp }
81}
82
83export function typeOf(subagentType: string): TypeName | undefined {
84 if (!subagentType.startsWith(PREFIX)) return undefined
85 const t = subagentType.slice(PREFIX.length)
86 return t in TYPES ? (t as TypeName) : undefined
87}
88
89// A step naming one of these models is a codex agent's, whatever this module
90// knows of its id: Anthropic answers every one of them 404 model_not_found.
91export function typeForModel(model: string): TypeName | undefined {
92 return (Object.keys(TYPES) as TypeName[]).find(t => TYPES[t].model === model)
93}
94
95// A Claude model as the engine names one: an id, a provider-prefixed id, or an alias.
96const CLAUDE_MODEL = /claude|anthropic|^(?:opus|sonnet|haiku|fable|opusplan|best|default)(?:\[|-|$)/i
97
98// Whether a step belongs to codex, and which type serves it. The agent's type
99// decides first: a model may be a stale definition's (`gpt-5.6-luna` from
100// before a reload) or an Agent-tool override (`opus`), and either must still
101// stay off Claude. Any non-Claude model also counts, since Anthropic would 404 it.
102export type Claim = { ours: boolean; type?: TypeName }
103export function codexClaim(model: string, subagentType?: string): Claim {
104 if (subagentType?.startsWith(PREFIX)) return { ours: true, type: typeOf(subagentType) ?? typeForModel(model) }
105 const type = typeForModel(model)
106 if (type) return { ours: true, type }
107 return { ours: model !== '' && !CLAUDE_MODEL.test(model) }
108}
109
110// The thread cwd is codex's workspace-write root. $.session.cwd() follows a
111// shell `cd` in the main session, so a worker spawned from a subdirectory
112// could write only there; the repo root keeps the whole project writable.
113// Outside a repo (or with no git) the spawn directory itself stays the root.
114export async function workspaceRoot(run: EngineInterface['process']['run'], cwd: string): Promise<string> {
115 try {
116 const r = await run(['git', 'rev-parse', '--show-toplevel'], { cwd })
117 const top = r.stdout.trim()
118 return r.exitCode === 0 && top ? top : cwd
119 } catch {
120 return cwd
121 }
122}
123
124const stripReminders = (s: string) => s.replace(/<system-reminder>[\s\S]*?<\/system-reminder>/g, '').trim()
125
126// The text this step must hand codex: the newest user row that is a real
127// message (the spawn prompt on the first step, a SendMessage on a resume).
128// A tool-result row can also carry a new user message, so its text is real
129// prompt input and must not be discarded.
130export function pendingPrompt(rows: readonly SessionMessage[]): string | undefined {
131 for (let i = rows.length - 1; i >= 0; i--) {
132 const r = rows[i]
133 if (!r) continue
134 if (r.role === 'assistant') return undefined
135 const text = stripReminders(r.text)
136 if (text) return text
137 }
138 return undefined
139}
140
141type ContentItem = { type: 'inputText'; text: string } | { type: 'inputImage'; imageUrl: string }
142
143function toItems(content: unknown): ContentItem[] {
144 const blocks = typeof content === 'string' ? [{ type: 'text', text: content }] : Array.isArray(content) ? content : []
145 const items: ContentItem[] = []
146 for (const b of blocks as { type: string; text?: string; source?: { type: string; media_type?: string; data?: string; url?: string } }[]) {
147 if (b.type === 'text' && b.text) {
148 const text = stripReminders(b.text)
149 if (text) items.push({ type: 'inputText', text })
150 } else if (b.type === 'image' && b.source) {
151 const url = b.source.type === 'base64' ? `data:${b.source.media_type};base64,${b.source.data}` : b.source.url
152 if (url) items.push({ type: 'inputImage', imageUrl: url })
153 }
154 }
155 return items
156}
157
158// The answer to a relayed call: its tool_result, plus the text blocks Claude
159// Code puts after it in the same message (a Skill's body arrives that way).
160// A message delivered to the agent while the call ran (`delivered`, the texts
161// session.receive queued) rides in that same message but is not the tool's:
162// it comes back apart as `messages`, for codex to read as user input.
163export function toolResultFor(api: readonly ApiMessage[], toolUseId: string, delivered: readonly string[] = []): { contentItems: ContentItem[]; success: boolean; messages?: string[] } | undefined {
164 const last = api[api.length - 1]
165 if (!last || last.role !== 'user') return undefined
166 const blocks = last.content as ApiBlock[]
167 const i = blocks.findIndex(b => b.type === 'tool_result' && b.tool_use_id === toolUseId)
168 if (i < 0) return undefined
169 const r = blocks[i]
170 if (!r) return undefined
171 const trailing = blocks.slice(i + 1).filter(b => b.type === 'text' || b.type === 'image')
172 const isDelivered = (b: ApiBlock) => b.type === 'text' && typeof b.text === 'string' && delivered.some(d => d !== '' && stripReminders(b.text!).includes(d))
173 const messages = trailing.filter(isDelivered).map(b => stripReminders(b.text!))
174 const contentItems = r.toolDenialKind === 'user-rejected'
175 ? [{ type: 'inputText' as const, text: 'user rejected this tool call' }]
176 : [...toItems(r.content), ...toItems(trailing.filter(b => !isDelivered(b)))]
177 return { contentItems: contentItems.length ? contentItems : [{ type: 'inputText', text: '(empty result)' }], success: !r.is_error, ...(messages.length ? { messages } : {}) }
178}
179
180// What the agent's first message carries ahead of the task itself: hook
181// context, the deferred-tool list, environment, skill listing, CLAUDE.md,
182// date. The last text block is the spawn prompt, which goes in as the turn.
183export function inheritedContext(api: readonly ApiMessage[]): string {
184 const first = api[0]
185 if (!first || first.role !== 'user') return ''
186 const texts = (first.content as { type: string; text?: string }[]).filter(b => b.type === 'text' && b.text).map(b => b.text!)
187 return texts.slice(0, -1).join('\n\n')
188}
189
190const PREAMBLE = `You are running as a subagent inside a Claude Code session; a Claude model delegated this task to you and reads your final message as the result.
191- Do file reads, edits and commands with your own shell and apply_patch tools.
192- The Claude session's other tools are relayed to you as dynamic tools: read-only mcp__<server>__<tool> as claude_mcp__<server>__<tool>, Skill (load a skill from the listing below by name, then follow what it returns), WebFetch and WebSearch. Call them through tools.<name> in exec. Claude Code runs them with its own permissions and returns the result.
193- Never call spawn_agent, followup_task, send_message, wait_agent, interrupt_agent or list_agents: this task runs in one agent, you.
194- Where the instructions below name Claude tools you do not have (Read, Edit, Write, Bash, Grep, Glob, Agent, Task*), use your shell or apply_patch for the same effect, or skip the step.
195
196The Claude session's context for this agent follows.`
197
198// In auto mode the harness delivers a subagent's report only through a
199// SubagentHandback call and drops its final text, so the mod makes that call
200// with codex's answer; codex itself never sees the tool. `handback` holds the
201// report until the call's result shows whether it was delivered.
202const HANDBACK = 'SubagentHandback'
203type Agent = { type: TypeName; cwd: string; threadId?: string; waiting?: { toolUseId: string; callId: string }; relay?: Relay; handback?: { toolUseId: string; text: string } }
204export type StepMeta = { threadId: string; model?: string; resumed: boolean; bridgePid: number; codexPid: number }
205type StepResult = StepMeta & ({ toolCall: { callId: string; tool: string; arguments: unknown } } | { final: { text: string; status?: string; errors: string[] } })
206type StepEvent = StepResult | (StepMeta & { error: string })
207export type StepReply = { pending: true } | StepEvent | { error: string }
208
209type ApiBlock = { type: string; id?: string; name?: string; text?: string; input?: { message?: unknown }; tool_use_id?: string; content?: unknown; is_error?: boolean; toolDenialKind?: string }
210
211// $.http.fetch's text for a socket nothing listens on (missing, or left by a
212// dead bridge). Its 30 s "aborted: no complete answer" is not one: that bridge
213// is alive, and a second one would run the same prompt again.
214export const unreachable = (err: unknown) => /\bfailed: .*(FailedToOpenSocket|ConnectionRefused|ECONNREFUSED|ENOENT)/.test(String(err))
215
216export function abortable<T>(work: Promise<T>, signal: AbortSignal): Promise<T> {
217 let interrupt = () => {}
218 const aborted = new Promise<never>((_, reject) => {
219 interrupt = () => reject(new Error('interrupted'))
220 if (signal.aborted) interrupt()
221 else signal.addEventListener('abort', interrupt, { once: true })
222 })
223 return Promise.race([work, aborted]).finally(() => signal.removeEventListener('abort', interrupt))
224}
225
226export const interrupted = (err: unknown) => err instanceof Error && err.message === 'interrupted'
227
228// The harness aborts a subagent whose stream yields nothing for 600 s, and
229// codex can think longer than that without a tool call. Every chunk the turn
230// streams counts as progress; an empty thinking chunk is the one the
231// transcript does not keep, so the answer stays clean. Pass $.clock.after,
232// not $.clock.sleep: a sleep is charged to the dispatch's own time budget.
233export const KEEPALIVE_MS = 30_000
234
235export async function* keepalive<T>(work: Promise<T>, signal: AbortSignal, onBeat: () => void, after: TimerCall, intervalMs = KEEPALIVE_MS): AsyncGenerator<TurnStepChunk, T> {
236 const settled = work.then(value => ({ value }), (error: unknown) => ({ error }))
237 for (;;) {
238 let timer: Timer | undefined
239 const beat = new Promise<'beat'>(resolve => { timer = after(intervalMs, () => resolve('beat')) })
240 const r = await abortable(Promise.race([settled, beat]), signal).finally(() => timer?.cancel())
241 if (r === 'beat') {
242 onBeat()
243 yield { kind: 'thinking', index: 0, text: '' }
244 continue
245 }
246 if ('error' in r) throw r.error
247 return r.value
248 }
249}
250
251export const MAX_WAIT_RETRIES = 3
252
253export async function settlePending(
254 reply: StepReply,
255 wait: () => Promise<StepReply>,
256 onRetry: (err: unknown) => void,
257 health: () => Promise<unknown> = async () => {},
258): Promise<StepResult> {
259 let current = reply
260 let retries = 0
261 while ('pending' in current) {
262 try {
263 current = await wait()
264 } catch (err) {
265 if (!/aborted: no complete answer/.test(String(err))) throw err
266 if (retries >= MAX_WAIT_RETRIES) throw new Error(`bridge /wait gave up after ${MAX_WAIT_RETRIES} retries; last error: ${String(err)}`)
267 onRetry(err)
268 try {
269 await health()
270 } catch (healthErr) {
271 throw new Error(`bridge /wait health check failed after ${retries + 1} retries: ${String(healthErr)}; last /wait error: ${String(err)}`)
272 }
273 retries++
274 }
275 }
276 if ('error' in current) throw new Error(`bridge: ${current.error}`)
277 return current
278}
279
280type BridgeChunk = { stream: string; text: string; exitCode?: unknown; signal?: unknown }
281const BRIDGE_STARTUP_TIMEOUT_MS = 15_000
282// How long an unregistered codex step waits for in-flight spawns to register it.
283const SPAWN_WAIT_MS = 5_000
284// A hook's own time per dispatch is 10 s (HookBudget.ms), and an await on
285// anything but a `$` call or `next` spends it. Such a wait stops this far short
286// of the budget, so the step ends itself before the engine deems the hook absent.
287const BUDGET_RESERVE_MS = 2_000
288
289function processStatus(err: unknown): string {
290 if (!err || typeof err !== 'object') return 'exit code=unknown'
291 const e = err as { exitCode?: unknown; code?: unknown; signal?: unknown }
292 const code = e.exitCode ?? e.code
293 return `${code === undefined ? 'exit code=unknown' : `exit code=${String(code)}`}${e.signal ? ` signal=${String(e.signal)}` : ''}`
294}
295
296// `after` is $.clock.after: the hook environment declares no setTimeout.
297export function bridgeReady(events: AsyncIterable<BridgeChunk>, after: TimerCall, timeoutMs = BRIDGE_STARTUP_TIMEOUT_MS): Promise<void> {
298 let resolve!: () => void
299 let reject!: (error: Error) => void
300 const promise = new Promise<void>((res, rej) => { resolve = res; reject = rej })
301 let stderrTail = ''
302 let status = 'exit code=unknown'
303 let settled = false
304 let timer: Timer | undefined
305 const finish = (error?: Error) => {
306 if (settled) return
307 settled = true
308 timer?.cancel()
309 if (error) reject(error)
310 else resolve()
311 }
312 timer = after(timeoutMs, () => finish(new Error(`codex bridge startup timed out after ${timeoutMs}ms; exit code=unknown; stderr tail: ${stderrTail || '(empty)'}`)))
313 void (async () => {
314 try {
315 for await (const { stream, text } of events) {
316 if (stream === 'stderr') stderrTail = (stderrTail + text).slice(-2000)
317 if (stream === 'exit') status = `${text || 'exit'}${status === 'exit code=unknown' ? '' : `; ${status}`}`
318 if (stream === 'stdout' && text.includes('"ready":true')) finish()
319 }
320 finish(new Error(`codex bridge exited before ready; ${status}; stderr tail: ${stderrTail || '(empty)'}`))
321 } catch (err) {
322 finish(new Error(`codex bridge failed before ready; ${processStatus(err)}; ${String(err)}; stderr tail: ${stderrTail || '(empty)'}`))
323 }
324 })()
325 return promise
326}
327
328// cyrb53: the fingerprint only has to tell one bridge source from another.
329export function contentHash(text: string): string {
330 let h1 = 0xdeadbeef, h2 = 0x41c6ce57
331 for (let i = 0; i < text.length; i++) {
332 const c = text.charCodeAt(i)
333 h1 = Math.imul(h1 ^ c, 2654435761)
334 h2 = Math.imul(h2 ^ c, 1597334677)
335 }
336 h1 = Math.imul(h1 ^ (h1 >>> 16), 2246822507) ^ Math.imul(h2 ^ (h2 >>> 13), 3266489909)
337 h2 = Math.imul(h2 ^ (h2 >>> 16), 2246822507) ^ Math.imul(h1 ^ (h1 >>> 13), 3266489909)
338 return (h2 >>> 0).toString(16).padStart(8, '0') + (h1 >>> 0).toString(16).padStart(8, '0')
339}
340
341// The bridge's source as this module sees it on disk. A bridge outlives a
342// module reload, so a bridge whose /health reports another fingerprint runs
343// code this module was not written against.
344export const bridgeSource = ($: EngineInterface) => `${$.plugin.root}/daemon/codex-bridge.mjs`
345export const UNREAD_FINGERPRINT = 'unread'
346export async function bridgeFingerprint($: EngineInterface): Promise<string> {
347 try {
348 return contentHash(await $.fs.read(bridgeSource($)))
349 } catch (err) {
350 // An unreadable source must not keep the bridge from starting; it only skips the staleness check.
351 $.ui.log(`[codex-subagent] cannot fingerprint ${bridgeSource($)}: ${String(err)}`, { to: 'debug' })
352 return UNREAD_FINGERPRINT
353 }
354}
355
356// A bridge for the session's life; resolves once it listens. A step calls this
357// again when nothing listens on the socket (the bridge or codex died), or when
358// the bridge runs stale code; the new bridge takes over the socket (the old one
359// exits once its socket is replaced) and resumes threads from the threadId the step sends.
360function startBridge($: EngineInterface, sock: string, onFail: (ready: Promise<void>) => void): Promise<void> {
361 let stop: (() => void) | undefined
362 async function* output() {
363 const bridge = await $.process.spawn({ argv: ['node', bridgeSource($), sock, await bridgeFingerprint($)] })
364 stop = () => {
365 try { (bridge as unknown as { kill?: (signal?: string) => void }).kill?.('SIGTERM') } catch (err) { $.ui.log(`[codex-bridge] failed to stop after startup failure: ${String(err)}`, { to: 'debug' }) }
366 }
367 for await (const chunk of bridge) {
368 const { stream, text } = chunk
369 // $.ui.log drops any text over 4096 characters, so long bridge lines go in pieces.
370 for (const line of text.trimEnd().split('\n'))
371 for (let i = 0; i < line.length; i += 3900) $.ui.log(`[codex-bridge ${stream}] ${line.slice(i, i + 3900)}`, { to: 'debug' })
372 yield chunk
373 }
374 const status = bridge as unknown as { exitCode?: unknown; code?: unknown; signal?: unknown }
375 yield { stream: 'exit', text: processStatus(status), exitCode: status.exitCode ?? status.code, signal: status.signal }
376 $.ui.log('[codex-bridge] exited', { to: 'debug' })
377 }
378 const ready = bridgeReady(output(), (ms, fn) => $.clock.after(ms, fn))
379 void ready.catch(() => { stop?.(); onFail(ready) })
380 return ready
381}
382
383async function socketPath($: EngineInterface): Promise<string> {
384 const home = (await $.env.get('HOME')) ?? '/tmp'
385 return `${home}/.claude-work/run/codex-subagent-${await $.session.id()}.sock`
386}
387
388// What one load of this module keeps: a hot reload starts it over.
389type ModState = {
390 hookLog?: { path: string; text: string }
391 logChain: Promise<void>
392 registration?: Promise<void>
393 registrationLogged: boolean
394 // Each agent's definition (`codex-subagent:luna`, `general-purpose`, ...), by id.
395 subagentTypes: Map<string, string>
396}
397
398// The hook-side log beside the bridge's: `$.fs.write` replaces a whole file,
399// so the module keeps the text and writes it back, cut to its newest half
400// once it passes LOG_MAX.
401const LOG_MAX = 1 << 20
402function diag(st: ModState, $: EngineInterface, line: string): void {
403 try { $.ui.log(`[codex-subagent] ${line}`, { to: 'debug' }) } catch {}
404 st.logChain = st.logChain.then(async () => {
405 if (!st.hookLog) {
406 const path = (await socketPath($)).replace(/\.sock$/, '.hook.log')
407 st.hookLog = { path, text: await $.fs.read(path).catch(() => '') }
408 }
409 st.hookLog.text += `${new Date().toISOString()} ${line}\n`
410 if (st.hookLog.text.length > LOG_MAX) st.hookLog.text = st.hookLog.text.slice(-(LOG_MAX >> 1))
411 await $.fs.write(st.hookLog.path, st.hookLog.text)
412 }).catch(() => {})
413}
414
415// A hot reload runs register() again but not session.start, so the first
416// hook of a fresh module registers the types once more; a re-registered name
417// is replaced, so a TYPES model change reaches the running session.
418function ensureRegistered(st: ModState, $: EngineInterface): Promise<void> {
419 return st.registration ??= (async () => {
420 if ((await $.env.get('CODEX_SUBAGENT_HARVEST')) === '1') return
421 for (const [name, t] of Object.entries(TYPES)) {
422 await $.agent.register({
423 name,
424 description: `Runs the task on OpenAI ${t.model} at ${t.effort} effort in the ${t.sandbox} sandbox through Codex (codex's own shell and edit tools, plus this session's read-only MCP tools and skills relayed; ${t.blurb}). SendMessage continues the same codex thread.`,
425 prompt: 'This agent is served by Codex; this prompt is never sent to a Claude model.',
426 model: t.model, // never called: every step of this agent is answered by codex
427 })
428 }
429 })().catch(err => {
430 st.registration = undefined
431 if (!st.registrationLogged) { st.registrationLogged = true; diag(st, $, `agent type registration failed (retried on the next hook): ${String(err)}`) }
432 })
433}
434
435async function subagentTypeOf(st: ModState, $: EngineInterface, agents: ReadonlyMap<string, Agent>, agentId: string | undefined): Promise<string | undefined> {
436 if (!agentId) return undefined
437 const known = agents.get(agentId)
438 if (known) return `${PREFIX}${known.type}`
439 if (!st.subagentTypes.has(agentId)) {
440 try {
441 for (const a of await $.agent.list()) st.subagentTypes.set(a.id, a.type)
442 } catch (err) {
443 diag(st, $, `agent list failed while classifying ${agentId}: ${String(err)}`)
444 }
445 }
446 return st.subagentTypes.get(agentId)
447}
448
449export const register: Register = (on) => {
450 const st: ModState = { logChain: Promise.resolve(), registrationLogged: false, subagentTypes: new Map() }
451 const agents = new Map<string, Agent>()
452 // Messages delivered to a codex agent and not yet handed to codex, by agent id.
453 const inbox = new Map<string, string[]>()
454 // Codex spawns whose next(e) has not resolved yet.
455 const spawning = new Set<Promise<void>>()
456 // The first block index each codex step has not yielded yet: the turn.step
457 // .catch handler writes its error there, past the chunks the failed hook left.
458 const freeBlock = new Map<string, number>()
459 const stepKey = (e: { agentId?: string; turnId: string; index: number }) => `${e.agentId ?? ''}|${e.turnId}|${e.index}`
460 // Agents a step has claimed for codex: the .catch handler reads it without awaiting.
461 const claimed = new Set<string>()
462
463 on('session.receive', async ($, e, next) => {
464 const r = await next(e)
465 const text = r.text?.trim()
466 if (e.agentId && text && agents.has(e.agentId)) inbox.set(e.agentId, [...(inbox.get(e.agentId) ?? []), text])
467 return r
468 })
469 let sock = ''
470 let bridgeReady: Promise<void> | undefined
471 const dropFailedBridge = (failed: Promise<void>) => { if (bridgeReady === failed) bridgeReady = undefined }
472 // The bridgeReady whose /health matched this module's bridge source.
473 let verifiedBridge: Promise<void> | undefined
474
475 on('session.start', async ($, e, next) => {
476 const started = await next(e)
477 // The bridge's schema-harvest child is a `claude -p` that never reaches the
478 // API; were this plugin installed there, it would start a bridge of its own.
479 if ((await $.env.get('CODEX_SUBAGENT_HARVEST')) === '1') return started
480 sock = await socketPath($)
481 await ensureRegistered(st, $)
482 bridgeReady = startBridge($, sock, dropFailedBridge)
483 return started
484 })
485
486 on('agent.spawn', async ($, e, next) => {
487 void ensureRegistered(st, $)
488 const type = typeOf(e.subagentType)
489 if (!type) return next(e)
490 // A background agent's loop starts before next(e) resolves, so its first
491 // step may already be running: the cwd is resolved first and the agent
492 // registered the moment its id exists. turn.step waits on `spawning`.
493 let done!: () => void
494 const spawned = new Promise<void>(resolve => { done = resolve })
495 spawning.add(spawned)
496 let r: Awaited<ReturnType<typeof next>>
497 let cwd: string
498 try {
499 cwd = await workspaceRoot((argv, init) => $.process.run(argv, init), e.cwd ?? (await $.session.cwd()))
500 r = await next(e)
501 if (r.agentId) { agents.set(r.agentId, { type, cwd }); st.subagentTypes.set(r.agentId, e.subagentType) }
502 } finally {
503 spawning.delete(spawned)
504 done()
505 }
506 if (r.agentId) {
507 try {
508 sock ||= await socketPath($)
509 bridgeReady ??= startBridge($, sock, dropFailedBridge)
510 const ready = bridgeReady
511 await abortable(ready, next.signal)
512 const res = await $.http.fetch('http://codex-bridge/remember', {
513 method: 'POST', socketPath: sock, headers: { 'content-type': 'application/json' },
514 body: JSON.stringify({ key: r.agentId, type, cwd }),
515 })
516 if (!res.ok) throw new Error(`bridge POST /remember ${res.status}: ${res.text.slice(0, 2000)}`)
517 } catch (err) {
518 diag(st, $, `could not persist ${r.agentId}: ${err instanceof Error ? err.message : String(err)}`)
519 }
520 }
521 return r
522 })
523
524 on('turn.step', async function* ($, e, next) {
525 const agentId = e.agentId
526 const key = stepKey(e)
527 freeBlock.delete(key)
528 for (const k of freeBlock.keys()) { if (freeBlock.size < 64) break; freeBlock.delete(k) }
529 const used = (index: number) => { freeBlock.set(key, Math.max(freeBlock.get(key) ?? 0, index + 1)) }
530 const finish = (text: string, index = 0) => { used(index); return endTurn(e, text, index) }
531 // `work` is no `$` call, so waiting on it is charged to this hook's budget:
532 // once what the budget spares has passed, `onExpiry` settles it instead.
533 const withinBudget = <T>(work: Promise<T>, capMs: number, onExpiry: () => T): Promise<T> => {
534 let timer: Timer | undefined
535 const expired = new Promise<T>((resolve, reject) => {
536 timer = $.clock.after(Math.max(0, Math.min(capMs, next.budget.remainingMs - BUDGET_RESERVE_MS)), () => {
537 try { resolve(onExpiry()) } catch (err) { reject(err) }
538 })
539 })
540 return Promise.race([work, expired]).finally(() => timer?.cancel())
541 }
542 const bridgeUp = (ready: Promise<void>) => withinBudget(abortable(ready, next.signal), Infinity, () => {
543 throw new Error(`the codex bridge is still starting after this step's time budget (${next.budget.ms} ms)`)
544 })
545 const post = (route: string, payload: unknown) =>
546 $.http.fetch(`http://codex-bridge${route}`, { method: 'POST', socketPath: sock, headers: { 'content-type': 'application/json' }, body: JSON.stringify(payload) })
547 void ensureRegistered(st, $)
548 let agent = agentId ? agents.get(agentId) : undefined
549 if (!agent && agentId && spawning.size) {
550 // This may be the first step of a codex spawn still inside next(e); a
551 // spawn that never resolves must not hold the step, so the wait is bounded.
552 await withinBudget(Promise.all(spawning).then(() => {}), SPAWN_WAIT_MS, () => {})
553 agent = agents.get(agentId)
554 }
555 const subagentType = await subagentTypeOf(st, $, agents, agentId)
556 const claim: Claim = agent ? { ours: true, type: agent.type } : codexClaim(e.model, subagentType)
557 if (claim.ours && agentId) claimed.add(agentId)
558 const cause = `model=${e.model} type=${subagentType ?? '(unknown)'}`
559 let rebuilt = false
560 if (!agent && agentId) {
561 let found = false
562 try {
563 sock ||= await socketPath($)
564 bridgeReady ??= startBridge($, sock, dropFailedBridge)
565 const ready = bridgeReady
566 await bridgeUp(ready)
567 const res = await abortable(post('/recover', { key: agentId }), next.signal)
568 if (!res.ok) throw new Error(`bridge POST /recover ${res.status}: ${res.text.slice(0, 2000)}`)
569 const recovered = JSON.parse(res.text) as { found?: boolean; type?: TypeName; cwd?: string; threadId?: string }
570 found = recovered.found === true
571 if (found && recovered.type && recovered.cwd)
572 agent = { type: recovered.type, cwd: recovered.cwd, threadId: recovered.threadId }
573 } catch (err) {
574 diag(st, $, `recovery lookup failed for ${agentId} (${cause}): ${String(err)}`)
575 if (claim.ours) return yield* finish(`[codex-subagent error] cannot reach the bridge to recover ${agentId}: ${err instanceof Error ? err.message : String(err)}; send the message again`)
576 }
577 // Neither this module nor the bridge has seen the spawn: the agent's
578 // type (or its model) names the codex type, and the session's repo root
579 // is the cwd the spawn would pick.
580 if (!agent && !found && claim.ours) {
581 if (!claim.type) {
582 diag(st, $, `${agentId} is a codex step no type serves (${cause})`)
583 return yield* finish(`[codex-subagent error] cannot serve ${agentId} on codex: neither its agent type nor its model names a codex-subagent type (${cause}); this step was not sent to Claude`)
584 }
585 try {
586 agent = { type: claim.type, cwd: await workspaceRoot((argv, init) => $.process.run(argv, init), await $.session.cwd()) }
587 rebuilt = true
588 diag(st, $, `${agentId} unknown to spawn and bridge; serving it as ${claim.type} in ${agent.cwd} (${cause})`)
589 } catch (err) {
590 return yield* finish(`[codex-subagent error] cannot recover ${agentId ??'(unknown agent)'}: ${err instanceof Error ? err.message : String(err)}`)
591 }
592 }
593 if (agent && agentId) agents.set(agentId, agent)
594 }
595 if (!agent || !agentId) {
596 if (claim.ours) {
597 diag(st, $, `ending an unservable codex step for ${agentId ?? '(no agent id)'} (${cause})`)
598 return yield* finish(agentId ? `[codex-subagent error] cannot recover ${agentId} (${cause})` : `[codex-subagent error] ${e.model} step carries no agent id`)
599 }
600 return yield* next(e)
601 }
602 if (next.signal.aborted) return yield* finish('[codex-subagent error] interrupted')
603 const onAbort = () => { void post('/interrupt', { key: agentId }).catch(err => diag(st, $, `interrupt failed: ${String(err)}`)) }
604 next.signal.addEventListener('abort', onAbort)
605 // Block 0 holds the keepalive's thinking once it beats; the answer then streams in block 1.
606 let block = 0
607 const ask = async (route: string, payload: unknown) => {
608 const res = await abortable(post(route, payload), next.signal)
609 if (!res.ok) throw new Error(`bridge POST ${route} ${res.status}: ${res.text.slice(0, 2000)}`)
610 const r = JSON.parse(res.text) as StepReply
611 if ('error' in r) throw new Error(`bridge: ${r.error}`)
612 return r
613 }
614
615 // The bridge routes /health for GET only; every other route is a POST.
616 const getHealth = () => abortable($.http.fetch('http://codex-bridge/health', { method: 'GET', socketPath: sock }), next.signal)
617
618 const settle = (reply: StepReply): Promise<StepResult> => settlePending(
619 reply,
620 () => ask('/wait', { key: agentId }),
621 err => diag(st, $, `/wait failed (${String(err)}); retrying`),
622 async () => {
623 const res = await getHealth()
624 if (!res.ok) throw new Error(`bridge GET /health ${res.status}: ${res.text.slice(0, 2000)}`)
625 }
626 )
627 // True once the listening bridge runs this module's bridge source. A stale
628 // bridge is replaced only while no other agent has a turn in it and this
629 // step holds none of its tool calls; otherwise the next step checks again.
630 const bridgeCurrent = async (body: Record<string, unknown>): Promise<boolean> => {
631 const ready = bridgeReady
632 if (!ready || verifiedBridge === ready) return true
633 const want = await bridgeFingerprint($)
634 if (want === UNREAD_FINGERPRINT) return true
635 const res = await getHealth()
636 if (!res.ok) {
637 diag(st, $, `bridge GET /health ${res.status}: ${res.text.slice(0, 2000)}; cannot check its code, stepping through it`)
638 return true
639 }
640 const health = JSON.parse(res.text) as { fingerprint?: string; busy?: string[]; threads?: [string, string | null][] }
641 if (health.fingerprint === want) { verifiedBridge = ready; return true }
642 // A bridge from before /health reported `busy` still lists its threads;
643 // any of those agents the engine shows running may be mid-turn in it.
644 const known = new Set((health.threads ?? []).map(([key]) => key))
645 const others = Array.isArray(health.busy)
646 ? health.busy.filter(key => key !== agentId)
647 : (await $.agent.list()).filter(a => a.id !== agentId && a.status === 'running' && known.has(a.id)).map(a => a.id)
648 const holdsCall = body.toolResult !== undefined
649 if (others.length || holdsCall) {
650 diag(st, $, `bridge runs stale code (fingerprint ${health.fingerprint ?? '(none)'}, want ${want}); restart deferred: ${holdsCall ? `${agentId} answers a held tool call` : `turns running for ${others.join(', ')}`}`)
651 return false
652 }
653 diag(st, $, `bridge runs stale code (fingerprint ${health.fingerprint ?? '(none)'}, want ${want}); replacing it`)
654 const fresh = bridgeReady = startBridge($, sock, dropFailedBridge)
655 await bridgeUp(fresh)
656 verifiedBridge = fresh
657 return true
658 }
659 const step = async (body: Record<string, unknown>): Promise<StepResult> => {
660 let submitted = false
661 const roundTrip = async () => {
662 const ready = bridgeReady ?? (bridgeReady = startBridge($, sock, dropFailedBridge))
663 await bridgeUp(ready)
664 if (!(await bridgeCurrent(body)) && body.toolResult && body.prompt !== undefined) {
665 // A stale bridge may ignore `prompt` beside `toolResult`; the message
666 // rides in the tool's output instead, the one place it reaches codex.
667 const result = body.toolResult as { contentItems: ContentItem[] }
668 body.toolResult = { ...result, contentItems: [...result.contentItems, { type: 'inputText', text: body.prompt as string }] }
669 delete body.prompt
670 }
671 const reply = await ask('/step', body)
672 submitted = true
673 return settle(reply)
674 }
675 try {
676 return await roundTrip()
677 } catch (err) {
678 if (!unreachable(err) || submitted) throw err
679 diag(st, $, `bridge unreachable (${String(err)}); restarting it`)
680 const fresh = bridgeReady = startBridge($, sock, dropFailedBridge)
681 await bridgeUp(fresh)
682 return roundTrip()
683 }
684 }
685
686 // A failure from here on reaches the .catch handler below, which ends the
687 // turn as text: a failed hook would otherwise fall through to a Claude request.
688 try {
689 const api = await $.session.messages({ agentId, as: 'api' })
690 if (!Array.isArray(api)) throw new Error(`cannot read agent ${agentId}'s messages: ${api.deny}`)
691 const rows = await $.session.messages({ agentId })
692 const prompt = Array.isArray(rows) ? pendingPrompt(rows) : undefined
693 if (!agent.handback) {
694 const last = api[api.length - 1]
695 const handback = last?.role === 'assistant' ? (last.content as ApiBlock[]).find(b => b.type === 'tool_use' && b.name === HANDBACK && b.id && typeof b.input?.message === 'string') : undefined
696 if (handback?.id) agent.handback = { toolUseId: handback.id, text: handback.input!.message as string }
697 }
698 // A delivered SubagentHandback ends the run with no step after it, so the
699 // handback outlives it here; a user message after it is a SendMessage
700 // resume, which must start a codex turn rather than replay the old ending.
701 if (agent.handback && prompt !== undefined) agent.handback = undefined
702 if (agent.handback) {
703 // A refused or failed handback ends as plain text, which the harness then delivers.
704 const { toolUseId, text } = agent.handback
705 agent.handback = undefined
706 return yield* finish(toolResultFor(api, toolUseId)?.success ? `(report delivered through ${HANDBACK})` : text)
707 }
708 agent.relay ??= relaySet(await $.tool.list())
709 if (!agent.waiting) {
710 const last = api[api.length - 1]
711 const call = last?.role === 'assistant' ? (last.content as ApiBlock[]).find(b => b.type === 'tool_use' && b.id?.startsWith('toolu_codex_') && !b.id.startsWith('toolu_codex_handback_')) : undefined
712 if (call?.id) agent.waiting = { toolUseId: call.id, callId: call.id.slice('toolu_codex_'.length) }
713 }
714 const type = TYPES[agent.type]
715 const body: Record<string, unknown> = { key: agentId, agentType: agent.type, threadId: agent.threadId, cwd: agent.cwd, model: type.model, effort: type.effort, sandbox: type.sandbox, mcpTools: agent.relay.mcp }
716 const delivered = inbox.get(agentId) ?? []
717 const answered = agent.waiting && toolResultFor(api, agent.waiting.toolUseId, delivered)
718 if (agent.waiting && answered) {
719 const { messages, ...result } = answered
720 body.toolResult = { callId: agent.waiting.callId, ...result }
721 // The bridge answers the call first, then steers this text into the running turn.
722 if (messages) {
723 body.prompt = messages.join('\n\n')
724 const left = delivered.filter(d => !messages.some(m => m.includes(d)))
725 if (left.length) inbox.set(agentId, left)
726 else inbox.delete(agentId)
727 }
728 } else {
729 inbox.delete(agentId)
730 if (!prompt) throw new Error(`agent ${agentId} has neither a pending prompt nor the tool_result for ${agent.waiting?.toolUseId ?? '(none)'}`)
731 if (!agent.threadId) {
732 body.dynamicTools = agent.relay.tools
733 body.developerInstructions = `${PREAMBLE}\n\n${inheritedContext(api)}`
734 }
735 body.prompt = prompt
736 }
737 agent.waiting = undefined
738 sock ||= await socketPath($)
739 bridgeReady ??= startBridge($, sock, dropFailedBridge)
740 let reply = yield* keepalive(step(body), next.signal, () => { block = 1; used(0) }, (ms, fn) => $.clock.after(ms, fn))
741 agent.threadId = reply.threadId
742
743 for (;;) {
744 if (!('toolCall' in reply)) break
745 const { callId, tool, arguments: input } = reply.toolCall
746 const name = agent.relay?.claudeName.get(tool)
747 if (!name) {
748 const failed: Record<string, unknown> = { key: agentId, agentType: agent.type, threadId: agent.threadId, cwd: agent.cwd, model: type.model, effort: type.effort, sandbox: type.sandbox, toolResult: { callId, contentItems: [{ type: 'inputText', text: `tool not relayed: ${tool}` }], success: false } }
749 reply = yield* keepalive(step(failed), next.signal, () => { block = 1; used(0) }, (ms, fn) => $.clock.after(ms, fn))
750 agent.threadId = reply.threadId
751 continue
752 }
753 const toolUseId = `toolu_codex_${callId.replace(/[^A-Za-z0-9_-]/g, '_')}`
754 agent.waiting = { toolUseId, callId }
755 used(block)
756 const chunks: TurnStepChunk[] = [
757 { kind: 'tool', index: block, id: toolUseId, name },
758 { kind: 'input', index: block, json: JSON.stringify(input ?? {}) },
759 { kind: 'stop', stopReason: 'tool_use', usage: null },
760 ]
761 for (const c of chunks) yield c
762 return { turnId: e.turnId, index: e.index, answer: '', toolUses: [{ name, input }], stopReason: 'tool_use', usage: null } satisfies TurnStepResult
763 }
764
765 const f = 'final' in reply ? reply.final : { text: '', status: undefined, errors: ['bridge reply had neither toolCall nor final'] }
766 const errors = f.errors.length ? `\n\n[codex errors] ${f.errors.join('\n')}` : ''
767 const notice = rebuilt && api.some(m => m.role === 'assistant') ? '[codex-subagent: the earlier codex thread was not recoverable; this answer comes from a fresh thread]\n\n' : ''
768 const text = `${notice}${f.text}${errors}\n\n[codex-subagent: model=${reply.model ?? '?'} thread=${reply.threadId}${reply.resumed ? ' (resumed)' : ''} codexPid=${reply.codexPid} bridgePid=${reply.bridgePid} status=${f.status}]`
769 const toolUseId = `toolu_codex_handback_${reply.threadId.replace(/[^A-Za-z0-9_-]/g, '_')}_${e.turnId.replace(/[^A-Za-z0-9_-]/g, '_')}_${e.index}`
770 agent.handback = { toolUseId, text }
771 const input = { message: text }
772 used(block)
773 yield { kind: 'tool', index: block, id: toolUseId, name: HANDBACK }
774 yield { kind: 'input', index: block, json: JSON.stringify(input) }
775 yield { kind: 'stop', stopReason: 'tool_use', usage: null }
776 return { turnId: e.turnId, index: e.index, answer: '', toolUses: [{ name: HANDBACK, input }], stopReason: 'tool_use', usage: null } satisfies TurnStepResult
777 } finally {
778 next.signal.removeEventListener('abort', onAbort)
779 }
780 }).catch(async function* ($, e, next) {
781 // A codex step must never reach `next`: beneath it is the Claude request.
782 // The hook classifies the step before anything in it can fail, so this
783 // reads that verdict without awaiting; only a step that failed before then
784 // asks $.agent.list, bounded well inside the 1 s grace.
785 const agentId = e.agentId
786 const known = agentId !== undefined && (claimed.has(agentId) || agents.has(agentId))
787 let ours = known || codexClaim(e.model, agentId ? st.subagentTypes.get(agentId) : undefined).ours
788 if (!ours && agentId && !st.subagentTypes.has(agentId)) {
789 let timer: Timer | undefined
790 const listed = await Promise.race([
791 $.agent.list().catch(() => [] as { id: string; type: string }[]),
792 new Promise<{ id: string; type: string }[]>(resolve => { timer = $.clock.after(300, () => resolve([])) }),
793 ]).finally(() => timer?.cancel())
794 ours = codexClaim(e.model, listed.find(a => a.id === agentId)?.type).ours
795 }
796 if (!ours) return yield* next(e)
797 const key = stepKey(e)
798 const index = freeBlock.get(key) ?? 0
799 freeBlock.delete(key)
800 const { kind, message } = next.error
801 try { diag(st, $, `turn.step ${kind} for ${e.agentId ?? '(no agent id)'}: ${message ?? '(no message)'}`) } catch {}
802 return yield* endTurn(e, `[codex-subagent error] step hook ${kind}: ${message ?? '(no message)'}; send the message again`, index)
803 })
804}
805
806async function* endTurn(e: { turnId: string; index: number }, text: string, block = 0): AsyncGenerator<TurnStepChunk, TurnStepResult> {
807 yield { kind: 'text', index: block, text }
808 yield { kind: 'stop', stopReason: 'end_turn', usage: null }
809 return { turnId: e.turnId, index: e.index, answer: text, toolUses: [], stopReason: 'end_turn', usage: null }
810}
811