SLOPSHOPPER

codex-subagent

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…

newprocessnetworktimeragents
A shopper browsing a rack in a slop shop
README

.dotfiles

Installation

Windows CMD

Requirements:

  • Windows CMD
  • curl executable
  • Winget executable (>= v1.6.2631, follow instruction from winget repository if you don't have or have a lower version)
curl -L dotfiles.ranolp.dev/setup | cmd /Q

Windows PowerShell (TODO)

Requirements:

  • Windows PowerShell
  • Winget executable (>= v1.6.2631, follow instruction from winget repository if you don't have or have a lower version)
curl dotfiles.ranolp.dev/setup | iex

Windows ArchWSL (WIP)

Requirements:

  • ArchWSL Bash
  • curl executable
curl -L dotfiles.ranolp.dev/setup | sh

macOS (TODO)

Requirements:

  • macOS zsh
  • curl executable
curl -L dotfiles.ranolp.dev/setup | sh
Source 1 files
hooks/register.ts 811 lines
1import 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