Official
codex
Codex agent provider: runs the Codex CLI's app-server for each account.
The app opens the listing; nothing installs until an agent in your Plugins workspace has read the files and you enable the plugin. In a terminal: cvg install convergence/codex@0.2.0
Permissions in 0.2.0
Files
agent.ts82.1 KB
// The Codex provider: one `codex app-server` process per account, one
// Codex thread per Convergence session (a session id is a thread id).
// See NOTES.md for the protocol facts behind every choice here.
import * as Effect from "effect/Effect";
import * as Deferred from "effect/Deferred";
import * as Fiber from "effect/Fiber";
import * as Data from "effect/Data";
import * as Queue from "effect/Queue";
import * as Result from "effect/Result";
import * as Scope from "effect/Scope";
import * as Stream from "effect/Stream";
import * as z from "zod";
import { Host, Process, Kernel, parse } from "convergence/effect";
import type { PluginServices, EffectAgent } from "convergence/effect";
import { RpcTransport, request, shutdown, runProcess, callHostToolEffect } from "../sdk/effect.ts";
import type { TransportError } from "../sdk/effect.ts";
import type { RpcProcess, ProcessExit } from "../sdk/process.ts";
import type { JsonRpcConnection, RpcRequest, Responder, RequestOptions } from "../sdk/jsonrpc.ts";
import { errorMessage } from "../sdk/errors.ts";
import * as wire from "./wire.ts";
import type { Instance } from "./instances.ts";
import type {
Child,
Task,
Session,
Facts,
Event,
History,
ToolCall,
ToolContent,
Usage,
QuestionField,
QuestionAnswer,
SessionOptions,
PromptInput,
UsageLimits,
AgentInfo,
} from "./types.ts";
import { sessionSchema, answerSchema, promptSchema } from "./types.ts";
export class ProviderError extends Data.TaggedError("ProviderError")<{ readonly message: string }> {}
interface Pending<A> {
session: string;
request: string;
answer: (value: A | null) => void;
}
type Maintenance = NonNullable<AgentInfo["maintenance"]>;
import { methodNotFound, RpcError } from "../sdk/jsonrpc.ts";
import { newId, compact, outcome } from "../sdk/agent.ts";
import { installerOf, npmPrefix, isNewer, latestNpmEffect, latestBrewEffect } from "../sdk/maintenance.ts";
import * as map from "./map.ts";
import {
OPTION_MODEL,
OPTION_REASONING,
buildOptions,
startParams,
resumeParams,
forkParams,
turnParams,
userInput,
} from "./params.ts";
import { agentId, displayName, homeOf } from "./instances.ts";
import { ICON } from "./icon.ts";
export const FAMILY = "codex";
/// Reported to the app-server as the client's version.
export const CLIENT_VERSION = "0.2.0";
/// The npm package that publishes the CLI.
const NPM_PACKAGE = "@openai/codex";
/// Turns read per `thread/turns/list` call, and the point at which a very
/// long thread stops being loaded.
const TURN_PAGE = 100;
const MAX_TURNS = 2000;
/// Subagent threads read for one `read_session`, at every depth together.
const MAX_SUBAGENT_THREADS = 64;
/// How long a Stop waits on each subagent's interrupt, and on all of them.
const CHILD_INTERRUPT = 3000;
const CHILDREN_INTERRUPT = 10_000;
const APPROVAL_OPTIONS: import("convergence/protocol").ApprovalOption[] = [
{ id: "accept", name: "Allow once", kind: "allow_once" },
{ id: "acceptForSession", name: "Allow for this session", kind: "allow_always" },
{ id: "decline", name: "Reject", kind: "reject_once" },
{ id: "cancel", name: "Reject and stop", kind: "reject_always" },
];
/// The Codex version out of the app-server's user agent, which reads
/// `<client>/<codex version> (<platform>) …`.
export function versionOf(userAgent: unknown) {
if (typeof userAgent !== "string") return null;
const slash = userAgent.indexOf("/");
if (slash < 0) return null;
const version =
userAgent
.slice(slash + 1)
.trim()
.split(/\s+/)[0] ?? "";
return version || null;
}
function now() {
return new Date().toISOString();
}
// These native JSON-RPC validation/precondition replies prove no admission.
// INTERNAL_ERROR is also used for locally lost writes/timeouts: never replay it.
function nativeRejection(error: unknown) {
const wrapped = z.object({ cause: z.unknown() }).safeParse(error);
const cause = wrapped.success ? wrapped.data.cause : null;
return cause instanceof RpcError && [-32600, -32601, -32602].includes(cause.code);
}
export function newSession(workspace = "", options: Record<string, unknown> = {}): Session {
return {
workspace,
options: { ...options },
run: null,
turn: null,
inputs: new Map(),
endedTurns: new Set(),
// The host's id of each user message, with the turn it started. A
// rollback names the host's id; Codex only knows its turns.
turns: new Map(),
// Items that received a delta, so the completed item does not repeat
// the streamed text.
streamed: new Set(),
startedTools: new Set(),
// Reasoning held back while it could still be only the stdin notice.
held: new Map(),
};
}
/// Moves a subagent's task to `status`, stamping the end of its work.
/// Returns whether anything changed.
function setStatus(child: Child, status: import("convergence/protocol").TaskStatus) {
if (child.task.status === status) return false;
child.task.status = status;
if (map.isLive(status)) {
delete child.task.endedAt;
} else {
delete child.task.activity;
child.task.endedAt = now();
}
return true;
}
/// The thread that started a subagent.
function parentOf(child: Child) {
return child.task.parentTaskId ?? child.root;
}
/// What a subagent's own thread record says about it.
function factsOfThread(thread: wire.Thread | null | undefined): Facts {
const created =
typeof thread?.createdAt === "number" && thread.createdAt > 0 ? map.timestamp(thread.createdAt) : null;
return {
prompt: typeof thread?.preview === "string" && thread.preview.trim() ? thread.preview : null,
taskName: map.threadTaskName(thread),
name: map.threadNickname(thread) ?? map.threadRole(thread),
model: typeof thread?.model === "string" ? thread.model : null,
effort: typeof thread?.reasoningEffort === "string" ? thread.reasoningEffort : null,
startedAt: created,
};
}
/// Adds what is known about a subagent to its snapshot. Every field fills
/// only what the snapshot does not have yet, except the name, which a
/// later `thread/started` knows better than a spawn path.
function applyFacts(facts: Facts, child: Child) {
const task = child.task;
const fill = (key: "toolCallId" | "prompt" | "model" | "effort", value: string | null | undefined) => {
if ((task[key] === undefined || task[key] === null) && typeof value === "string" && value) task[key] = value;
};
fill("toolCallId", facts.toolCallId);
fill("prompt", facts.prompt);
fill("model", facts.model);
fill("effort", facts.effort);
if (!child.taskName && typeof facts.taskName === "string" && facts.taskName) child.taskName = facts.taskName;
if (typeof facts.name === "string" && facts.name) task.name = facts.name;
if (facts.startedAt) {
task.startedAt =
task.startedAt && Date.parse(task.startedAt) <= Date.parse(facts.startedAt) ? task.startedAt : facts.startedAt;
}
task.title = map.taskTitle(task.prompt ?? null, child.taskName, task.name ?? null);
}
/// Whether reasoning text could still turn out to be the stdin notice.
function couldBeStdinNotice(text: string) {
return map.STDIN_NOTICE.startsWith(text.trimStart()) || map.isStdinNotice(text);
}
export class CodexAgent {
instance: Instance;
transport: RpcTransport["Service"];
scope: Scope.Scope;
env: Record<string, string>;
emit: (event: Event) => void;
id: string;
name: string;
proc: RpcProcess | null;
connecting: Fiber.Fiber<RpcProcess, ProviderError | TransportError> | null;
userAgent: string | null;
models: wire.Model[] | null;
sessions: Map<string, Session>;
approvals: Map<string, Pending<string>>;
/// Each file change item's diffs while it runs: Codex's approval names
/// the item only, and the card shows what it changes.
fileChanges: Map<string, ToolContent[]>;
questions: Map<string, Pending<QuestionAnswer>>;
skills: Map<string, string>;
subagents: Map<string, Child>;
accountDirty: boolean;
quotaFacts: import("convergence/protocol").UsageRecovery | null = null;
accountGeneration = 0;
jobs: Queue.Queue<Effect.Effect<unknown, unknown, PluginServices>>;
enqueue(job: Effect.Effect<unknown, unknown, PluginServices>) {
Queue.offerUnsafe(this.jobs, job);
}
background = Effect.fn("Codex.background")(function* (this: CodexAgent) {
yield* Stream.runForEach(Stream.fromQueue(this.jobs), (job) => job.pipe(Effect.forkIn(this.scope))).pipe(
Effect.forkIn(this.scope),
);
});
request = Effect.fn("Codex.request")(
<M extends keyof typeof wire.responses>(
connection: JsonRpcConnection,
method: M,
params?: unknown,
options: Omit<RequestOptions, "signal"> = {},
) => request<z.infer<(typeof wire.responses)[M]>>(connection, method, wire.responses[method], params, options),
);
latest = Effect.fn("Codex.latest")(function* (
this: CodexAgent,
manager: string,
name: string,
prefix: string | null,
) {
return yield* manager === "brew" ? latestBrewEffect(name) : latestNpmEffect(name, prefix);
});
/// `api`: the plugin API. `instance`: the account. `env`: the login
/// environment's `HOME`, `CODEX_HOME` and `CODEX_ARGS`. `emit`: sends an
/// agent event (set once the agent is registered).
constructor({
instance,
scope,
transport,
jobs,
env = {},
emit = () => {},
}: {
instance: Instance;
scope: Scope.Scope;
transport: RpcTransport["Service"];
jobs: CodexAgent["jobs"];
env?: Record<string, string>;
emit?: (event: Event) => void;
}) {
this.scope = scope;
this.transport = transport;
this.jobs = jobs;
this.instance = instance;
this.env = env;
this.emit = emit;
this.id = agentId(instance);
this.name = displayName(instance);
this.proc = null;
this.connecting = null;
this.userAgent = null;
this.models = null;
this.sessions = new Map();
this.approvals = new Map();
this.fileChanges = new Map();
this.questions = new Map();
// Skill name to the path Codex knows it by, filled by `list_skills`.
this.skills = new Map();
// Subagent thread id to its chat and task, at any depth.
this.subagents = new Map();
// The account changed, so a cached status must be read again.
this.accountDirty = false;
}
/// The handlers the host calls, by `agent/<method>` name.
definition(
before: Effect.Effect<Record<string, string>, never, PluginServices> = Effect.succeed(this.env),
): EffectAgent {
const call = <S extends z.ZodType, A, E>(
schema: S,
handler: (params: z.infer<S>) => Effect.Effect<A, E, PluginServices>,
) =>
Effect.fn("Codex.handler")(
function* (this: CodexAgent, params: unknown) {
this.env = yield* before;
return yield* handler(yield* parse("codex handler", schema, params));
}.bind(this),
);
const session = z.object({ sessionId: z.string() });
return {
id: this.id,
name: this.name,
initialize: () =>
before.pipe(
Effect.tap((env) =>
Effect.sync(() => {
this.env = env;
}),
),
Effect.andThen(this.initialize()),
),
list_options: call(z.object({ workspace: z.string() }), (params) => this.listOptions(params)),
list_commands: call(z.unknown(), () => this.listCommands()),
list_sessions: call(z.object({ workspace: z.string() }), (params) => this.listSessions(params)),
read_session: call(session, (params) => this.readSession(params)),
list_skills: call(z.object({ workspace: z.string() }), (params) => this.listSkills(params)),
create_session: call(sessionSchema, (params) => this.createSession(params)),
resume_session: call(sessionSchema.extend({ sessionId: z.string() }), (params) => this.resumeSession(params)),
close_session: call(session, (params) => this.closeSession(params)),
fork_session: call(sessionSchema.extend({ sessionId: z.string() }), (params) => this.forkSession(params)),
prompt: call(session.extend({ input: promptSchema }), (params) => this.prompt(params)),
cancel: call(session, (params) => this.cancel(params)),
cancel_task: call(session.extend({ taskId: z.string() }), (params) => this.cancelTask(params)),
set_option: call(session.extend({ optionId: z.string(), value: z.unknown() }), (params) =>
this.setOption(params),
),
respond_to_approval: call(z.object({ approvalId: z.string(), optionId: z.string() }), (params) =>
this.respondToApproval(params),
),
respond_to_question: call(z.object({ questionId: z.string(), answer: answerSchema }), (params) =>
this.respondToQuestion(params),
),
rollback: call(session.extend({ itemId: z.string() }), (params) => this.rollback(params)),
compact: call(session, (params) => this.compact(params)),
usage_limits: call(z.unknown(), () => this.usageLimits()),
update: call(z.unknown(), () => this.update()),
authenticate: call(z.unknown(), () => this.authenticate()),
logout: call(z.unknown(), () => this.logout()),
};
}
// --- connection -------------------------------------------------------------
home() {
return homeOf(this.instance, this.env);
}
/// The live app-server, started on first use and again when it died.
client = Effect.fn("Codex.client")(function* (
this: CodexAgent,
): Effect.fn.Return<RpcProcess, ProviderError | TransportError, PluginServices> {
if (this.proc?.alive) return this.proc;
let connecting = this.connecting;
if (!connecting) {
connecting = yield* this.connect().pipe(
Effect.ensuring(
Effect.sync(() => {
this.connecting = null;
}),
),
Effect.forkIn(this.scope),
);
this.connecting = connecting;
}
return yield* Fiber.join(connecting);
});
connect = Effect.fn("Codex.connect")(function* (this: CodexAgent) {
const transport = this.transport;
const launch = (this.env.CODEX_ARGS ?? "").split(/\s+/).filter(Boolean);
const args = ["app-server", ...launch, ...this.instance.args];
const home = this.home();
let proc: RpcProcess | undefined;
proc = yield* transport
.spawn("codex", args, {
name: "codex",
env: home ? { CODEX_HOME: home } : undefined,
onNotification: ({ method, params }) => this.onNotification(method, params),
onRequest: (request) => this.onRequest(request),
onExit: (status, tail) => this.onExit(proc, status, tail),
})
.pipe(
Effect.provideService(Scope.Scope, this.scope),
Effect.mapError(
(error) =>
new ProviderError({
message: `the codex CLI could not be started (it must be on the login PATH): ${errorMessage(error)}`,
}),
),
);
const initialized = yield* Effect.result(
this.request(proc.connection, "initialize", {
clientInfo: { name: "convergence", title: "Divergence", version: CLIENT_VERSION },
capabilities: { experimentalApi: true, requestAttestation: false },
}),
);
if (Result.isFailure(initialized)) {
const tail = proc.stderrTail().join("\n");
yield* shutdown(proc, 500).pipe(Effect.catchTag("TransportFailed", () => Effect.void));
return yield* new ProviderError({
message: `codex app-server rejected initialize: ${errorMessage(initialized.failure)}${tail ? `\n${tail}` : ""}`,
});
}
proc.connection.notify("initialized", null);
this.userAgent = initialized.success.userAgent ?? "";
this.proc = proc;
return proc;
});
/// The app-server ended: every subagent still at work and every run
/// still going fails, so nothing waits on a process that is gone.
onExit(proc: RpcProcess | undefined, status: ProcessExit | null, tail: string[]) {
if (proc && this.proc !== proc && this.proc !== null) return;
const detail = tail?.length ? `: ${tail.slice(-5).join(" | ")}` : "";
console.info(`codex app-server exited (${JSON.stringify(status)})${detail}`);
this.proc = null;
this.invalidateQuota();
for (const [thread, child] of this.subagents) {
if (!map.isLive(child.task.status)) continue;
this.updateChild(thread, (c) => {
c.turn = null;
c.task.summary ??= "The codex app-server exited";
return setStatus(c, "failed");
});
}
for (const [session, state] of this.sessions) {
// A thread is loaded in the app-server that started or resumed it;
// the next one knows none of them until they are resumed again.
state.unloaded = true;
this.finishRun(session, outcome.failed("the codex app-server exited"));
}
}
/// The connection with `sessionId`'s thread loaded in it: a thread that
/// was open in an app-server that exited (an update, a crash) is resumed
/// in the new one first, or `turn/start` fails with "thread not found".
loadedConnection = Effect.fn("Codex.loadedConnection")(function* (this: CodexAgent, sessionId: string) {
const { connection } = yield* this.client();
const state = this.sessions.get(sessionId);
if (state?.unloaded) {
yield* this.request(connection, "thread/resume", resumeParams(sessionId, state)).pipe(
Effect.mapError((error) => new ProviderError({ message: `thread/resume failed: ${errorMessage(error)}` })),
);
state.unloaded = false;
}
return connection;
});
/// Stops the app-server. The next call starts a fresh one.
shutdownProcess = Effect.fn("Codex.shutdownProcess")(function* (this: CodexAgent) {
if (this.proc) yield* shutdown(this.proc);
});
// --- events ---------------------------------------------------------------------
send(event: Event) {
this.emit(event);
}
runOf(session: string) {
return this.sessions.get(session)?.run ?? null;
}
/// Emits an event for a Codex thread. A subagent runs in its own thread;
/// its events belong to the chat that spawned it, tagged with the task,
/// so the transcript nests them instead of dropping them.
emitFor(thread: string, kind: import("convergence/protocol").AgentEventKind) {
const child = this.subagents.get(thread);
const session = child ? child.root : thread;
const event: { sessionId: string; runId?: string; taskId?: string } = { sessionId: session };
const runId = this.runOf(session);
if (runId) event.runId = runId;
if (child) event.taskId = thread;
this.send(Object.assign(event, kind));
}
/// The host session a thread belongs to: itself, or the chat of the
/// subagent it is.
rootOf(thread: string) {
return this.subagents.get(thread)?.root ?? thread;
}
/// Sends a subagent's snapshot to its chat. The host places it inside
/// `parentTaskId`, so the event itself is untagged.
emitTask(thread: string) {
const child = this.subagents.get(thread);
if (!child) return;
const event: { sessionId: string; runId?: string } = { sessionId: child.root };
const runId = this.runOf(child.root);
if (runId) event.runId = runId;
this.send(Object.assign(event, { event: "task" as const }, compact(child.task)));
}
/// Changes a known subagent and sends its snapshot when that changed
/// anything. `false` for a thread that is not a subagent.
updateChild(thread: string, change: (child: Child) => boolean) {
const child = this.subagents.get(thread);
if (!child) return false;
if (change(child)) this.emitTask(thread);
return true;
}
/// Registers a subagent `parent` started, or adds what is now known
/// about one already registered, and sends its snapshot. A grandchild
/// reports into the same chat, inside its parent's task.
registerChild(thread: string, parent: string, facts: Facts) {
let child = this.subagents.get(thread);
if (child) {
applyFacts(facts, child);
} else {
// `thread/started` reports every thread of the app-server; one whose
// parent is neither a chat here nor a known subagent is not ours.
const parentChild = this.subagents.get(parent);
let root;
let parentTask = null;
if (parentChild) {
root = parentChild.root;
parentTask = parent;
} else if (this.sessions.has(parent)) {
root = parent;
} else {
return;
}
const task: Task = { id: thread, title: "", status: "running", startedAt: now() };
if (parentTask) task.parentTaskId = parentTask;
child = {
root,
turn: null,
endedTurns: new Set(),
task,
taskName: null,
tools: new Set(),
streamed: new Set(),
held: new Map(),
};
applyFacts(facts, child);
this.subagents.set(thread, child);
}
this.emitTask(thread);
}
/// Registers a subagent a collab call addresses but this process never
/// saw start (it was spawned before a restart), so its events reach the
/// chat. Its names come from its own thread record.
adoptChild(thread: string, parent: string) {
// A chat is never a subagent, even when a child addresses it.
if (this.subagents.has(thread) || this.sessions.has(thread)) return;
this.registerChild(thread, parent, {});
this.readFacts(thread);
}
/// Fills a registered subagent's prompt and names from its own thread
/// record, in the background, on the running connection.
readFacts(thread: string) {
const proc = this.proc;
if (!proc?.alive) return;
this.enqueue(
this.request(proc.connection, "thread/read", { threadId: thread, includeTurns: false }).pipe(
Effect.tap((response) =>
Effect.sync(() => {
if (response.thread) this.registerChild(thread, "", factsOfThread(response.thread));
}),
),
Effect.catchTag(["TransportFailed", "ParseFailed"], (error) =>
Effect.sync(() => console.debug(`could not read the subagent thread ${thread}: ${errorMessage(error)}`)),
),
),
);
}
notice(thread: string, level: import("convergence/protocol").NoticeLevel, message: unknown) {
this.emitFor(thread, { event: "notice", level, message: String(message) });
}
/// Ends the active run of `session`. A second call does nothing, so a
/// run always ends with exactly one `run_finished`.
finishRun(session: string, result: import("../sdk/agent.ts").RunOutcome) {
const state = this.sessions.get(session);
if (!state) return;
if (state.turn) state.endedTurns.add(state.turn);
if (state.starting) Deferred.doneUnsafe(state.starting, Effect.succeed(null));
state.starting = undefined;
state.turn = null;
const runId = state.run;
state.run = null;
if (runId) this.send({ sessionId: session, runId, event: "run_finished", outcome: result });
}
/// The streaming state of a thread: its session's, or its subagent's.
streamState(thread: string) {
return this.subagents.get(thread) ?? this.sessions.get(thread) ?? null;
}
markStreamed(thread: string, item: string) {
this.streamState(thread)?.streamed.add(item);
}
wasStreamed(thread: string, item: string) {
return this.streamState(thread)?.streamed.has(item) ?? false;
}
/// A tool item already reported as started.
toolStarted(thread: string, item: string) {
const child = this.subagents.get(thread);
if (child) return child.tools.has(item);
return this.sessions.get(thread)?.startedTools.has(item) ?? false;
}
/// A reasoning delta. The first text of an item is held while it could
/// still be only `codex exec`'s stdin notice, which is not reasoning.
reasoningDelta(thread: string, item: string, delta: string) {
const state = this.streamState(thread);
const first = state ? !state.streamed.has(item) : false;
if (state) state.streamed.add(item);
const held = state?.held.get(item);
if (held !== undefined) {
const text = held + delta;
if (couldBeStdinNotice(text)) {
state?.held.set(item, text);
return;
}
state?.held.delete(item);
this.emitFor(thread, { event: "reasoning_delta", itemId: item, text, mode: "append" });
return;
}
if (first && delta && couldBeStdinNotice(delta)) {
state?.held.set(item, delta);
return;
}
this.emitFor(thread, { event: "reasoning_delta", itemId: item, text: delta, mode: "append" });
}
/// A completed reasoning item's held text: dropped when it is the stdin
/// notice, sent otherwise.
releaseHeld(thread: string, item: string) {
const state = this.streamState(thread);
const held = state?.held.get(item);
if (held === undefined) return;
state?.held.delete(item);
if (!map.isStdinNotice(held))
this.emitFor(thread, { event: "reasoning_delta", itemId: item, text: held, mode: "append" });
}
// --- notifications -----------------------------------------------------------------
onNotification(method: string, params: unknown) {
try {
this.routeNotification(method, wire.notification.parse(params ?? {}));
} catch (error) {
console.debug(`could not read the codex notification ${method}: ${errorMessage(error)}`);
}
}
routeNotification(method: string, p: wire.Notification) {
const thread = typeof p.threadId === "string" ? p.threadId : null;
const stream = thread ? this.streamState(thread) : null;
// Correlated late input evidence belongs to its original run even after
// retirement. Fence known retired turns, not every differing ID: child
// activity and native turn/started are independent lifecycle evidence.
const nativeTurn = p.turnId ?? p.turn?.id;
const item = method === "item/started" || method === "item/completed" ? p.item : null;
if (thread && item) this.consumeInput(thread, item);
if (nativeTurn && stream?.endedTurns.has(nativeTurn)) {
// Child snapshots outlive the turn that spawned/addressed them; these
// bookkeeping items cannot finish the parent or revive its run.
if (thread && item?.type === "collabAgentToolCall") this.onCollabItem(thread, item);
if (thread && item?.type === "subAgentActivity") this.onActivity(thread, item);
return;
}
switch (method) {
case "item/agentMessage/delta":
if (thread === null || typeof p.itemId !== "string") return;
this.markStreamed(thread, p.itemId);
this.emitFor(thread, {
event: "text_delta",
itemId: p.itemId,
text: typeof p.delta === "string" ? p.delta : "",
mode: "append",
});
return;
case "item/reasoning/textDelta":
case "item/reasoning/summaryTextDelta":
if (thread === null || typeof p.itemId !== "string") return;
this.reasoningDelta(thread, p.itemId, typeof p.delta === "string" ? p.delta : "");
return;
case "item/started":
if (thread === null || !p.item) return;
this.onItemStarted(thread, p.item);
return;
case "item/completed":
if (thread === null || !p.item) return;
this.onItemCompleted(thread, p.item);
return;
case "item/commandExecution/outputDelta":
case "item/fileChange/outputDelta":
if (thread === null || typeof p.itemId !== "string") return;
this.emitFor(thread, {
event: "tool_call_updated",
id: p.itemId,
outputDelta: typeof p.delta === "string" ? p.delta : "",
});
return;
case "item/fileChange/patchUpdated": {
if (thread === null || typeof p.itemId !== "string") return;
const content = map.diffContent(Array.isArray(p.changes) ? p.changes : []);
this.fileChanges.set(p.itemId, content);
this.emitFor(thread, { event: "tool_call_updated", id: p.itemId, content });
return;
}
case "turn/started": {
if (thread === null || typeof p.turn?.id !== "string") return;
const turn = p.turn.id;
if (stream?.turn && stream.turn !== turn) stream.endedTurns.add(stream.turn);
const isChild = this.updateChild(thread, (child) => {
child.turn = turn;
child.streamed.clear();
child.held.clear();
return setStatus(child, "running");
});
if (!isChild) {
const state = this.sessions.get(thread);
if (state) {
if (!state.run) {
state.run = newId("codex-run");
this.emitFor(thread, { event: "run_started" });
}
state.turn = turn;
state.streamed.clear();
state.startedTools.clear();
state.held.clear();
if (state.starting) Deferred.doneUnsafe(state.starting, Effect.succeed(turn));
}
}
return;
}
case "turn/completed": {
if (thread === null || typeof p.turn?.id !== "string") return;
if (!["completed", "interrupted", "failed"].includes(p.turn.status ?? "")) return;
stream?.endedTurns.add(p.turn.id);
// A terminal event cannot settle a different live native turn.
if (stream?.turn && stream.turn !== p.turn.id) return;
// A subagent's turn is not the chat's run.
const status = map.turnTaskStatus(p.turn.status);
const message = typeof p.turn.error?.message === "string" && p.turn.error.message ? p.turn.error.message : null;
const isChild = this.updateChild(thread, (child) => {
child.turn = null;
let changed = false;
if (status === "failed" && message) {
child.task.summary = message;
changed = true;
}
// The parent may already have been told it is done for good; a
// finished turn does not reopen it.
if (status === "idle" && child.task.status === "completed") return changed;
return setStatus(child, status) || changed;
});
if (isChild) return;
switch (p.turn.status) {
case "completed":
this.finishRun(thread, outcome.completed());
return;
case "interrupted":
this.finishRun(thread, outcome.cancelled());
return;
case "failed":
this.quotaFailure(thread, p.turn.error);
this.finishRun(
thread,
outcome.failed(p.turn.error ? (p.turn.error.message ?? "") : "the codex turn failed"),
);
return;
default:
return;
}
}
case "turn/plan/updated":
if (thread === null) return;
this.emitFor(thread, { event: "plan", ...map.planFromSteps(p.plan) });
return;
case "thread/tokenUsage/updated": {
if (thread === null || !p.tokenUsage) return;
const tokens = (value: unknown) =>
typeof value === "number" && Number.isInteger(value) && value >= 0 ? value : 0;
// A subagent's usage is what it spent in total; the chat's is the
// size of its context window.
const total: Usage = { usedTokens: tokens(p.tokenUsage.total?.totalTokens) };
const isChild = this.updateChild(thread, (child) => {
child.task.usage = total;
return false;
});
let usage: Usage = total;
if (!isChild) {
usage = { usedTokens: tokens(p.tokenUsage.last?.totalTokens) };
if (Number.isInteger(p.tokenUsage.modelContextWindow))
usage.contextWindow = p.tokenUsage.modelContextWindow ?? undefined;
}
this.emitFor(thread, { event: "usage", ...usage });
return;
}
case "thread/started": {
const spawn = map.threadSpawn(p.thread);
// Only spawned subagents; review, compaction and guardian threads
// are Codex's own business.
if (spawn && typeof p.thread?.id === "string")
this.registerChild(p.thread.id, spawn.parentThreadId, factsOfThread(p.thread));
return;
}
case "thread/status/changed": {
if (thread === null || p.status?.type !== "active") return;
const status = map.threadWaiting(p.status) ? "waiting" : "running";
this.updateChild(thread, (child) => map.isLive(child.task.status) && setStatus(child, status));
return;
}
case "thread/name/updated":
if (thread === null) return;
this.emitFor(
thread,
compact({ event: "session_info", title: typeof p.threadName === "string" ? p.threadName : null }),
);
return;
case "error": {
if (thread === null) return;
const message = p.error
? typeof p.error.message === "string"
? p.error.message
: ""
: "codex reported an error";
this.notice(thread, "error", message);
this.quotaFailure(thread, p.error);
// Even a non-retry error is not the terminal turn notification.
// Native retries/grace/tool settlement retain ownership until completed.
return;
}
case "warning":
if (thread !== null) this.notice(thread, "warning", typeof p.message === "string" ? p.message : "");
return;
// Diagnostics Codex sends out of band. Dropping them leaves the user
// with no explanation when the agent misbehaves.
case "deprecationNotice":
case "configWarning":
case "guardianWarning":
case "windows/worldWritableWarning":
this.notice(thread ?? "", "warning", typeof p.message === "string" ? p.message : `codex reported ${method}`);
return;
case "model/rerouted": {
const message =
typeof p.fromModel === "string" && typeof p.toModel === "string"
? `Codex switched from ${p.fromModel} to ${p.toModel}.`
: "Codex switched to a different model.";
this.notice(thread ?? "", "info", message);
return;
}
// A summary part is a section break in the reasoning; without it the
// parts run together as one paragraph.
case "item/reasoning/summaryPartAdded":
if (thread === null || typeof p.itemId !== "string") return;
if (this.wasStreamed(thread, p.itemId)) this.reasoningDelta(thread, p.itemId, "\n\n");
return;
case "item/mcpToolCall/progress":
if (thread === null || typeof p.itemId !== "string") return;
if (typeof p.message === "string" && p.message) {
this.emitFor(thread, { event: "tool_call_updated", id: p.itemId, outputDelta: `${p.message}\n` });
}
return;
case "account/rateLimits/updated": {
const first = this.sessions.keys().next();
const recovery = map.usageRecovery(p, "codex/account/rateLimits/updated", now(), this.quotaFacts?.identity);
this.quotaFacts = recovery;
const limits = map.usageLimits(p, recovery);
this.emitFor(first.done ? "" : first.value, { event: "usage_limits", ...limits });
return;
}
case "account/updated":
this.accountDirty = true;
this.invalidateQuota();
return;
// Another Codex client answered a request we are also showing.
// Without this the card waits for an answer that has been given.
case "serverRequest/resolved":
if (p.requestId === undefined || p.requestId === null) return;
this.withdrawRequest(JSON.stringify(p.requestId));
return;
default:
return;
}
}
/// Closes the approval or question `request` asked, telling the UI to
/// drop the card, and leaves the app-server unanswered by us.
withdrawRequest(request: string) {
for (const [id, entry] of this.approvals) {
if (entry.request !== request) continue;
this.approvals.delete(id);
entry.answer(null);
this.emitFor(entry.session, { event: "approval_resolved", id });
return;
}
for (const [id, entry] of this.questions) {
if (entry.request !== request) continue;
this.questions.delete(id);
entry.answer(null);
this.emitFor(entry.session, { event: "question_resolved", id });
return;
}
}
/// A collab tool call (multi-agent v1) becomes subagent tasks.
/// Registering the child thread lets `emitFor` route the child's own
/// events into the chat instead of dropping them.
onCollabItem(thread: string, item: wire.Item) {
const sender = typeof item.senderThreadId === "string" && item.senderThreadId ? item.senderThreadId : thread;
const receivers = (Array.isArray(item.receiverThreadIds) ? item.receiverThreadIds : []).filter(
(id) => typeof id === "string",
);
const states = item.agentsStates && typeof item.agentsStates === "object" ? item.agentsStates : {};
if (item.tool === "spawnAgent") {
for (const child of receivers) {
this.registerChild(child, sender, {
toolCallId: typeof item.id === "string" ? item.id : null,
prompt: typeof item.prompt === "string" ? item.prompt : null,
model: typeof item.model === "string" ? item.model : null,
effort: typeof item.reasoningEffort === "string" ? item.reasoningEffort : null,
});
}
} else if (item.tool !== "listAgents") {
// A call to a subagent this process never saw start: one from
// before a restart that the agent talks to again.
for (const child of [...receivers, ...Object.keys(states)]) this.adoptChild(child, sender);
}
for (const [child, state] of Object.entries(states)) {
const status = map.collabStatus(state?.status);
const message = typeof state?.message === "string" && state.message.trim() ? state.message : null;
this.updateChild(child, (c) => {
let changed = false;
// The child's answer is its result, never its title.
if (message !== null && c.task.summary !== message) {
c.task.summary = message;
changed = true;
}
// `pendingInit` and `running` are what the child's own turns
// report better.
return status ? setStatus(c, status) || changed : changed;
});
}
}
/// A multi-agent v2 subagent event. A v2 spawn appears only as a
/// `started` activity.
onActivity(thread: string, item: wire.Item) {
const target = typeof item.agentThreadId === "string" ? item.agentThreadId : "";
if (!target) return;
const path = typeof item.agentPath === "string" ? item.agentPath : "";
if (item.kind === "started") {
// The path's last segment is the task name its parent gave it; v2
// sends the instruction itself encrypted.
this.registerChild(target, thread, { taskName: map.pathTaskName(path) });
const child = this.subagents.get(target);
const described = !child || child.task.prompt || child.task.name;
if (!described) this.readFacts(target);
return;
}
// A child the chat started before a restart and addresses again. Only
// the chat's own activities adopt one: inside a child, an activity may
// name an ancestor (`/root` is the chat itself).
if (path !== "/root" && this.sessions.has(thread)) this.adoptChild(target, thread);
const status = map.activityStatus(item.kind);
if (!status) return;
// Only the thread that started the child speaks for it.
this.updateChild(target, (child) => parentOf(child) === thread && setStatus(child, status));
}
/// Notes a tool a subagent ran: its count and its live activity line.
childTool(thread: string, call: ToolCall, running: boolean) {
this.updateChild(thread, (child) => {
const added = !child.tools.has(call.id);
child.tools.add(call.id);
child.task.toolUses = child.tools.size;
const activity = running && call.title ? call.title : null;
const active = running && (child.task.activity ?? null) !== activity;
if (active) {
if (activity === null) delete child.task.activity;
else child.task.activity = activity;
}
return added || active;
});
}
consumeInput(thread: string, item: wire.Item) {
if (item.type !== "userMessage" || !item.clientId) return;
const state = this.sessions.get(thread);
const input = state?.inputs.get(item.clientId);
if (!input) return;
state?.inputs.delete(item.clientId);
this.send({
sessionId: thread,
runId: input.runId,
event: "input_consumed",
inputId: input.inputId,
nativeInputId: item.id ?? item.clientId,
});
}
quotaFailure(thread: string, error: z.infer<typeof wire.turnError> | null | undefined) {
if (error?.codexErrorInfo !== "usageLimitExceeded" || !this.sessions.get(thread)?.run) return;
const facts = this.quotaFacts;
this.emitFor(thread, {
event: "usage_blocked",
recovery: {
...facts,
availability: "blocked",
observedAt: now(),
source: "codex/usageLimitExceeded",
reason: error.message ?? "Codex reported usage quota exhaustion",
},
});
}
invalidateQuota() {
this.accountGeneration += 1;
this.quotaFacts = null;
}
onItemStarted(thread: string, item: wire.Item) {
switch (map.itemType(item)) {
case "collabAgentToolCall":
this.onCollabItem(thread, item);
return;
case "subAgentActivity":
this.onActivity(thread, item);
return;
case null:
case "userMessage":
case "agentMessage":
case "reasoning":
case "plan":
case "contextCompaction":
return;
default: {
const call = map.toolCallFromItem(item);
if (!call) return;
if (map.itemType(item) === "fileChange") this.fileChanges.set(call.id, call.content ?? []);
this.emitFor(thread, { event: "tool_call_started", ...call });
this.childTool(thread, call, true);
this.sessions.get(thread)?.startedTools.add(call.id);
}
}
}
onItemCompleted(thread: string, item: wire.Item) {
const question = map.backgroundQuestion(item);
if (question) this.emitFor(thread, { event: "question", ...question });
switch (map.itemType(item)) {
case "collabAgentToolCall":
this.onCollabItem(thread, item);
return;
case "subAgentActivity":
this.onActivity(thread, item);
return;
case null:
case "userMessage":
case "plan":
return;
case "agentMessage": {
const text = typeof item.text === "string" ? item.text : "";
// A subagent's latest message is its progress, and its last one is
// its result.
if (text.trim()) {
this.updateChild(thread, (child) => {
const changed = child.task.summary !== text;
child.task.summary = text;
return changed;
});
}
if (!this.wasStreamed(thread, item.id ?? "") && text) {
this.emitFor(thread, { event: "text_delta", itemId: item.id ?? "", text, mode: "replace" });
}
return;
}
case "reasoning": {
this.releaseHeld(thread, item.id ?? "");
const text = map.reasoningText(item.summary, item.content);
if (!this.wasStreamed(thread, item.id ?? "") && text.trim() && !map.isStdinNotice(text)) {
this.emitFor(thread, { event: "reasoning_delta", itemId: item.id ?? "", text, mode: "replace" });
}
return;
}
case "contextCompaction":
this.emitFor(thread, { event: "compacted" });
return;
default: {
const call = map.toolCallFromItem(item);
if (!call) return;
this.fileChanges.delete(call.id);
const started = this.toolStarted(thread, call.id);
this.childTool(thread, call, false);
if (started) {
this.emitFor(thread, {
event: "tool_call_updated",
id: call.id,
status: call.status,
title: call.title,
kind: call.kind,
input: call.input ?? null,
content: call.content ?? [],
locations: call.locations ?? [],
});
} else {
this.emitFor(thread, { event: "tool_call_started", ...call });
}
}
}
}
// --- server requests ------------------------------------------------------------------
onRequest({ id, method, params, reply }: RpcRequest) {
// Canonical JSON, so the id matches the one `serverRequest/resolved`
// reports for the same request.
const request = JSON.stringify(id);
const p = wire.notification.parse(params ?? {});
switch (method) {
case "item/commandExecution/requestApproval":
return this.commandApproval(p, reply, request);
case "item/fileChange/requestApproval":
return this.fileChangeApproval(p, reply, request);
case "item/tool/requestUserInput":
return this.userInputRequest(p, reply, request);
case "mcpServer/elicitation/request":
return this.elicitation(p, reply, request);
case "item/tool/call":
return this.toolCall(p, reply);
default:
// `item/permissions/requestApproval`,
// `account/chatgptAuthTokens/refresh`, …: answered at once so the
// app-server never waits on us.
reply.err(methodNotFound(method));
return undefined;
}
}
/// Runs one of the plugin tools given at thread start. The host runs
/// it; a tool can take minutes, and code mode's `execute` calls other
/// tools while it runs, so nothing here waits for it.
toolCall(p: wire.Notification, reply: Responder) {
if (typeof p.threadId !== "string" || typeof p.tool !== "string")
throw new Error("item/tool/call without a thread or a tool");
const session = this.rootOf(p.threadId);
this.enqueue(
callHostToolEffect({
agentId: this.id,
sessionId: session,
name: p.tool,
input: p.arguments ?? null,
callId: typeof p.callId === "string" ? p.callId : null,
}).pipe(Effect.tap((result) => Effect.sync(() => reply.ok(map.dynamicToolResponse(result))))),
);
}
commandApproval(p: wire.Notification, reply: Responder, request: string) {
if (typeof p.threadId !== "string" || typeof p.itemId !== "string")
throw new Error("an approval without a thread or an item");
const command = typeof p.command === "string" ? p.command : "";
const cwd = typeof p.cwd === "string" ? p.cwd : null;
const terminal: ToolContent = { type: "terminal", command, output: "" };
if (cwd !== null) terminal.cwd = cwd;
const call = map.toolCall({
id: p.itemId,
name: "shell",
kind: map.commandToolKind(p.commandActions),
title: command,
status: "pending",
input: { command, cwd },
content: [terminal],
});
// The card draws the command under its headline; saying it in the
// headline too would show it twice.
const title = typeof p.reason === "string" ? p.reason : "Run a command";
this.askApproval(p.threadId, title, call, reply, request);
}
fileChangeApproval(p: wire.Notification, reply: Responder, request: string) {
if (typeof p.threadId !== "string" || typeof p.itemId !== "string")
throw new Error("an approval without a thread or an item");
const title =
typeof p.reason === "string"
? p.reason
: typeof p.grantRoot === "string"
? `Allow writes under ${p.grantRoot}`
: "Apply file changes";
const content = this.fileChanges.get(p.itemId) ?? [];
const call = map.toolCall({
id: p.itemId,
name: "apply_patch",
kind: "edit",
title,
status: "pending",
content,
});
this.askApproval(p.threadId, title, call, reply, request);
}
/// Emits the approval and answers the app-server once the user decides.
/// Both decision enums share the four option ids.
askApproval(session: string, title: string, call: ToolCall, reply: Responder, request: string) {
const id = newId("codex-approval");
const entry = {
session,
request,
// `null` withdraws the request without answering it.
answer: (decision: string | null) => {
if (this.approvals.get(id) === entry) this.approvals.delete(id);
if (decision === null) reply.discard();
else reply.ok({ decision });
},
};
this.approvals.set(id, entry);
this.emitFor(session, { event: "approval", id, title, toolCall: call, options: APPROVAL_OPTIONS });
}
/// An MCP server asks through Codex: an MCP tool approval (Computer Use
/// and other MCP tools), an external page to open, or a form.
elicitation(p: wire.Notification, reply: Responder, request: string) {
if (typeof p.threadId !== "string") throw new Error("an elicitation without a thread");
const session = p.threadId;
const message = typeof p.message === "string" ? p.message : "";
const meta = p._meta && typeof p._meta === "object" && !Array.isArray(p._meta) ? p._meta : null;
const refuse = { action: "cancel", content: null, _meta: null };
if (meta?.codex_approval_kind === "mcp_tool_call") {
// The advertised persistence scopes become choices and return in the
// response metadata, never in the form content.
const supports = (scope: string) =>
Array.isArray(meta.persist) ? meta.persist.includes(scope) : meta.persist === scope;
const options: import("convergence/protocol").ApprovalOption[] = [
{ id: "accept", name: "Allow once", kind: "allow_once" },
];
if (supports("session"))
options.push({ id: "acceptForSession", name: "Allow for this session", kind: "allow_always" });
if (supports("always")) options.push({ id: "acceptAlways", name: "Always allow", kind: "allow_always" });
options.push({ id: "cancel", name: "Cancel", kind: "reject_always" });
const id = newId("codex-elicitation");
const entry = {
session,
request,
// `null` withdraws the request without answering it.
answer: (decision: string | null) => {
if (this.approvals.get(id) === entry) this.approvals.delete(id);
if (decision === null) return reply.discard();
if (decision === "accept") return reply.ok({ action: "accept", content: {}, _meta: null });
if (decision === "acceptForSession")
return reply.ok({ action: "accept", content: {}, _meta: { persist: "session" } });
if (decision === "acceptAlways")
return reply.ok({ action: "accept", content: {}, _meta: { persist: "always" } });
if (decision === "decline") return reply.ok({ action: "decline", content: null, _meta: null });
return reply.ok(refuse);
},
};
this.approvals.set(id, entry);
this.emitFor(session, { event: "approval", id, title: message, options });
return;
}
const mode = typeof p.mode === "string" ? p.mode : "";
if (mode === "url") {
const url = typeof p.url === "string" && /^https?:\/\//.test(p.url) ? p.url : null;
if (url === null) {
this.notice(session, "error", "Codex requested an invalid external URL");
reply.ok(refuse);
return;
}
this.askElicitation(session, request, reply, { message, url, fields: [] }, (answer) => ({
action: answer.cancelled ? "cancel" : "accept",
content: null,
_meta: null,
}));
return;
}
if (!["form", "openai/form", "openaiForm"].includes(mode)) {
reply.err(methodNotFound("mcpServer/elicitation/request"));
return;
}
const schema = p.requestedSchema ?? {};
this.askElicitation(session, request, reply, { message, fields: map.elicitationFields(schema) }, (answer) =>
answer.cancelled
? refuse
: { action: "accept", content: map.elicitationContent(schema, answer.values ?? {}), _meta: null },
);
}
/// Emits an elicitation as a question and replies with `respond(answer)`.
askElicitation(
session: string,
request: string,
reply: Responder,
question: { message: string; url?: string; fields: QuestionField[] },
respond: (answer: QuestionAnswer) => unknown,
) {
const id = newId("codex-elicitation");
const entry = {
session,
request,
// `null` withdraws the question without answering it.
answer: (answer: QuestionAnswer | null) => {
if (this.questions.get(id) === entry) this.questions.delete(id);
if (answer === null) reply.discard();
else reply.ok(respond(answer));
},
};
this.questions.set(id, entry);
this.emitFor(session, { event: "question", responseMode: "tool", id, ...question });
}
userInputRequest(p: wire.Notification, reply: Responder, request: string) {
if (typeof p.threadId !== "string") throw new Error("a question without a thread");
const questions = Array.isArray(p.questions) ? p.questions : [];
const id = newId("codex-question");
const questionIds = questions.map((question) => String(question?.id ?? ""));
const fields = questions.map((question) => {
const header = typeof question?.header === "string" ? question.header : "";
const text = typeof question?.question === "string" ? question.question : "";
const options = (Array.isArray(question?.options) ? question.options : []).map((option) => {
const out: import("./types.ts").QuestionOption = {
value: String(option?.label ?? ""),
label: String(option?.label ?? ""),
};
if (typeof option?.description === "string" && option.description) out.description = option.description;
return out;
});
const field: QuestionField = {
id: String(question?.id ?? ""),
label: header || text,
kind: options.length ? "select" : "text",
allowOther: question?.isOther === true,
required: true,
};
if (text) field.description = text;
if (options.length) field.options = options;
return field;
});
const entry = {
session: p.threadId,
request,
// `null` withdraws the question without answering it.
answer: (answer: QuestionAnswer | null) => {
if (this.questions.get(id) === entry) this.questions.delete(id);
if (answer === null) {
reply.discard();
return;
}
const values = answer?.values ?? {};
const answers: Record<string, { answers: string[] }> = {};
for (const question of questionIds) {
const value = values[question];
let list: string[] = [];
if (typeof value === "string") list = [value];
else if (Array.isArray(value)) list = value.filter((item) => typeof item === "string");
else if (typeof value === "boolean") list = [String(value)];
answers[question] = { answers: list };
}
reply.ok({ answers });
},
};
this.questions.set(id, entry);
this.emitFor(p.threadId, { event: "question", responseMode: "tool", id, fields });
}
// --- reading threads -------------------------------------------------------------------
/// Every turn of a thread, oldest first, following `nextCursor`.
/// `itemsView: "full"`: the default summary view drops the items the
/// transcript is built from.
allTurns = Effect.fn("Codex.allTurns")(function* (
this: CodexAgent,
connection: JsonRpcConnection,
thread: string,
): Effect.fn.Return<wire.Turn[], TransportError> {
const turns: wire.Turn[] = [];
let cursor: string | null = null;
for (;;) {
const response: z.infer<(typeof wire.responses)["thread/turns/list"]> = yield* this.request(
connection,
"thread/turns/list",
{ threadId: thread, cursor, limit: TURN_PAGE, sortDirection: "asc", itemsView: "full" },
);
turns.push(...(response.data ?? []));
const next: string | null | undefined = response.nextCursor;
if (!next) break;
cursor = next;
if (turns.length >= MAX_TURNS) {
console.warn(`stopping after ${MAX_TURNS} turns of history of ${thread}`);
break;
}
}
return turns;
});
/// A thread's record and every turn of it. `thread/read` with
/// `includeTurns` returns one page and is deprecated for `paginated`
/// threads, so the turns are paged separately; only a legacy thread on
/// an app-server without `thread/turns/list` falls back to it.
readThread = Effect.fn("Codex.readThread")(function* (
this: CodexAgent,
connection: JsonRpcConnection,
thread: string,
): Effect.fn.Return<History, TransportError> {
const response = yield* this.request(connection, "thread/read", { threadId: thread, includeTurns: false });
const record = response.thread ?? { id: thread };
const paginated = record.historyMode === "paginated";
let turns: wire.Turn[] | null = null;
const listed = yield* Effect.result(this.allTurns(connection, thread));
if (Result.isSuccess(listed)) {
if (listed.success.length || paginated) turns = listed.success;
} else
console.debug(`thread/turns/list is unavailable for ${thread}; using thread/read: ${listed.failure.message}`);
if (turns === null) {
const full = yield* this.request(connection, "thread/read", { threadId: thread, includeTurns: true });
turns = full.thread?.turns ?? [];
}
return { thread: record, turns };
});
/// Every subagent thread the turns start, and theirs in turn, read one
/// generation at a time. A child that cannot be read is left out; its
/// row is still rebuilt from what the parent's items say about it.
readChildren = Effect.fn("Codex.readChildren")(function* (
this: CodexAgent,
connection: JsonRpcConnection,
turns: wire.Turn[],
) {
const children = new Map<string, History>();
let generation = map.childIds(turns);
while (generation.length) {
const room = Math.max(0, MAX_SUBAGENT_THREADS - children.size);
if (generation.length > room) {
console.warn(`reading only ${MAX_SUBAGENT_THREADS} subagent threads`);
generation = generation.slice(0, room);
}
const read = yield* Effect.forEach(generation, (id) => Effect.result(this.readThread(connection, id)), {
concurrency: "unbounded",
});
const next: string[] = [];
generation.forEach((id, index) => {
// forEach preserves the generation's length and order.
const result = read[index]!;
if (Result.isSuccess(result)) {
next.push(...map.childIds(result.success.turns));
children.set(id, result.success);
} else console.debug(`could not read the subagent thread ${id}: ${result.failure.message}`);
});
const seen = new Set<string>();
generation = next.filter((id) => !children.has(id) && !seen.has(id) && Boolean(seen.add(id)));
}
return children;
});
/// The subagents below `thread` at every depth, `thread` included when
/// it is one itself.
subtree(thread: string) {
const found = new Set<string>();
if (this.subagents.has(thread)) found.add(thread);
const parents = [thread];
while (parents.length) {
const parent = parents.pop();
for (const [id, child] of this.subagents) {
if (parentOf(child) === parent && !found.has(id)) {
found.add(id);
parents.push(id);
}
}
}
return found;
}
/// Interrupts the live turn of each of `threads`, in parallel, giving
/// each a few seconds and all of them a bound, so a Stop never hangs on
/// a subagent that does not answer.
interruptChildren = Effect.fn("Codex.interruptChildren")(function* (
this: CodexAgent,
connection: JsonRpcConnection,
threads: Set<string>,
) {
const turns: [string, string][] = [];
for (const id of threads) {
const turn = this.subagents.get(id)?.turn;
if (turn) turns.push([id, turn]);
}
const interrupts = Effect.forEach(
turns,
([thread, turn]) =>
this.request(
connection,
"turn/interrupt",
{ threadId: thread, turnId: turn },
{ timeout: CHILD_INTERRUPT },
).pipe(
// Stop is best effort: the original provider also continued when a child refused or timed out.
Effect.catchTag(["TransportFailed", "ParseFailed"], () =>
Effect.sync(() => console.debug(`the subagent ${thread} did not stop in time`)),
),
),
{ concurrency: "unbounded" },
);
yield* Effect.raceFirst(interrupts, Effect.sleep(CHILDREN_INTERRUPT));
});
/// Cancels every approval `threads` wait on, so the app-server stops
/// waiting for the user.
releaseApprovals(threads: Set<string>) {
for (const [id, entry] of [...this.approvals]) {
if (!threads.has(entry.session)) continue;
this.approvals.delete(id);
entry.answer("cancel");
}
}
loadModels = Effect.fn("Codex.loadModels")(function* (
this: CodexAgent,
): Effect.fn.Return<wire.Model[], ProviderError | TransportError, PluginServices> {
if (this.models) return this.models;
const { connection } = yield* this.client();
const models: wire.Model[] = [];
let cursor: string | null = null;
for (;;) {
const response: z.infer<(typeof wire.responses)["model/list"]> = yield* this.request(connection, "model/list", {
cursor,
limit: 100,
});
models.push(...(response.data ?? []));
const next: string | null | undefined = response.nextCursor;
if (!next) break;
cursor = next;
}
this.models = models.filter((model) => typeof model.id === "string" && !model.hidden);
return this.models;
});
config = Effect.fn("Codex.config")(function* (this: CodexAgent, workspace: string) {
const result = yield* Effect.result(
Effect.gen(
function* (this: CodexAgent) {
const { connection } = yield* this.client();
const response = yield* this.request(connection, "config/read", { cwd: workspace });
return response.config ?? {};
}.bind(this),
),
);
if (Result.isSuccess(result)) return result.success;
console.warn(`config/read failed; using built-in defaults: ${result.failure.message}`);
return {};
});
optionsFor = Effect.fn("Codex.optionsFor")(function* (
this: CodexAgent,
workspace: string,
overrides: Record<string, unknown>,
) {
const models = yield* this.loadModels();
const config = yield* this.config(workspace);
return buildOptions(models, config, overrides);
});
/// The account's subscription limits, read on demand. Later changes
/// arrive as `account/rateLimits/updated`.
limits = Effect.fn("Codex.limits")(function* (this: CodexAgent) {
const result = yield* Effect.result(
Effect.gen(
function* (this: CodexAgent) {
const { connection } = yield* this.client();
const generation = this.accountGeneration;
// `account/read` exposes the active routing account, not just an
// email or CODEX_HOME. Older servers may omit it; the usage read's
// own native accountId remains valid evidence when present.
const account = yield* Effect.result(this.request(connection, "account/read", {}));
const response = yield* this.request(connection, "account/rateLimits/read", {
excludeResetCreditDetails: true,
});
if (generation !== this.accountGeneration) return null;
const identity = Result.isSuccess(account) ? account.success.workspaceRouting?.chatgptAccountId : null;
this.quotaFacts = map.usageRecovery(response, "codex/account/rateLimits/read", now(), identity);
return map.usageLimits(response, this.quotaFacts);
}.bind(this),
),
);
return Result.isSuccess(result) ? result.success : null;
});
/// Who owns the installed `codex`, and so whether Convergence may
/// upgrade it. Proven from the resolved real path; anything unproven
/// reports the version but stays manual.
maintenance = Effect.fn("Codex.maintenance")(function* (this: CodexAgent, installedVersion: string | null) {
const process = yield* Process;
const found = yield* process
.which("codex")
.pipe(
Effect.catchTag(["HostCallFailed", "NeedsReview", "PermissionNotGranted"], () =>
Effect.succeed({ path: null, realPath: null }),
),
);
const realPath = found.realPath;
const manager = installerOf(realPath);
let latest: string | null = null;
if (manager === "npm") latest = yield* this.latest("npm", NPM_PACKAGE, npmPrefix(realPath));
else if (manager === "homebrew") latest = yield* this.latest("brew", "codex", null);
else if (manager === "native") latest = yield* this.latest("npm", NPM_PACKAGE, null);
const out: Maintenance = { canUpdate: manager !== "" };
if (installedVersion) out.installedVersion = installedVersion;
if (installedVersion && latest && isNewer(installedVersion, latest)) out.latestVersion = latest;
if (manager) out.manager = manager;
return { out, realPath, manager };
});
// --- agent methods ------------------------------------------------------------------------
initialize = Effect.fn("Codex.initialize")(function* (this: CodexAgent) {
const { connection } = yield* this.client();
const version = versionOf(this.userAgent);
const account = yield* Effect.result(this.request(connection, "account/read", {}));
const status: AgentInfo["status"] = Result.isSuccess(account)
? (account.success.account !== null && account.success.account !== undefined) ||
!account.success.requiresOpenaiAuth
? { state: "ready" }
: { state: "auth_required", message: "Sign in with `codex login`" }
: { state: "unavailable", message: account.failure.message };
const [usageLimits, maintenance] = yield* Effect.all(
[this.limits(), this.maintenance(version).pipe(Effect.map((found) => found.out))],
{ concurrency: "unbounded" },
);
const info: AgentInfo = {
id: this.id,
name: this.name,
description: "OpenAI Codex, over the app-server protocol",
icon: ICON,
capabilities: {
reasoning: true,
images: true,
approvals: true,
questions: true,
sessionList: true,
sessionHistory: true,
resume: true,
slashCommands: false,
steer: true,
cancel: true,
skills: true,
rollback: true,
fork: true,
compact: true,
usageLimits: true,
subagents: true,
cancelTask: true,
},
authMethods: [
{
id: "chatgpt",
name: "Sign in with ChatGPT",
description: "Opens the browser, the same flow as `codex login`",
},
],
status,
maintenance,
family: FAMILY,
continuationKey: this.home() ?? "",
};
if (usageLimits) info.usageLimits = usageLimits;
if (version) info.version = version;
return info;
});
listOptions = Effect.fn("Codex.listOptions")(function* (this: CodexAgent, { workspace }: { workspace?: string }) {
return { options: yield* this.optionsFor(String(workspace ?? ""), {}) };
});
listCommands = Effect.fn("Codex.listCommands")(function* (this: CodexAgent) {
return { commands: [] };
});
listSessions = Effect.fn("Codex.listSessions")(function* (this: CodexAgent, { workspace }: { workspace?: string }) {
const { connection } = yield* this.client();
const response = yield* this.request(connection, "thread/list", { cwd: String(workspace ?? ""), limit: 100 });
return {
sessions: (response.data ?? [])
.filter((thread) => typeof thread.id === "string" && map.threadParent(thread) === null)
.map((thread) =>
compact({
id: thread.id,
title: map.threadTitle(thread.name, thread.preview, thread.turns),
createdAt: map.timestamp(thread.createdAt),
updatedAt: map.timestamp(thread.updatedAt),
}),
),
};
});
listSkills = Effect.fn("Codex.listSkills")(function* (this: CodexAgent, { workspace }: { workspace?: string }) {
const { connection } = yield* this.client();
const response = yield* this.request(connection, "skills/list", { cwds: [String(workspace ?? "")] });
const skills: { name: string; description: string; source: string; manualOnly: boolean }[] = [];
const paths = new Map<string, string>();
for (const entry of response.data ?? []) {
for (const skill of entry.skills ?? []) {
if (typeof skill.name !== "string" || typeof skill.path !== "string" || skill.enabled === false) continue;
paths.set(skill.name, skill.path);
skills.push({
name: skill.name,
description: typeof skill.shortDescription === "string" ? skill.shortDescription : (skill.description ?? ""),
source: typeof skill.scope === "string" ? skill.scope : "",
manualOnly: false,
});
}
}
this.skills = paths;
return { skills };
});
readSession = Effect.fn("Codex.readSession")(function* (
this: CodexAgent,
{ sessionId }: { sessionId: string; workspace?: string },
) {
const { connection } = yield* this.client();
const { thread, turns } = yield* this.readThread(connection, sessionId);
const title = map.threadTitle(thread.name, thread.preview, turns);
if (title) this.send({ sessionId, event: "session_info", title });
const children = yield* this.readChildren(connection, turns);
return { items: map.transcriptWithSubagents(turns, children, null) };
});
createSession = Effect.fn("Codex.createSession")(function* (this: CodexAgent, params: SessionOptions) {
const { connection } = yield* this.client();
const response = yield* this.request(connection, "thread/start", startParams(params)).pipe(
Effect.mapError((error) => new ProviderError({ message: `thread/start failed: ${errorMessage(error)}` })),
);
const id = response.thread?.id;
if (typeof id !== "string") return yield* new ProviderError({ message: "thread/start answered without a thread" });
this.sessions.set(id, newSession(params.workspace ?? "", params.options ?? {}));
return { sessionId: id };
});
resumeSession = Effect.fn("Codex.resumeSession")(function* (
this: CodexAgent,
params: SessionOptions & { sessionId: string },
) {
const { connection } = yield* this.client();
yield* this.request(connection, "thread/resume", resumeParams(params.sessionId, params)).pipe(
Effect.mapError((error) => new ProviderError({ message: `thread/resume failed: ${errorMessage(error)}` })),
);
const state = this.sessions.get(params.sessionId) ?? newSession();
state.workspace = params.workspace ?? "";
state.options = { ...(params.options ?? {}) };
state.unloaded = false;
this.sessions.set(params.sessionId, state);
return {};
});
closeSession = Effect.fn("Codex.closeSession")(function* (this: CodexAgent, { sessionId }: { sessionId: string }) {
if (this.proc?.alive)
yield* Effect.result(this.request(this.proc.connection, "thread/unsubscribe", { threadId: sessionId }));
this.sessions.delete(sessionId);
return {};
});
forkSession = Effect.fn("Codex.forkSession")(function* (
this: CodexAgent,
params: SessionOptions & { sessionId: string },
) {
const { connection } = yield* this.client();
const response = yield* this.request(connection, "thread/fork", forkParams(params.sessionId, params)).pipe(
Effect.mapError((error) => new ProviderError({ message: `thread/fork failed: ${errorMessage(error)}` })),
);
const id = response.thread?.id;
if (typeof id !== "string") return yield* new ProviderError({ message: "thread/fork answered without a thread" });
this.sessions.set(id, newSession(params.workspace ?? "", params.options ?? {}));
return { sessionId: id };
});
prompt = Effect.fn("Codex.prompt")(function* (
this: CodexAgent,
{ sessionId, input }: { sessionId: string; input?: PromptInput },
) {
const delivery = input?.delivery;
const rejected = (reason: string) => {
if (delivery)
this.send({
sessionId,
event: "input_rejected",
inputId: delivery.inputId,
attemptId: delivery.attemptId,
reason,
});
return new ProviderError({ message: reason });
};
if (delivery?.intent === "queue") return yield* rejected("Queued input must remain held by the host");
const connection = yield* this.loadedConnection(sessionId);
const known = this.sessions.get(sessionId);
const workspace = known?.workspace ?? "";
const options = { ...(known?.options ?? {}) };
const blocks = userInput(input?.blocks, this.skills);
if (known?.run) {
const runId = known.run;
const turn = known.turn ?? (known.starting ? yield* Deferred.await(known.starting) : null);
if (!turn || known.run !== runId)
return yield* rejected("The original Codex turn ended before steering could be submitted");
if (delivery) known.inputs.set(delivery.attemptId, { ...delivery, runId });
const result = yield* Effect.result(
this.request(
connection,
"turn/steer",
{
threadId: sessionId,
input: blocks,
expectedTurnId: turn,
...(delivery ? { clientUserMessageId: delivery.attemptId } : {}),
},
{ timeout: null },
),
);
if (Result.isFailure(result)) {
const reason = `turn/steer failed: ${errorMessage(result.failure)}`;
if (nativeRejection(result.failure)) {
if (delivery) known.inputs.delete(delivery.attemptId);
return yield* rejected(reason);
}
return yield* new ProviderError({ message: reason });
}
return {
runId,
receipt: { evidence: "native_admission" as const, ...(delivery ? { nativeInputId: delivery.attemptId } : {}) },
};
}
const runId = newId("codex-run");
const hostItem = typeof input?.itemId === "string" ? input.itemId : null;
const state = known ?? newSession(workspace, options);
this.sessions.set(sessionId, state);
if (!state.workspace) state.workspace = workspace;
state.run = runId;
state.turn = null;
state.streamed.clear();
state.startedTools.clear();
state.held.clear();
const starting = Deferred.makeUnsafe<string | null>();
state.starting = starting;
if (delivery) state.inputs.set(delivery.attemptId, { ...delivery, runId });
const fiber = yield* this.request(
connection,
"turn/start",
{
...turnParams(sessionId, blocks, workspace, options),
...(delivery ? { clientUserMessageId: delivery.attemptId } : {}),
},
{
timeout: null,
},
).pipe(
Effect.tap((response) =>
Effect.sync(() => {
const current = this.sessions.get(sessionId);
const turn = response.turn?.id;
if (!current || typeof turn !== "string") {
Deferred.doneUnsafe(starting, Effect.succeed(null));
return;
}
if (hostItem) current.turns.set(hostItem, turn);
if (current.run === runId) current.turn ??= turn;
else if (current.turn !== turn) current.endedTurns.add(turn);
Deferred.doneUnsafe(starting, Effect.succeed(turn));
}),
),
Effect.tapError((error) =>
Effect.sync(() => {
Deferred.doneUnsafe(starting, Effect.succeed(null));
if (nativeRejection(error)) {
if (delivery) state.inputs.delete(delivery.attemptId);
rejected(errorMessage(error));
}
if (state.run === runId) {
this.notice(sessionId, "error", errorMessage(error));
this.finishRun(sessionId, outcome.failed(errorMessage(error)));
}
}),
),
Effect.result,
Effect.forkIn(this.scope),
);
if (delivery) {
const result = yield* Fiber.join(fiber);
if (Result.isFailure(result)) return yield* new ProviderError({ message: errorMessage(result.failure) });
if (!result.success.turn?.id)
return yield* new ProviderError({
message: "turn/start answered without a turn identity; delivery is uncertain",
});
return { runId, receipt: { evidence: "native_admission" as const, nativeInputId: delivery.attemptId } };
}
return { runId };
});
cancel = Effect.fn("Codex.cancel")(function* (this: CodexAgent, { sessionId }: { sessionId: string }) {
const threads = this.subtree(sessionId);
if (this.proc?.alive) {
yield* this.interruptChildren(this.proc.connection, threads);
const turn = this.sessions.get(sessionId)?.turn;
if (turn)
yield* Effect.result(
this.request(this.proc.connection, "turn/interrupt", { threadId: sessionId, turnId: turn }),
);
}
threads.add(sessionId);
this.releaseApprovals(threads);
this.finishRun(sessionId, outcome.cancelled());
return {};
});
cancelTask = Effect.fn("Codex.cancelTask")(function* (
this: CodexAgent,
{ sessionId, taskId }: { sessionId: string; taskId: string },
) {
const child = this.subagents.get(taskId);
if (!child || child.root !== sessionId)
return yield* new ProviderError({ message: `no codex subagent ${taskId} in this chat` });
if (!child.turn) return yield* new ProviderError({ message: "the subagent is not running" });
const threads = this.subtree(taskId);
const { connection } = yield* this.client();
yield* this.interruptChildren(connection, threads);
this.releaseApprovals(threads);
return {};
});
setOption = Effect.fn("Codex.setOption")(function* (
this: CodexAgent,
{ sessionId, optionId, value }: { sessionId: string; optionId: string; value: unknown },
) {
const state = this.sessions.get(sessionId) ?? newSession();
this.sessions.set(sessionId, state);
state.options[optionId] = value;
if (optionId === OPTION_MODEL) delete state.options[OPTION_REASONING];
return { options: yield* this.optionsFor(state.workspace, state.options) };
});
respondToApproval = Effect.fn("Codex.respondToApproval")(function* (
this: CodexAgent,
{ approvalId, optionId }: { approvalId: string; optionId: string },
) {
const entry = this.approvals.get(approvalId);
if (!entry) return yield* new ProviderError({ message: `no codex approval is waiting for ${approvalId}` });
this.approvals.delete(approvalId);
const decision = ["accept", "acceptForSession", "acceptAlways", "decline", "cancel"].includes(optionId)
? optionId
: "decline";
entry.answer(decision);
return {};
});
respondToQuestion = Effect.fn("Codex.respondToQuestion")(function* (
this: CodexAgent,
{ questionId, answer }: { questionId: string; answer?: QuestionAnswer },
) {
const entry = this.questions.get(questionId);
if (!entry) return yield* new ProviderError({ message: `no codex question is waiting for ${questionId}` });
this.questions.delete(questionId);
entry.answer(answer ?? { values: {}, cancelled: false });
return {};
});
/// Rewinds the thread so `itemId` and everything after it is gone. Codex
/// reverts by turn, so the turn holding that item is found first;
/// `thread/revert` rewrites the history only, which is why the host
/// pairs it with its own checkpoint.
rollback = Effect.fn("Codex.rollback")(function* (
this: CodexAgent,
{ sessionId, itemId }: { sessionId: string; itemId: string },
) {
const connection = yield* this.loadedConnection(sessionId);
let turnId = this.sessions.get(sessionId)?.turns.get(itemId) ?? null;
if (!turnId) {
const turns = yield* this.allTurns(connection, sessionId);
turnId = turns.find((turn) => (turn.items ?? []).some((item) => map.itemId(item) === itemId))?.id ?? null;
if (!turnId) return yield* new ProviderError({ message: `codex does not know message ${itemId}` });
}
yield* this.request(connection, "thread/revert", { threadId: sessionId, beforeTurnId: turnId }).pipe(
Effect.mapError((error) => new ProviderError({ message: `thread/revert failed: ${errorMessage(error)}` })),
);
return {};
});
compact = Effect.fn("Codex.compact")(function* (this: CodexAgent, { sessionId }: { sessionId: string }) {
const connection = yield* this.loadedConnection(sessionId);
yield* this.request(connection, "thread/compact/start", { threadId: sessionId }, { timeout: null }).pipe(
Effect.mapError((error) => new ProviderError({ message: `thread/compact/start failed: ${errorMessage(error)}` })),
);
return {};
});
usageLimits = Effect.fn("Codex.usageLimits")(function* (this: CodexAgent) {
return { limits: yield* this.limits() };
});
/// Starts the ChatGPT sign-in: the app-server answers with the page to
/// open, and finishes the login itself when the browser comes back.
authenticate = Effect.fn("Codex.authenticate")(function* (this: CodexAgent, _params?: { method: string }) {
const { connection } = yield* this.client();
const response = yield* this.request(
connection,
"account/login/start",
{ type: "chatgpt" },
{ timeout: null },
).pipe(
Effect.mapError((error) => new ProviderError({ message: `account/login/start failed: ${errorMessage(error)}` })),
);
this.accountDirty = true;
this.invalidateQuota();
const url = response.authUrl ?? response.verificationUrl;
if (typeof url === "string" && url) {
const kernel = yield* Kernel;
const result = yield* kernel
.call("open_url", { url })
.pipe(
Effect.mapError(
(error) => new ProviderError({ message: `open ${url} in a browser to sign in (${errorMessage(error)})` }),
),
);
if (result.opened === false)
return yield* new ProviderError({
message: `open ${url} in a browser to sign in (${result.reason ?? "the link was not opened"})`,
});
}
return {};
});
logout = Effect.fn("Codex.logout")(function* (this: CodexAgent) {
const { connection } = yield* this.client();
yield* this.request(connection, "account/logout", {}).pipe(
Effect.mapError((error) => new ProviderError({ message: `account/logout failed: ${errorMessage(error)}` })),
);
this.accountDirty = true;
this.invalidateQuota();
return {};
});
/// Upgrades the installed Codex through whichever installer owns it.
update = Effect.fn("Codex.update")(function* (this: CodexAgent) {
const { manager, realPath } = yield* this.maintenance(null);
let result;
if (manager === "npm") {
const prefix = npmPrefix(realPath);
if (!prefix) return yield* new ProviderError({ message: "could not find the npm prefix that owns codex" });
result = yield* runProcess("npm", ["install", "-g", "--prefix", prefix, `${NPM_PACKAGE}@latest`]).pipe(
Effect.scoped,
);
} else if (manager === "homebrew") result = yield* runProcess("brew", ["upgrade", "codex"]).pipe(Effect.scoped);
else if (manager === "native") result = yield* runProcess("codex", ["update"]).pipe(Effect.scoped);
else
return yield* new ProviderError({
message: "Divergence cannot prove which installer owns this codex, so it must be updated by hand",
});
if (result.code !== 0)
return yield* new ProviderError({
message: result.stderr.trim() || `the update ended with ${result.code ?? result.signal}`,
});
yield* this.shutdownProcess();
this.models = null;
return {};
});
}Versions
| Version | Published | Plugin API | Size | Permissions | Status |
|---|---|---|---|---|---|
| 0.2.0latest | Oct 5, 2026 | >=2 <3 | 78.1 KB | 4 permissions | Listed |
No comments yet.