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

  • Provide agents agents.provideMediumAdds agents to the app.Provide the Codex agent and pass its tool calls to plugin tools
  • Run named programs processMediumStarts the listed programs.Run the Codex CLI, and ask or tell the installer that owns it (npm or Homebrew) about a newer versionPrograms: codexnpmbrew
  • Environment variables envMediumReads the listed environment variables.Find each account's Codex home, and pass extra app-server arguments set in CODEX_ARGSVariables: HOMECODEX_HOMECODEX_ARGS
  • Read files fs.readMediumReads files in the listed places.Read, once, the accounts the previous Codex provider keptPlaces: its own data folder

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

VersionPublishedPlugin APISizePermissionsStatus
0.2.0latestOct 5, 2026>=2 <378.1 KB4 permissionsListed

Reviews and comments

0 threads · 0 reviews

No comments yet.