Official

opencode-v2

OpenCode 2 agent provider: runs opencode serve (2.x) and talks to its /api over HTTP and server-sent events.

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/opencode-v2@0.1.0

Permissions in 0.1.0

  • Provide agents agents.provideMediumAdds agents to the app.Provide the OpenCode 2 agent, and serve plugin tools to it through the host's loopback MCP server
  • Run named programs processMediumStarts the listed programs.Run the OpenCode 2 server (`opencode serve`, or the binary you choose in Settings) and read its versionPrograms: opencode${settings.binaryPath}
  • Network access netMediumConnects to the listed hosts.Talk to the OpenCode server it starts on this computer, at the port it picks for each launch, and to the external server you choose in SettingsHosts: localhost:*${settings.serverUrl}
  • Environment variables envMediumReads the listed environment variables.Read the user name and password of the OpenCode serverVariables: OPENCODE_SERVER_USERNAMEOPENCODE_SERVER_PASSWORDCONVERGENCE_OPENCODE_SERVER_PASSWORD

Files

agent.ts61.1 KB
// The OpenCode 2 provider: one `opencode serve` (or a server the user
// runs), one event stream for every session (`GET /api/event`), and the
// session calls of the 2.x API. See NOTES.md for the API facts behind the
// choices here.
import * as Deferred from "effect/Deferred";
import * as Effect from "effect/Effect";
import * as Fiber from "effect/Fiber";
import * as Result from "effect/Result";
import * as Scope from "effect/Scope";
import * as Semaphore from "effect/Semaphore";
import * as Stream from "effect/Stream";
import * as z from "zod";
import { Host } from "convergence/effect";
import type { EffectAgent, PluginServices } from "convergence/effect";
import type {
  AgentEvent,
  AgentEventKind,
  AgentInfo,
  ConfigOption,
  NoticeLevel,
  PromptInput,
  QuestionAnswer,
  RunOutcome,
  SessionOptions,
  SessionSummary,
  Skill,
  SlashCommand,
  TaskInfo,
  TaskStatus,
  ToolSpec,
  TranscriptItem,
} from "convergence/protocol";
import { newId, outcome, permissionMode } from "../sdk/agent.ts";
import { sseStream } from "../sdk/effect.ts";
import { errorMessage } from "../sdk/errors.ts";
import { MCP_SERVER } from "../sdk/mcp.ts";
import * as api from "./api.ts";
import type { Form, ModelRef, ServerEvent, Tokens } from "./api.ts";
import { ICON } from "./icon.ts";
import * as map from "./map.ts";
import {
  OPTION_AGENT,
  OPTION_MODEL,
  OPTION_VARIANT,
  buildOptions,
  chosenModel,
  contextWindow,
  usableModels,
} from "./options.ts";
import type { Catalog } from "./options.ts";
import { approvalTitle, repliesAutomatically, ruleset } from "./permission.ts";
import { messageIdFor, promptBody, slashCommand } from "./prompt.ts";
import type { PromptBody } from "./prompt.ts";
import {
  DEFAULT_USERNAME,
  LocalServer,
  ServerError,
  atLocation,
  basicAuth,
  generatedPassword,
  isLoopback,
  passwordFor,
  request,
} from "./server.ts";
import type { Auth, RequestOptions } from "./server.ts";

export const AGENT_ID = "opencode-v2";
export const FAMILY = "opencode-v2";

/// Reconnect backoff for the event stream, and how long an external server
/// may stay unreachable before its runs are failed.
const RECONNECT_MIN = 250;
const RECONNECT_MAX = 10_000;
const RECONNECT_RESET = 30_000;
const RECONNECT_GIVE_UP = 60_000;
/// How long the first connection of the event stream may take.
const STREAM_OPEN_TIMEOUT = 15_000;
/// How deep `read_session` follows subagents of subagents.
const MAX_SUBAGENT_DEPTH = 8;
/// Page size for sessions and messages.
const PAGE = 200;
/// How long a catalog read waits for the server to finish building its
/// catalog after a start, and how often it looks again meanwhile.
const CATALOG_WAIT = 20_000;
const CATALOG_POLL = 500;
/// The instruction entry that carries the session's plugin rules.
const INSTRUCTIONS_KEY = "convergence";

/// Where the server comes from: the binary to start, or the address of one
/// the user runs (which wins).
export interface Config {
  binaryPath: string | null;
  serverUrl: string | null;
}

/// A chat the host knows. `applied` is what the server's session was last
/// told, so a prompt only switches what changed.
export interface SessionState {
  workspace: string;
  options: Record<string, unknown>;
  tools: readonly ToolSpec[];
  instructions: string | null;
  applied: { model: string | null; agent: string | null; mode: string | null; instructions: string | null };
}

/// A subagent: a session of its own whose events show in the chat that
/// started it.
export interface Child {
  root: string;
  parent: string;
  task: TaskInfo;
}

/// A tool call seen on the stream, by its host item id.
interface CallState {
  session: string;
  name: string;
  input: Record<string, unknown>;
}

interface EventStream {
  ready: Deferred.Deferred<void, unknown>;
  fiber: Fiber.Fiber<void, unknown> | null;
}

const modelKey = (ref: ModelRef | null) => (ref ? `${ref.providerID}/${ref.id}#${ref.variant ?? ""}` : null);
const now = () => new Date().toISOString();

export class OpenCodeV2Agent {
  readonly scope: Scope.Scope;
  config: Config;
  env: Record<string, string>;
  emit: (event: AgentEvent) => void = () => {};

  server: LocalServer | null = null;
  starting: Fiber.Fiber<LocalServer, unknown> | null = null;
  password: string | null = null;
  stream: EventStream | null = null;

  sessions = new Map<string, SessionState>();
  runs = new Map<string, string>();
  promptLocks = new Map<string, Semaphore.Semaphore>();
  inputs = new Map<
    string,
    {
      session: string;
      inputId: string;
      attemptId: string;
      runId: string;
      consumed: boolean;
      admission?: api.AdmittedInput;
    }
  >();
  children = new Map<string, Child>();
  /// Sessions looked up and found to belong to no chat here (another
  /// client's), so their events are skipped without asking again.
  strangers = new Set<string>();
  calls = new Map<string, CallState>();
  started = new Set<string>();
  /// Subagent calls whose session has not appeared yet, per parent session.
  unlinked = new Map<string, string[]>();
  approvals = new Map<string, string>();
  forms = new Map<string, Form>();
  /// The tokens of each session's last step, for the usage that follows.
  steps = new Map<string, Tokens>();
  ranOn = new Map<string, ModelRef>();
  catalogs = new Map<string, Catalog>();
  /// Completed (and replaced) whenever a catalog changes on the server.
  catalogSignal = Deferred.makeUnsafe<void>();
  commands = new Map<string, readonly { name: string; description?: string }[]>();
  /// The tools server each workspace was given: its address, the tools it
  /// serves, and the server process that was told.
  mcp = new Map<string, { url: string; key: string; server: LocalServer | null }>();

  constructor(scope: Scope.Scope, config: Config, env: Record<string, string> = {}) {
    this.scope = scope;
    this.config = config;
    this.env = env;
  }

  /// The handlers the host calls, by `agent/<method>` name.
  definition(): EffectAgent {
    return {
      id: AGENT_ID,
      name: "OpenCode",
      icon: ICON,
      description: "The OpenCode 2 coding agent.",
      initialize: () => this.initialize(),
      list_options: (params) => this.listOptions(params.workspace),
      list_commands: (params) => this.listCommands(params.sessionId),
      list_sessions: (params) => this.listSessions(params.workspace),
      read_session: (params) => this.readSession(params.sessionId),
      list_skills: (params) => this.listSkills(params.workspace),
      create_session: (params) => this.createSession(params),
      resume_session: (params) => this.resumeSession(params),
      close_session: (params) => this.closeSession(params.sessionId),
      fork_session: (params) => this.forkSession(params),
      prompt: (params) => this.prompt(params.sessionId, params.input),
      cancel: (params) => this.cancel(params.sessionId),
      cancel_task: (params) => this.cancelTask(params.sessionId, params.taskId),
      set_option: (params) => this.setOption(params.sessionId, params.optionId, params.value),
      respond_to_approval: (params) => this.respondToApproval(params.approvalId, params.optionId),
      respond_to_question: (params) => this.respondToQuestion(params.questionId, params.answer),
      compact: (params) => this.compact(params.sessionId),
    };
  }

  // --- the server --------------------------------------------------------------------

  get external(): string | null {
    return this.config.serverUrl ? this.config.serverUrl.replace(/\/+$/, "") : null;
  }

  /// The credentials. A server this plugin starts always has a password:
  /// the configured one, or one made up for this plugin's lifetime.
  get auth(): Auth | null {
    const username = this.env["OPENCODE_SERVER_USERNAME"] || DEFAULT_USERNAME;
    const external = this.external !== null;
    const password = passwordFor(
      external,
      this.env["CONVERGENCE_OPENCODE_SERVER_PASSWORD"] ?? null,
      this.env["OPENCODE_SERVER_PASSWORD"] ?? null,
    );
    if (password) return { username, password };
    if (external) return null;
    this.password ??= generatedPassword();
    return { username, password: this.password };
  }

  /// The server's address, starting a local server the first time and again
  /// after it exited.
  base = Effect.fn("OpenCodeV2.base")(function* (this: OpenCodeV2Agent) {
    const external = this.external;
    if (external) return external;
    if (this.server && !this.server.exited) return this.server.base;
    if (this.server) {
      console.warn(`opencode-v2: the server exited${this.server.logTail()}; starting it again`);
      this.server = null;
      this.catalogs.clear();
      this.commands.clear();
    }
    const auth = this.auth ?? { username: DEFAULT_USERNAME, password: generatedPassword() };
    const fiber =
      this.starting ??
      (yield* LocalServer.start(this.config.binaryPath ?? "opencode", auth).pipe(
        Effect.tap((server) =>
          Effect.sync(() => {
            this.server = server;
            console.info(`opencode-v2: server ${server.version} ready at ${server.base}`);
          }),
        ),
        Effect.forkIn(this.scope),
      ));
    this.starting = fiber;
    const server = yield* Fiber.join(fiber).pipe(
      Effect.ensuring(
        Effect.sync(() => {
          if (this.starting === fiber) this.starting = null;
        }),
      ),
    );
    return server.base;
  });

  /// One request to the server, its answer parsed with `schema`. A server
  /// address no grant covers is explained, not reported as "not granted".
  call = Effect.fn("OpenCodeV2.call")(function* <S extends z.ZodType>(
    this: OpenCodeV2Agent,
    method: string,
    path: string,
    schema: S,
    options: RequestOptions = {},
  ) {
    const base = yield* this.base();
    return yield* request(base, this.auth, method, path, schema, options).pipe(
      Effect.catchTag("PermissionNotGranted", (error) =>
        Effect.fail(
          this.external
            ? new ServerError({
                message: `the OpenCode server at ${base} is not allowed yet: set it as the External server in Settings > OpenCode 2 and allow the connection on its card`,
              })
            : error,
        ),
      ),
    );
  });

  /// A request whose answer is not read (most writes answer 204).
  send(method: string, path: string, options: RequestOptions = {}) {
    return this.call(method, path, z.unknown(), options).pipe(Effect.asVoid);
  }

  /// Stops the event stream and the server this plugin started.
  shutdown = Effect.fn("OpenCodeV2.shutdown")(function* (this: OpenCodeV2Agent) {
    const stream = this.stream;
    this.stream = null;
    if (stream?.fiber) yield* Fiber.interrupt(stream.fiber);
    const starting = this.starting;
    this.starting = null;
    if (starting) yield* Fiber.interrupt(starting);
    const server = this.server;
    this.server = null;
    if (server) yield* server.stop();
    this.mcp.clear();
    this.catalogs.clear();
    this.commands.clear();
  });

  /// A new binary or server address applies from the next call: the old
  /// server and its stream stop, and the runs on them end.
  reconfigure = Effect.fn("OpenCodeV2.reconfigure")(function* (this: OpenCodeV2Agent, config: Config) {
    if (config.binaryPath === this.config.binaryPath && config.serverUrl === this.config.serverUrl) return;
    this.config = config;
    this.password = null;
    this.failRuns("the OpenCode server changed in Settings");
    yield* this.shutdown();
  });

  // --- routing and events ----------------------------------------------------------------

  /// The chat state of a session the host created or resumed here.
  known = Effect.fn("OpenCodeV2.known")(function* (this: OpenCodeV2Agent, session: string) {
    const state = this.sessions.get(session);
    if (!state) return yield* new ServerError({ message: `unknown OpenCode session ${session}` });
    return state;
  });

  /// The chat a session's events go to, and the subagent it is.
  route(session: string): { root: string; task: string | null } | null {
    if (this.sessions.has(session)) return { root: session, task: null };
    const child = this.children.get(session);
    return child ? { root: child.root, task: session } : null;
  }

  emitFor(session: string, event: AgentEventKind): void {
    const at = this.route(session);
    if (!at) return;
    const runId = this.runs.get(at.root);
    this.emit({ sessionId: at.root, ...(runId ? { runId } : {}), ...(at.task ? { taskId: at.task } : {}), ...event });
  }

  notice(session: string, level: NoticeLevel, message: string): void {
    this.emitFor(session, { event: "notice", level, message });
  }

  /// Ends the run of a chat, when one is active. Removing it first makes
  /// exactly one `run_finished` per run.
  finish(session: string, result: RunOutcome): void {
    const runId = this.runs.get(session);
    if (!runId) return;
    this.runs.delete(session);
    this.forgetCalls(session);
    this.emit({ sessionId: session, runId, event: "run_finished", outcome: result });
  }

  /// Drops what the stream told about the calls of a session whose turn is
  /// over. A subagent's calls go when it settles: it can outlive the turn
  /// that started it.
  forgetCalls(session: string): void {
    for (const [id, call] of this.calls)
      if (call.session === session) {
        this.calls.delete(id);
        this.started.delete(id);
      }
    this.unlinked.delete(session);
    this.steps.delete(session);
  }

  /// Fails every active run: the server they ran on is gone.
  failRuns(message: string): void {
    for (const session of [...this.runs.keys()]) {
      this.notice(session, "error", message);
      this.finish(session, outcome.failed(message));
    }
    for (const [id, session] of this.approvals) this.emitFor(session, { event: "approval_resolved", id });
    for (const [id, form] of this.forms) this.emitFor(form.sessionID, { event: "question_resolved", id });
    this.approvals.clear();
    this.forms.clear();
  }

  /// Subscribes to the event stream, and waits until it is open, so the
  /// events of a prompt sent right after are not missed.
  ensureStream = Effect.fn("OpenCodeV2.ensureStream")(function* (this: OpenCodeV2Agent) {
    let stream = this.stream;
    if (!stream) {
      const ready = yield* Deferred.make<void, unknown>();
      const opened: EventStream = { ready, fiber: null };
      this.stream = opened;
      opened.fiber = yield* this.supervise(opened).pipe(
        Effect.ensuring(
          Effect.sync(() => {
            if (this.stream === opened) this.stream = null;
          }),
        ),
        Effect.forkIn(this.scope),
      );
      stream = opened;
    }
    yield* Deferred.await(stream.ready).pipe(
      Effect.timeoutOrElse({
        duration: STREAM_OPEN_TIMEOUT,
        orElse: () => Effect.fail(new ServerError({ message: "the OpenCode event stream did not open within 15s" })),
      }),
    );
  });

  /// Keeps the event stream open while the server lives: a dropped
  /// connection reconnects, because a transport blip is not the end of a
  /// run. Only a server that is gone fails the runs.
  supervise = Effect.fn("OpenCodeV2.supervise")(function* (this: OpenCodeV2Agent, stream: EventStream) {
    let backoff = RECONNECT_MIN;
    let failingSince: number | null = null;
    let reattach = false;
    for (;;) {
      const startedAt = Date.now();
      const read = yield* this.read(stream, reattach).pipe(Effect.result);
      const message = Result.isSuccess(read) ? "the OpenCode event stream closed" : errorMessage(read.failure);
      // The first connection failed: whoever waits is told, and the next
      // call tries again from the start.
      if (!(yield* Deferred.isDone(stream.ready))) {
        yield* Deferred.fail(stream.ready, Result.isFailure(read) ? read.failure : new ServerError({ message }));
        return;
      }
      reattach = true;
      if (Date.now() - startedAt >= RECONNECT_RESET) {
        backoff = RECONNECT_MIN;
        failingSince = null;
      }
      failingSince ??= Date.now();
      if (this.server?.exited) {
        this.failRuns(`the OpenCode server exited${this.server.logTail()}`);
        return;
      }
      if (this.external && Date.now() - failingSince >= RECONNECT_GIVE_UP) {
        this.failRuns(message);
        return;
      }
      console.warn(`opencode-v2: ${message}; reconnecting in ${backoff} ms`);
      yield* Effect.sleep(backoff);
      backoff = Math.min(backoff * 2, RECONNECT_MAX);
    }
  });

  /// Reads `GET /api/event` until it ends. The first frame opens the
  /// stream; after a reconnect, what the server did while we were away is
  /// recovered first.
  read = Effect.fn("OpenCodeV2.read")(function* (this: OpenCodeV2Agent, stream: EventStream, reattach: boolean) {
    const base = yield* this.base();
    const authorization = basicAuth(this.auth);
    let first = true;
    yield* sseStream(`${base}/api/event`, { headers: authorization ? { authorization } : {} }).pipe(
      Stream.runForEach((frame) => {
        const opening = first;
        first = false;
        return this.onFrame(frame.data, opening ? stream : null, reattach);
      }),
    );
  });

  /// One frame of the stream. `opened` is the stream the frame opens (the
  /// first frame of a connection), whose waiters are released.
  onFrame = Effect.fn("OpenCodeV2.onFrame")(function* (
    this: OpenCodeV2Agent,
    data: string,
    opened: EventStream | null,
    reattach: boolean,
  ) {
    if (opened) {
      yield* Deferred.succeed(opened.ready, undefined);
      if (reattach) yield* this.rehydrate().pipe(Effect.forkIn(this.scope));
    }
    const raw = yield* Effect.try((): unknown => JSON.parse(data)).pipe(Effect.orElseSucceed(() => null));
    const change = api.catalogChange(raw);
    if (change) return yield* this.onCatalogChange(change.directory);
    const event = api.parseEvent(raw);
    if (!event) return;
    const handled = yield* this.dispatch(event).pipe(Effect.result);
    if (Result.isFailure(handled))
      console.warn(`opencode-v2: the ${event.type} event could not be handled: ${errorMessage(handled.failure)}`);
  });

  /// After a reconnect: requests asked while the stream was down are
  /// announced and requests the server forgot are withdrawn, and a run
  /// whose turn ended meanwhile ends too.
  rehydrate = Effect.fn("OpenCodeV2.rehydrate")(function* (this: OpenCodeV2Agent) {
    const active = yield* this.call("GET", "/api/session/active", api.active);
    for (const session of [...this.runs.keys()]) {
      if (Object.hasOwn(active.data, session)) continue;
      const info = yield* this.call(
        "GET",
        `/api/session/${session}`,
        api.single(z.looseObject({ outcome: z.string().optional() })),
      );
      const ended = info.data.outcome;
      this.finish(
        session,
        ended === "interrupted"
          ? outcome.cancelled()
          : ended === "failed"
            ? outcome.failed("the OpenCode turn failed")
            : outcome.completed(),
      );
    }
    const workspaces = new Set([...this.sessions.values()].map((state) => state.workspace));
    const askedPermissions = new Set<string>();
    const askedForms = new Set<string>();
    for (const workspace of workspaces) {
      const permissions = yield* this.call("GET", "/api/permission/request", api.listed(api.permissionRequest), {
        query: atLocation(workspace),
      });
      for (const request of permissions.data) {
        askedPermissions.add(request.id);
        if (!this.approvals.has(request.id)) yield* this.onPermission(request);
      }
      const forms = yield* this.call("GET", "/api/form", api.listed(api.form), { query: atLocation(workspace) });
      for (const form of forms.data) {
        askedForms.add(form.id);
        if (!this.forms.has(form.id)) this.onForm(form);
      }
    }
    for (const [id, session] of [...this.approvals])
      if (!askedPermissions.has(id)) {
        this.approvals.delete(id);
        this.emitFor(session, { event: "approval_resolved", id });
      }
    for (const [id, form] of [...this.forms])
      if (!askedForms.has(id)) {
        this.forms.delete(id);
        this.emitFor(form.sessionID, { event: "question_resolved", id });
      }
  });

  /// Learns whether a session the stream mentions is a subagent of a chat
  /// here (one that started before the chat was resumed); others are
  /// remembered as strangers.
  adopt: (session: string) => Effect.Effect<void, unknown, PluginServices> = Effect.fn("OpenCodeV2.adopt")(function* (
    this: OpenCodeV2Agent,
    session: string,
  ) {
    if (this.route(session) || this.strangers.has(session)) return;
    this.strangers.add(session);
    const info = yield* this.call("GET", `/api/session/${session}`, api.single(api.session));
    const parent = info.data.parentID;
    if (!parent) return;
    yield* this.adopt(parent);
    const at = this.route(parent);
    if (!at) return;
    this.strangers.delete(session);
    this.addChild(session, parent, at.root, {
      title: info.data.title,
      name: info.data.agent,
      model: info.data.model?.id,
    });
  });

  addChild(
    session: string,
    parent: string,
    root: string,
    known: { title?: string; name?: string; model?: string },
  ): void {
    const waiting = this.unlinked.get(parent) ?? [];
    // The call that started it: the one with its description, else the
    // oldest call still waiting for its session.
    const call = waiting.find((id) => this.calls.get(id)?.input["description"] === known.title) ?? waiting[0] ?? null;
    if (call) waiting.splice(waiting.indexOf(call), 1);
    const input = call ? (this.calls.get(call)?.input ?? {}) : {};
    const prompt = input["prompt"];
    const task: TaskInfo = {
      id: session,
      title: known.title || (typeof prompt === "string" ? prompt.split("\n")[0] : "") || "Subagent",
      status: "running",
      startedAt: now(),
      toolUses: 0,
      ...(call ? { toolCallId: call } : {}),
      ...(this.children.has(parent) ? { parentTaskId: parent } : {}),
      ...(known.name ? { name: known.name } : {}),
      ...(known.model ? { model: known.model } : {}),
      ...(typeof prompt === "string" ? { prompt } : {}),
      ...(input["background"] === true ? { background: true } : {}),
    };
    this.children.set(session, { root, parent, task });
    this.emitTask(session);
  }

  emitTask(session: string): void {
    const child = this.children.get(session);
    if (!child) return;
    const runId = this.runs.get(child.root);
    this.emit({ sessionId: child.root, ...(runId ? { runId } : {}), event: "task", ...child.task });
  }

  updateTask(session: string, change: Partial<TaskInfo>): void {
    const child = this.children.get(session);
    if (!child) return;
    child.task = { ...child.task, ...change };
    this.emitTask(session);
  }

  settleTask(session: string, status: TaskStatus): void {
    const child = this.children.get(session);
    // A task settled once (a cancel) keeps that outcome.
    if (!child || child.task.status === "cancelled") return;
    const live = status === "running" || status === "waiting";
    if (!live) this.forgetCalls(session);
    this.updateTask(session, { status, ...(live ? {} : { endedAt: now() }) });
  }

  /// Applies one event.
  dispatch = Effect.fn("OpenCodeV2.dispatch")(function* (this: OpenCodeV2Agent, event: ServerEvent) {
    const session = event.type === "form.created" ? event.data.form.sessionID : event.data.sessionID;
    if (event.type === "session.created") {
      const parent = event.data.parentID;
      const at = parent ? this.route(parent) : null;
      if (parent && at && !this.route(session))
        this.addChild(session, parent, at.root, {
          title: event.data.title,
          name: event.data.agent,
          model: event.data.model?.id,
        });
      return;
    }
    if (!this.route(session)) yield* this.adopt(session).pipe(Effect.ignore);
    const at = this.route(session);
    if (!at) return;
    switch (event.type) {
      case "session.renamed":
        if (!at.task) this.emitFor(session, { event: "session_info", title: event.data.title });
        return;
      case "session.execution.started":
        if (at.task) return this.settleTask(session, "running");
        // A turn the server started by itself (a queued prompt, a
        // background result) is a run of its own.
        if (!this.runs.has(session)) {
          this.runs.set(session, newId("opencode-v2-run"));
          this.emitFor(session, { event: "run_started" });
        }
        return;
      case "session.execution.succeeded":
        return at.task ? this.settleTask(session, "completed") : this.finish(session, outcome.completed());
      case "session.execution.failed":
        if (at.task) return this.settleTask(session, "failed");
        this.notice(session, "error", event.data.error.message);
        return this.finish(session, outcome.failed(event.data.error.message));
      case "session.execution.interrupted":
        return at.task ? this.settleTask(session, "cancelled") : this.finish(session, outcome.cancelled());
      case "session.inbox.delivered":
      case "session.input.promoted": {
        const id = event.type === "session.inbox.delivered" ? event.data.inboxID : event.data.inputID;
        const input = this.inputs.get(id);
        if (!input || input.session !== session || input.consumed) return;
        input.consumed = true;
        // An admitted steer can fall through into a native continuation.
        // Promotion is evidence for its input, not authority to revive a run.
        this.emit({
          sessionId: session,
          runId: this.runs.get(session) ?? input.runId,
          event: "input_consumed",
          inputId: input.inputId,
          nativeInputId: id,
        });
        return;
      }
      case "session.step.started":
        if (event.data.model) this.ranOn.set(session, event.data.model);
        return;
      case "session.step.ended":
        if (event.data.tokens) this.steps.set(session, event.data.tokens);
        return;
      case "session.step.failed":
        return this.notice(session, "error", event.data.error.message);
      case "session.retry.scheduled":
        return this.notice(session, "warning", `Retrying (attempt ${event.data.attempt}): ${event.data.error.message}`);
      case "session.usage.updated": {
        const tokens = this.steps.get(session) ?? event.data.tokens;
        const model = this.ranOn.get(session);
        const catalog = this.catalogs.get(this.sessions.get(at.root)?.workspace ?? "");
        const window = model && catalog ? contextWindow(catalog, model) : null;
        const usage = map.usage(tokens, event.data.cost, window);
        if (at.task) return this.updateTask(session, { usage });
        return this.emitFor(session, { event: "usage", ...usage });
      }
      case "session.text.delta":
      case "session.reasoning.delta": {
        const kind = event.type === "session.text.delta" ? "text" : "reasoning";
        const itemId = map.blockItemId(event.data.assistantMessageID, kind, event.data.ordinal);
        const text = event.data.delta;
        return this.emitFor(session, {
          event: kind === "text" ? "text_delta" : "reasoning_delta",
          itemId,
          text,
          mode: "append",
        });
      }
      case "session.text.ended":
      case "session.reasoning.ended": {
        // The whole block, so a delta lost to a reconnect does not stay lost.
        const kind = event.type === "session.text.ended" ? "text" : "reasoning";
        const itemId = map.blockItemId(event.data.assistantMessageID, kind, event.data.ordinal);
        const text = event.data.text;
        return this.emitFor(session, {
          event: kind === "text" ? "text_delta" : "reasoning_delta",
          itemId,
          text,
          mode: "replace",
        });
      }
      case "session.tool.input.started": {
        const id = map.toolItemId(event.data.assistantMessageID, event.data.id);
        this.calls.set(id, { session, name: event.data.name, input: {} });
        this.started.add(id);
        return this.emitFor(session, {
          event: "tool_call_started",
          id,
          name: event.data.name,
          kind: map.toolKind(event.data.name),
          title: map.toolTitle(event.data.name, {}),
          status: "pending",
        });
      }
      case "session.tool.called":
        return this.onToolCalled(session, at.task, event.data);
      case "session.tool.progress": {
        const id = map.toolItemId(event.data.assistantMessageID, event.data.id);
        const child = event.data.metadata["sessionID"];
        if (map.toolKind(this.calls.get(id)?.name ?? "") !== "task" || typeof child !== "string") return;
        // A call that continues an earlier subagent creates no session: its
        // session is learned here.
        if (!this.children.has(child)) yield* this.adopt(child).pipe(Effect.ignore);
        const known = this.children.get(child);
        if (!known) return;
        this.unlinked.set(
          session,
          (this.unlinked.get(session) ?? []).filter((call) => call !== id),
        );
        this.updateTask(child, { toolCallId: id, status: "running" });
        return;
      }
      case "session.tool.success": {
        const id = map.toolItemId(event.data.assistantMessageID, event.data.id);
        const call = this.calls.get(id) ?? { session, name: "tool", input: {} };
        const content = map.toolContent(call.name, call.input, event.data.content, event.data.metadata);
        this.emitFor(session, { event: "tool_call_updated", id, status: "completed", content });
        const child = event.data.metadata?.["sessionID"];
        // A background subagent's call answers at once; its result comes
        // later as a synthetic message.
        if (
          map.toolKind(call.name) === "task" &&
          typeof child === "string" &&
          event.data.metadata?.["status"] === "completed"
        )
          this.updateTask(child, { summary: map.subagentResult(map.resultText(event.data.content)) });
        return;
      }
      case "session.synthetic": {
        const metadata = event.data.metadata ?? {};
        const child = metadata["childID"];
        if (metadata["source"] !== "subagent" || typeof child !== "string") return;
        const state = metadata["state"];
        const status = state === "completed" ? "completed" : state === "cancelled" ? "cancelled" : "failed";
        this.updateTask(child, { summary: map.subagentResult(event.data.text), status, endedAt: now() });
        return;
      }
      case "session.tool.failed": {
        const id = map.toolItemId(event.data.assistantMessageID, event.data.id);
        const text = map.resultText(event.data.content ?? []) || event.data.error.message;
        return this.emitFor(session, {
          event: "tool_call_updated",
          id,
          status: "failed",
          content: [{ type: "text", text }],
        });
      }
      case "session.compaction.ended":
        return this.emitFor(session, { event: "compacted", ...(event.data.text ? { summary: event.data.text } : {}) });
      case "session.compaction.failed":
        return this.notice(session, "error", `Compacting failed: ${event.data.error.message}`);
      case "permission.asked":
        return yield* this.onPermission(event.data);
      case "permission.replied": {
        if (at.task) this.settleTask(session, "running");
        if (!this.approvals.delete(event.data.requestID)) return;
        return this.emitFor(session, { event: "approval_resolved", id: event.data.requestID });
      }
      case "form.created":
        return this.onForm(event.data.form);
      case "form.replied":
      case "form.cancelled": {
        if (at.task) this.settleTask(session, "running");
        if (!this.forms.delete(event.data.id)) return;
        return this.emitFor(session, { event: "question_resolved", id: event.data.id });
      }
    }
  });

  onToolCalled(session: string, task: string | null, data: api.EventData<"session.tool.called">): void {
    const id = map.toolItemId(data.assistantMessageID, data.id);
    const name = this.calls.get(id)?.name ?? "tool";
    this.calls.set(id, { session, name, input: data.input });
    const fields = {
      id,
      kind: map.toolKind(name),
      title: map.toolTitle(name, data.input),
      status: "running" as const,
      input: data.input,
      locations: map.toolLocations(name, data.input),
    };
    if (this.started.has(id)) this.emitFor(session, { event: "tool_call_updated", ...fields });
    else {
      this.started.add(id);
      this.emitFor(session, { event: "tool_call_started", name, ...fields });
    }
    if (map.toolKind(name) === "task") this.unlinked.set(session, [...(this.unlinked.get(session) ?? []), id]);
    const child = task ? this.children.get(task) : null;
    if (task && child) this.updateTask(task, { toolUses: (child.task.toolUses ?? 0) + 1, activity: fields.title });
  }

  onPermission = Effect.fn("OpenCodeV2.onPermission")(function* (
    this: OpenCodeV2Agent,
    request: api.PermissionRequest,
  ) {
    const at = this.route(request.sessionID);
    if (!at) return;
    const mode = permissionMode.selected(this.sessions.get(at.root)?.options);
    if (repliesAutomatically(mode)) {
      yield* this.send("POST", `/api/session/${request.sessionID}/permission/${request.id}/reply`, {
        body: { decision: "once" },
      });
      return;
    }
    this.approvals.set(request.id, request.sessionID);
    if (at.task) this.settleTask(request.sessionID, "waiting");
    const source = request.source;
    const callId = source ? map.toolItemId(source.messageID, source.id) : null;
    const call = callId ? this.calls.get(callId) : null;
    this.emitFor(request.sessionID, {
      event: "approval",
      id: request.id,
      title: approvalTitle(request),
      options: map.approvalOptions(request),
      ...(callId && call
        ? {
            toolCall: {
              id: callId,
              name: call.name,
              kind: map.toolKind(call.name),
              title: map.toolTitle(call.name, call.input),
              input: call.input,
              content: map.toolContent(call.name, call.input, [], request.metadata),
            },
          }
        : {}),
    });
  });

  onForm(form: Form): void {
    const at = this.route(form.sessionID);
    if (!at) return;
    this.forms.set(form.id, form);
    if (at.task) this.settleTask(form.sessionID, "waiting");
    this.emitFor(form.sessionID, { event: "question", ...map.formQuestion(form) });
  }

  // --- catalog and tools ------------------------------------------------------------------

  /// The catalog of a folder. A server builds it in the background after it
  /// starts, so a catalog with no usable model is read again (sooner when a
  /// catalog event arrives) for a while before it counts as empty.
  catalog = Effect.fn("OpenCodeV2.catalog")(function* (this: OpenCodeV2Agent, workspace: string) {
    const cached = this.catalogs.get(workspace);
    if (cached) return cached;
    const deadline = Date.now() + CATALOG_WAIT;
    yield* this.ensureStream().pipe(Effect.ignore);
    for (;;) {
      const signal = this.catalogSignal;
      const catalog = yield* this.readCatalog(workspace);
      if (usableModels(catalog).length || Date.now() >= deadline) {
        if (usableModels(catalog).length) this.catalogs.set(workspace, catalog);
        return catalog;
      }
      yield* Deferred.await(signal).pipe(Effect.timeout(CATALOG_POLL), Effect.ignore);
    }
  });

  readCatalog = Effect.fn("OpenCodeV2.readCatalog")(function* (this: OpenCodeV2Agent, workspace: string) {
    const query = atLocation(workspace);
    const [providers, models, defaultModel, agents] = yield* Effect.all(
      [
        this.call("GET", "/api/provider", api.listed(api.provider), { query }),
        this.call("GET", "/api/model", api.listed(api.model), { query }),
        // A server with no default model answers an error; the first model
        // is the default then.
        this.call("GET", "/api/model/default", api.single(api.modelRef), { query }).pipe(
          Effect.map((answer) => answer.data),
          Effect.orElseSucceed(() => null),
        ),
        this.call("GET", "/api/agent", api.listed(api.agent), { query }),
      ],
      { concurrency: "unbounded" },
    );
    const catalog: Catalog = { providers: providers.data, models: models.data, defaultModel, agents: agents.data };
    return catalog;
  });

  /// A catalog changed on the server: the cached one is dropped, whoever
  /// waits for one reads again, and the chats in that folder get the new
  /// options.
  onCatalogChange = Effect.fn("OpenCodeV2.onCatalogChange")(function* (
    this: OpenCodeV2Agent,
    directory: string | null,
  ) {
    const workspaces = directory ? [directory] : [...this.catalogs.keys()];
    for (const workspace of workspaces) {
      this.catalogs.delete(workspace);
      this.commands.delete(workspace);
    }
    if (!directory) {
      this.catalogs.clear();
      this.commands.clear();
    }
    const signal = this.catalogSignal;
    this.catalogSignal = Deferred.makeUnsafe<void>();
    yield* Deferred.succeed(signal, undefined);
    const affected = new Set(
      [...this.sessions.values()]
        .map((state) => state.workspace)
        .filter((workspace) => !directory || workspace === directory),
    );
    for (const workspace of affected) yield* this.refreshOptions(workspace).pipe(Effect.forkIn(this.scope));
  });

  /// Sends the chats of a folder their options again, from a fresh catalog.
  refreshOptions = Effect.fn("OpenCodeV2.refreshOptions")(function* (this: OpenCodeV2Agent, workspace: string) {
    const catalog = yield* this.readCatalog(workspace);
    if (!usableModels(catalog).length) return;
    this.catalogs.set(workspace, catalog);
    for (const [session, state] of this.sessions)
      if (state.workspace === workspace)
        this.emitFor(session, { event: "config_options", options: buildOptions(catalog, state.options) });
  });

  commandsFor = Effect.fn("OpenCodeV2.commandsFor")(function* (this: OpenCodeV2Agent, workspace: string) {
    const cached = this.commands.get(workspace);
    if (cached) return cached;
    const listed = yield* this.call("GET", "/api/command", api.listed(api.command), { query: atLocation(workspace) });
    this.commands.set(workspace, listed.data);
    return listed.data;
  });

  /// Gives OpenCode the session's plugin tools: the host serves them over
  /// its loopback MCP server and OpenCode is told the address for the
  /// workspace. Added again when the tools change or the server restarted,
  /// because OpenCode lists a server's tools once, when it connects.
  offerTools = Effect.fn("OpenCodeV2.offerTools")(function* (this: OpenCodeV2Agent, session: string) {
    const state = yield* this.known(session);
    const known = this.mcp.get(state.workspace);
    if (!state.tools.length && !known) return;
    const base = yield* this.base();
    if (!isLoopback(base)) {
      if (state.tools.length)
        this.notice(session, "warning", "Plugin tools are offered only to an OpenCode server on this computer.");
      return;
    }
    const key = JSON.stringify(state.tools);
    if (known && known.key === key && known.server === this.server) return;
    const host = yield* Host;
    const served = yield* host.call("host/mcp.serve", {
      agentId: AGENT_ID,
      workspace: state.workspace,
      tools: state.tools,
      ...(known ? { url: known.url } : {}),
    });
    const query = atLocation(state.workspace);
    yield* this.send("PUT", `/api/experimental/mcp/${MCP_SERVER}`, {
      query,
      body: { config: { type: "remote", url: served.url, codemode: false } },
    });
    const servers = yield* this.call("GET", "/api/mcp", api.mcpServers, { query });
    const ours = servers.data.find((entry) => entry.name === MCP_SERVER);
    if (ours && ours.status.status !== "connected")
      return yield* new ServerError({
        message: `OpenCode could not connect to the Divergence tools server (${ours.status.status}${ours.status.error ? `: ${ours.status.error}` : ""})`,
      });
    this.mcp.set(state.workspace, { url: served.url, key, server: this.server });
  });

  /// Brings the server's session in line with the chat: model, mode,
  /// permission rules and plugin instructions, each only when it changed.
  applyContext = Effect.fn("OpenCodeV2.applyContext")(function* (this: OpenCodeV2Agent, session: string) {
    const state = yield* this.known(session);
    const catalog = yield* this.catalog(state.workspace);
    const path = `/api/session/${session}`;
    const model = chosenModel(catalog, state.options);
    if (model && modelKey(model) !== state.applied.model) {
      yield* this.send("POST", `${path}/model`, { body: { model } });
      state.applied.model = modelKey(model);
    }
    const agent = state.options[OPTION_AGENT];
    if (typeof agent === "string" && agent !== state.applied.agent) {
      yield* this.send("POST", `${path}/agent`, { body: { agent } });
      state.applied.agent = agent;
    }
    yield* this.applyRules(session);
    const instructions = state.instructions?.trim() || null;
    if (instructions !== state.applied.instructions) {
      const entry = `/api/experimental/session/${session}/instructions/entries/${INSTRUCTIONS_KEY}`;
      if (instructions) yield* this.send("PUT", entry, { body: { value: instructions } });
      else yield* this.send("DELETE", entry);
      state.applied.instructions = instructions;
    }
  });

  applyRules = Effect.fn("OpenCodeV2.applyRules")(function* (this: OpenCodeV2Agent, session: string) {
    const state = yield* this.known(session);
    const mode = permissionMode.selected(state.options);
    if (mode === state.applied.mode) return;
    yield* this.send("PATCH", `/api/session/${session}`, { body: { permissions: ruleset(mode) } });
    state.applied.mode = mode;
  });

  // --- agent methods --------------------------------------------------------------------

  initialize = Effect.fn("OpenCodeV2.initialize")(function* (this: OpenCodeV2Agent) {
    const external = this.external;
    const info: AgentInfo = {
      id: AGENT_ID,
      name: "OpenCode",
      description: external
        ? `The OpenCode 2 coding agent, on the server at ${external}.`
        : "The OpenCode 2 coding agent, running as a local HTTP server.",
      icon: ICON,
      capabilities: {
        reasoning: true,
        images: true,
        approvals: true,
        questions: true,
        sessionList: true,
        sessionHistory: true,
        resume: true,
        slashCommands: true,
        steer: true,
        cancel: true,
        skills: true,
        fork: true,
        compact: true,
        subagents: true,
        cancelTask: true,
      },
      authMethods: [],
      status: { state: "ready" },
      family: FAMILY,
      // Sessions live on the server: two agents on one server can continue
      // each other's chats.
      continuationKey: external ?? "",
    };
    const reached = yield* this.call("GET", "/api/info", api.info).pipe(Effect.result);
    if (Result.isFailure(reached)) {
      const error = reached.failure;
      const unauthorized = error instanceof ServerError && error.status === 401;
      info.status = unauthorized
        ? { state: "auth_required", message: error.message }
        : { state: "unavailable", message: errorMessage(error) };
      return info;
    }
    info.version = reached.success.version;
    return info;
  });

  listOptions = Effect.fn("OpenCodeV2.listOptions")(function* (this: OpenCodeV2Agent, workspace: string) {
    return { options: buildOptions(yield* this.catalog(workspace)) };
  });

  listCommands = Effect.fn("OpenCodeV2.listCommands")(function* (this: OpenCodeV2Agent, session: string) {
    const listed = yield* this.commandsFor((yield* this.known(session)).workspace);
    const commands: SlashCommand[] = listed.map((command) => ({
      name: command.name,
      ...(command.description ? { description: command.description } : {}),
      source: "agent",
    }));
    return { commands };
  });

  listSkills = Effect.fn("OpenCodeV2.listSkills")(function* (this: OpenCodeV2Agent, workspace: string) {
    const listed = yield* this.call("GET", "/api/skill", api.listed(api.skill), { query: atLocation(workspace) });
    const skills: Skill[] = listed.data.map((skill) => ({
      name: skill.id,
      ...(skill.description ? { description: skill.description } : {}),
    }));
    return { skills };
  });

  listSessions = Effect.fn("OpenCodeV2.listSessions")(function* (this: OpenCodeV2Agent, workspace: string) {
    const page = yield* this.call("GET", "/api/session", api.sessionPage, {
      // `null` asks for root sessions only: subagent sessions belong to the
      // chat that started them.
      query: { directory: workspace, limit: PAGE, parentID: "null" },
    });
    const sessions: SessionSummary[] = page.data
      .filter((session) => !session.parentID)
      .sort((a, b) => b.time.updated - a.time.updated)
      .map((session) => ({
        id: session.id,
        ...(session.title ? { title: session.title } : {}),
        createdAt: map.isoTime(session.time.created),
        updatedAt: map.isoTime(session.time.updated),
      }));
    return { sessions };
  });

  /// Every message of a session, oldest first.
  messages = Effect.fn("OpenCodeV2.messages")(function* (this: OpenCodeV2Agent, session: string) {
    const all: api.Message[] = [];
    // The first page sets the order; later pages follow the cursor, which
    // the server refuses together with an order.
    let query: RequestOptions["query"] = { order: "asc", limit: PAGE };
    for (;;) {
      const page: z.infer<typeof api.messagePage> = yield* this.call(
        "GET",
        `/api/session/${session}/message`,
        api.messagePage,
        { query },
      );
      all.push(...page.data);
      const next: string | null | undefined = page.cursor?.next;
      if (!next || page.data.length < PAGE) return all;
      query = { cursor: next, limit: PAGE };
    }
  });

  /// The transcript of a session, its subagents nested under the calls
  /// that started them.
  transcriptOf: (
    session: string,
    depth: number,
    parent: string | null,
  ) => Effect.Effect<TranscriptItem[], unknown, PluginServices> = Effect.fn("OpenCodeV2.transcriptOf")(function* (
    this: OpenCodeV2Agent,
    session: string,
    depth: number,
    parent: string | null,
  ) {
    const messages = yield* this.messages(session);
    const subagents = new Map<string, map.Subagent>();
    if (depth < MAX_SUBAGENT_DEPTH) {
      for (const message of messages) {
        const assistant = api.assistantMessage.safeParse(message);
        if (!assistant.success) continue;
        for (const entry of assistant.data.content) {
          if (entry.type !== "tool" || map.toolKind(entry.name) !== "task" || entry.state.status === "streaming")
            continue;
          const child = entry.state.metadata?.["sessionID"];
          if (typeof child !== "string") continue;
          const items = yield* this.transcriptOf(child, depth + 1, child).pipe(Effect.orElseSucceed(() => []));
          const status = entry.state.status;
          const input = entry.state.input;
          const description = input["description"];
          const prompt = input["prompt"];
          const name = input["agent"];
          const result = entry.state.status === "completed" ? map.resultText(entry.state.content) : "";
          const id = map.toolItemId(assistant.data.id, entry.id);
          subagents.set(id, {
            task: {
              id: child,
              title: typeof description === "string" && description ? description : "Subagent",
              status: status === "completed" ? "completed" : status === "error" ? "failed" : "running",
              toolCallId: id,
              ...(parent ? { parentTaskId: parent } : {}),
              ...(typeof name === "string" ? { name } : {}),
              ...(typeof prompt === "string" ? { prompt } : {}),
              ...(result ? { summary: map.subagentResult(result) } : {}),
              toolUses: items.filter((item) => item.role === "tool").length,
            },
            items,
          });
        }
      }
    }
    return map.transcript(messages, (messageId, callId) => subagents.get(map.toolItemId(messageId, callId)) ?? null);
  });

  readSession = Effect.fn("OpenCodeV2.readSession")(function* (this: OpenCodeV2Agent, session: string) {
    return { items: yield* this.transcriptOf(session, 0, null) };
  });

  remember(session: string, params: SessionOptions, applied: SessionState["applied"]): SessionState {
    const state: SessionState = {
      workspace: params.workspace,
      options: { ...params.options },
      // The host leaves the list out when there are none.
      tools: params.tools ?? [],
      instructions: params.instructions ?? null,
      applied,
    };
    this.sessions.set(session, state);
    this.strangers.delete(session);
    return state;
  }

  /// Offers the session's tools, telling the chat when they cannot reach
  /// OpenCode rather than failing the session over them.
  offerToolsOrWarn = Effect.fn("OpenCodeV2.offerToolsOrWarn")(function* (this: OpenCodeV2Agent, session: string) {
    const offered = yield* this.offerTools(session).pipe(Effect.result);
    if (Result.isFailure(offered))
      this.notice(session, "warning", `Plugin tools are not available to OpenCode: ${errorMessage(offered.failure)}`);
  });

  createSession = Effect.fn("OpenCodeV2.createSession")(function* (this: OpenCodeV2Agent, params: SessionOptions) {
    const catalog = yield* this.catalog(params.workspace);
    const model = chosenModel(catalog, params.options);
    const agent = params.options[OPTION_AGENT];
    const mode = permissionMode.selected(params.options);
    const created = yield* this.call("POST", "/api/session", api.single(api.session), {
      body: {
        location: { directory: params.workspace },
        ...(typeof agent === "string" ? { agent } : {}),
        ...(model ? { model } : {}),
        // The rules the chat runs under, not whatever the agent allows.
        permissions: ruleset(mode),
      },
    });
    const id = created.data.id;
    this.remember(id, params, {
      model: modelKey(model),
      agent: typeof agent === "string" ? agent : null,
      mode,
      instructions: null,
    });
    yield* this.ensureStream();
    yield* this.offerToolsOrWarn(id);
    return { sessionId: id };
  });

  resumeSession = Effect.fn("OpenCodeV2.resumeSession")(function* (
    this: OpenCodeV2Agent,
    params: SessionOptions & { sessionId: string },
  ) {
    const known = yield* this.call("GET", `/api/session/${params.sessionId}`, api.single(api.session)).pipe(
      Effect.mapError(
        (error) =>
          new ServerError({ message: `OpenCode does not know session ${params.sessionId}: ${errorMessage(error)}` }),
      ),
    );
    this.remember(params.sessionId, params, {
      model: modelKey(known.data.model ?? null),
      agent: known.data.agent ?? null,
      // The rules and instructions are sent again: the session may predate
      // the mode the user has since chosen.
      mode: null,
      instructions: null,
    });
    yield* this.ensureStream();
    yield* this.offerToolsOrWarn(params.sessionId);
    return {};
  });

  forkSession = Effect.fn("OpenCodeV2.forkSession")(function* (
    this: OpenCodeV2Agent,
    params: SessionOptions & { sessionId: string },
  ) {
    const forked = yield* this.call("POST", `/api/session/${params.sessionId}/fork`, api.single(api.session), {
      body: {},
    });
    const id = forked.data.id;
    this.remember(id, params, {
      model: modelKey(forked.data.model ?? null),
      agent: forked.data.agent ?? null,
      mode: null,
      instructions: null,
    });
    yield* this.ensureStream();
    yield* this.offerToolsOrWarn(id);
    return { sessionId: id };
  });

  closeSession = Effect.fn("OpenCodeV2.closeSession")(function* (this: OpenCodeV2Agent, session: string) {
    this.sessions.delete(session);
    this.promptLocks.delete(session);
    for (const [id, input] of this.inputs) if (input.session === session) this.inputs.delete(id);
    for (const [id, child] of [...this.children]) if (child.root === session) this.children.delete(id);
    return {};
  });

  prompt = Effect.fn("OpenCodeV2.prompt")((session: string, input: PromptInput) => {
    let lock = this.promptLocks.get(session);
    if (!lock) {
      lock = Semaphore.makeUnsafe(1);
      this.promptLocks.set(session, lock);
    }
    return lock.withPermit(this.sendPrompt(session, input));
  });

  sendPrompt = Effect.fn("OpenCodeV2.sendPrompt")(function* (
    this: OpenCodeV2Agent,
    session: string,
    input: PromptInput,
  ) {
    const delivery = input.delivery;
    if (delivery?.intent === "queue") {
      const reason = "Queued input must remain held by the host";
      this.emit({
        sessionId: session,
        event: "input_rejected",
        inputId: delivery.inputId,
        attemptId: delivery.attemptId,
        reason,
      });
      return yield* new ServerError({ message: reason });
    }
    const state = yield* this.known(session);
    const body = yield* Effect.try({ try: () => promptBody(input.blocks), catch: (error) => error });
    yield* this.ensureStream();
    // A server started again since the session began forgot the tools.
    yield* this.offerToolsOrWarn(session);
    const path = `/api/session/${session}`;
    const active = this.runs.get(session);
    if (!active) yield* this.applyContext(session);
    const commands = yield* this.commandsFor(state.workspace);
    const slash = slashCommand(body, commands);
    const runId = active ?? newId("opencode-v2-run");
    this.runs.set(session, runId);
    // Installed command responses have no input identity/receipt. Keep the
    // legacy path weak rather than misreading them as prompt admission.
    if (slash) {
      yield* this.runCommand(session, body, slash, active ? "steer" : undefined);
      return { runId };
    }
    const payload = { ...this.promptPayload(body, input), ...(active ? { delivery: "steer" } : {}) };
    const nativeId = payload.id;
    if (delivery && nativeId) this.inputs.set(nativeId, { session, ...delivery, runId, consumed: false });
    const sent = yield* this.call("POST", `${path}/prompt`, api.single(api.admittedInput), { body: payload }).pipe(
      Effect.result,
    );
    if (Result.isFailure(sent)) {
      // An active steer failure does not finish/replace its source execution.
      if (!active && this.runs.get(session) === runId) this.runs.delete(session);
      const error = sent.failure;
      if (
        delivery &&
        error instanceof ServerError &&
        ((error.status === 400 && error.refusalKind === "InvalidRequestError") ||
          (error.status === 404 && error.refusalKind === "SessionNotFoundError") ||
          (error.status === 401 && error.refusalKind === "UnauthorizedError"))
      ) {
        if (nativeId) this.inputs.delete(nativeId);
        this.emit({
          sessionId: session,
          runId,
          event: "input_rejected",
          inputId: delivery.inputId,
          attemptId: delivery.attemptId,
          reason: error.message,
        });
      }
      return yield* Effect.fail(sent.failure);
    }
    const admission = sent.success.data;
    if (admission.sessionID !== session || (nativeId && admission.id !== nativeId))
      return yield* new ServerError({
        message: "OpenCode answered with an unrelated admission identity; delivery is uncertain",
      });
    if (delivery) {
      const tracked = this.inputs.get(admission.id);
      if (tracked) tracked.admission = admission;
    }
    return { runId, receipt: { evidence: "native_admission" as const, nativeInputId: admission.id } };
  });

  promptPayload(body: PromptBody, input: PromptInput) {
    // Naming the message after the host's item keeps the two ids one.
    const id = input.delivery?.inputId ?? input.itemId;
    return { ...body, ...(id ? { id: messageIdFor(id) } : {}) };
  }

  /// A slash command. Its call may answer only once the command's turn is
  /// over, so it runs on its own; a failure ends the run with the reason.
  runCommand = Effect.fn("OpenCodeV2.runCommand")(function* (
    this: OpenCodeV2Agent,
    session: string,
    body: PromptBody,
    slash: { name: string; text: string },
    delivery?: "steer",
  ) {
    yield* this.send("POST", `/api/session/${session}/command`, {
      body: {
        name: slash.name,
        text: slash.text,
        ...(body.files ? { files: body.files } : {}),
        ...(delivery ? { delivery } : {}),
      },
    }).pipe(
      Effect.catch((error) =>
        Effect.sync(() => {
          const message = errorMessage(error);
          this.notice(session, "error", message);
          this.finish(session, outcome.failed(message));
        }),
      ),
      Effect.forkIn(this.scope),
    );
  });

  /// Every subagent session under a chat or a task, deepest last.
  subtree(session: string): string[] {
    const found: string[] = [];
    for (const [id, child] of this.children) if (child.parent === session) found.push(id, ...this.subtree(id));
    return found;
  }

  /// Rejects what the given sessions wait on, so nothing stays blocked on
  /// an answer nobody will give.
  rejectPending = Effect.fn("OpenCodeV2.rejectPending")(function* (this: OpenCodeV2Agent, sessions: Set<string>) {
    for (const [id, session] of [...this.approvals]) {
      if (!sessions.has(session)) continue;
      this.approvals.delete(id);
      this.emitFor(session, { event: "approval_resolved", id });
      yield* this.send("POST", `/api/session/${session}/permission/${id}/reply`, { body: { decision: "reject" } }).pipe(
        Effect.ignore,
      );
    }
    for (const [id, form] of [...this.forms]) {
      if (!sessions.has(form.sessionID)) continue;
      this.forms.delete(id);
      this.emitFor(form.sessionID, { event: "question_resolved", id });
      yield* this.send("DELETE", `/api/session/${form.sessionID}/form/${id}`).pipe(Effect.ignore);
    }
  });

  cancel = Effect.fn("OpenCodeV2.cancel")(function* (this: OpenCodeV2Agent, session: string) {
    yield* this.known(session);
    // Its subagents are stopped by the user, not failing: marked first.
    const subagents = this.subtree(session);
    for (const child of subagents) this.settleTask(child, "cancelled");
    yield* this.rejectPending(new Set([session, ...subagents]));
    const interrupted = yield* this.send("POST", `/api/session/${session}/interrupt`).pipe(Effect.result);
    for (const child of subagents) yield* this.send("POST", `/api/session/${child}/interrupt`).pipe(Effect.ignore);
    this.finish(session, outcome.cancelled());
    if (Result.isFailure(interrupted)) return yield* Effect.fail(interrupted.failure);
    return {};
  });

  cancelTask = Effect.fn("OpenCodeV2.cancelTask")(function* (this: OpenCodeV2Agent, session: string, task: string) {
    if (this.children.get(task)?.root !== session)
      return yield* new ServerError({ message: `no OpenCode subagent ${task} in session ${session}` });
    const tree = [task, ...this.subtree(task)];
    for (const child of tree) this.settleTask(child, "cancelled");
    yield* this.rejectPending(new Set(tree));
    // The parent's subagent call then fails and the parent turn goes on.
    for (const child of tree) yield* this.send("POST", `/api/session/${child}/interrupt`);
    return {};
  });

  setOption = Effect.fn("OpenCodeV2.setOption")(function* (
    this: OpenCodeV2Agent,
    session: string,
    optionId: string,
    value: unknown,
  ) {
    const state = yield* this.known(session);
    // Variants belong to a model, so a model change drops the choice.
    if (optionId === OPTION_MODEL) delete state.options[OPTION_VARIANT];
    if (value === null || value === undefined) delete state.options[optionId];
    else state.options[optionId] = value;
    // The model and mode switch before the next prompt; the permission
    // rules apply at once, to a run that is waiting on one.
    if (optionId === permissionMode.OPTION_ID) yield* this.applyRules(session);
    const options: ConfigOption[] = buildOptions(yield* this.catalog(state.workspace), state.options);
    return { options };
  });

  respondToApproval = Effect.fn("OpenCodeV2.respondToApproval")(function* (
    this: OpenCodeV2Agent,
    approval: string,
    option: string,
  ) {
    const session = this.approvals.get(approval);
    if (!session) return yield* new ServerError({ message: `no pending OpenCode approval ${approval}` });
    const decision = z.enum(["once", "always", "reject"]).safeParse(option);
    if (!decision.success) return yield* new ServerError({ message: `unknown approval choice ${option}` });
    yield* this.send("POST", `/api/session/${session}/permission/${approval}/reply`, {
      body: { decision: decision.data },
    });
    if (this.approvals.delete(approval)) this.emitFor(session, { event: "approval_resolved", id: approval });
    return {};
  });

  respondToQuestion = Effect.fn("OpenCodeV2.respondToQuestion")(function* (
    this: OpenCodeV2Agent,
    question: string,
    answer: QuestionAnswer,
  ) {
    const form = this.forms.get(question);
    if (!form) return yield* new ServerError({ message: `no pending OpenCode question ${question}` });
    const path = `/api/session/${form.sessionID}/form/${question}`;
    const values = map.formAnswer(form, answer.values);
    if (answer.cancelled || !Object.keys(values).length) yield* this.send("DELETE", path);
    else yield* this.send("POST", `${path}/reply`, { body: { answer: values } });
    if (this.forms.delete(question)) this.emitFor(form.sessionID, { event: "question_resolved", id: question });
    return {};
  });

  compact = Effect.fn("OpenCodeV2.compact")(function* (this: OpenCodeV2Agent, session: string) {
    yield* this.known(session);
    yield* this.ensureStream();
    yield* this.call("POST", `/api/session/${session}/compact`, z.unknown(), { body: {} });
    return {};
  });
}

Versions

VersionPublishedPlugin APISizePermissionsStatus
0.1.0latestOct 5, 2026>=2 <345.1 KB4 permissionsListed

Reviews and comments

0 threads · 0 reviews

No comments yet.