Official

chat

Built-in chat UI (JavaScript UI plugin).

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/chat@0.2.0

Permissions in 0.2.0

  • Control chats chats.controlHighCreates chats, sends prompts and cancels runs.Send prompts, answer agent questions, and create, change, fork, revert, stop and group chats
  • Run other plugins' commands commands.runHighRuns commands of other plugins.Run the sidebar's New chat command from the title bar
  • Read chats chats.readMediumReads your transcripts and chat lists.Show workspaces, tasks, chats and transcripts, and follow them as agents work
  • Read files fs.readMediumReads files in the listed places.Suggest workspace files for @ mentionsPlaces: the open workspace
  • Show panels ui.slotsLowShows views in the listed parts of the window.Draw the chat grid and the title barSlots: centertitle

Files

scheduler.ts21.7 KB
// Scoped orchestration for the stock scheduler. Only host snapshots own inputs;
// timers merely reconsider durable intent and every automatic claim carries the
// currently selected service token. No provider-specific recovery branches.
import * as Effect from "effect/Effect";
import * as Semaphore from "effect/Semaphore";
import * as z from "zod";
import { promise, parse } from "convergence/effect";
import type { ContentBlock, HostMethods, InputIntent, InputSnapshot, UsageRecovery } from "convergence";
import type { Controller } from "./controller.ts";
import { errorMessage } from "./data.ts";
import {
  autoSend,
  defaults,
  editable,
  fifoCandidate,
  firstQueued,
  fresh,
  matching,
  orderedWaits,
  pending,
  policyOf,
  policySchema,
  preferencesSchema,
  recoveryDecision,
  scopeOf,
  settlePolicy,
  uncertain,
  waitKey,
} from "./scheduler-policy.ts";
import type { PolicyState, Preferences } from "./scheduler-policy.ts";

export function unknownMethod(error: unknown): boolean {
  return (
    typeof error === "object" &&
    error !== null &&
    "_tag" in error &&
    error._tag === "HostCallFailed" &&
    (("code" in error && error.code === -32601) ||
      ("message" in error &&
        typeof error.message === "string" &&
        /^unknown host method host\/(inputs\.|agents\.refresh_usage)/i.test(error.message)))
  );
}
export class Scheduler {
  readonly controller: Controller;
  readonly snapshots: Record<string, InputSnapshot> = {};
  readonly errors: Record<string, string> = {};
  readonly checking = new Set<string>();
  readonly restored = new Set<string>();
  readonly manualStops = new Set<string>();
  readonly delivering = new Set<string>();
  readonly activeLimits: Record<string, { runId?: string; evidence: UsageRecovery }> = {};
  preferences: Preferences = defaults;
  available: boolean | undefined;
  ready = false;
  private readonly serial = Semaphore.makeUnsafe(1);
  // Serialize policy scans separately from user mutations. Provider admission
  // and availability reads may take arbitrarily long; host CAS and the live
  // policy token still fence every automatic claim against intervening edits.
  private readonly evaluation = Semaphore.makeUnsafe(1);
  constructor(controller: Controller) {
    this.controller = controller;
  }
  receive(snapshot: InputSnapshot) {
    const previous = this.snapshots[snapshot.chatId];
    if (
      previous &&
      (previous.generation > snapshot.generation ||
        (previous.generation === snapshot.generation && previous.revision > snapshot.revision))
    )
      return;
    this.snapshots[snapshot.chatId] = snapshot;
    const active = this.activeLimits[snapshot.chatId];
    if (active && ((!snapshot.running && !snapshot.settling) || (active.runId && active.runId !== snapshot.runId)))
      delete this.activeLimits[snapshot.chatId];
    this.controller.refresh(snapshot.chatId);
  }
  read = Effect.fn("Scheduler.read")((chatId: string) =>
    this.controller.host
      .call("host/inputs.list", { chatId })
      .pipe(Effect.tap((snapshot) => Effect.sync(() => this.receive(snapshot)))),
  );
  load = Effect.fn("Scheduler.load")(
    function* (this: Scheduler, chatId: string) {
      if (this.available === false) return;
      yield* this.read(chatId).pipe(
        Effect.tap(() =>
          Effect.sync(() => {
            this.available = true;
          }),
        ),
        Effect.catch((error) =>
          unknownMethod(error)
            ? Effect.sync(() => {
                this.available = false;
              })
            : this.report(chatId, error),
        ),
      );
    }.bind(this),
  );
  report = Effect.fn("Scheduler.report")((chatId: string, error: unknown) =>
    Effect.gen(
      function* (this: Scheduler) {
        const parsed = z
          .union([z.object({ message: z.string() }), z.object({ what: z.string(), issues: z.array(z.string()) })])
          .safeParse(error);
        const message = `Pending input: ${parsed.success ? errorMessage(parsed.data) : String(error)}`;
        if (this.errors[chatId] !== message) yield* this.controller.notify.error(message);
        this.errors[chatId] = message;
        this.controller.refresh(chatId);
      }.bind(this),
    ),
  );
  write = Effect.fn("Scheduler.write")((snapshot: InputSnapshot, state: PolicyState) =>
    this.controller.host
      .call("host/inputs.update", {
        chatId: snapshot.chatId,
        revision: snapshot.revision,
        policyState: state,
      })
      .pipe(Effect.tap((next) => Effect.sync(() => this.receive(next)))),
  );
  enqueue = Effect.fn("Scheduler.enqueue")(
    (chatId: string, id: string, blocks: readonly ContentBlock[], intent: InputIntent, origin = "user") =>
      this.serial.withPermit(
        Effect.gen(
          function* (this: Scheduler) {
            const before = yield* this.read(chatId);
            const state = policyOf(before);
            const initialized =
              intent === "queue" && !before.running && !before.settling && !state.wait
                ? { ...state, releasedRun: before.lastRun?.id }
                : state;
            if (
              !policySchema.safeParse(before.policyState).success ||
              JSON.stringify(initialized) !== JSON.stringify(before.policyState)
            )
              yield* this.write(before, initialized);
            const snapshot = yield* this.controller.host.call("host/inputs.enqueue", {
              chatId,
              id,
              blocks,
              intent,
              origin,
              append: intent === "steer",
            });
            this.receive(snapshot);
            delete this.errors[chatId];
            return snapshot;
          }.bind(this),
        ),
      ),
  );
  manualDispatch = Effect.fn("Scheduler.manualDispatch")(
    function* (this: Scheduler, chatId: string, inputId: string, revision: number) {
      // An amended reserved steer is already being dispatched. Appending never
      // requests a second cancellation during that cancel/settle handshake.
      if (this.delivering.has(inputId)) return yield* this.read(chatId);
      const prepared = yield* this.serial.withPermit(
        Effect.gen(
          function* (this: Scheduler) {
            const snapshot = this.snapshots[chatId];
            if (!snapshot || snapshot.revision !== revision) return revision;
            const state = policyOf(snapshot);
            const next = yield* this.controller.host.call("host/inputs.update", {
              chatId,
              revision,
              inputId,
              pauseReason: null,
              policyState: { ...state, pause: undefined, wait: undefined },
            });
            this.receive(next);
            this.manualStops.delete(chatId);
            return next.revision;
          }.bind(this),
        ),
      );
      this.delivering.add(inputId);
      // Do not hold the plugin semaphore across provider I/O: new contributions
      // must reach the host reservation before its replacement prompt is sent.
      const snapshot = yield* this.controller.host
        .call("host/inputs.dispatch", { chatId, inputId, revision: prepared })
        .pipe(Effect.ensuring(Effect.sync(() => this.delivering.delete(inputId))));
      this.receive(snapshot);
      delete this.errors[chatId];
      return snapshot;
    }.bind(this),
  );
  change = Effect.fn("Scheduler.change")((params: HostMethods["host/inputs.update"]["params"]) =>
    this.serial.withPermit(
      this.controller.host
        .call("host/inputs.update", params)
        .pipe(Effect.tap((snapshot) => Effect.sync(() => this.receive(snapshot)))),
    ),
  );
  promote = Effect.fn("Scheduler.promote")(
    function* (this: Scheduler, chatId: string, inputId: string, revision: number) {
      const snapshot = yield* this.change({ chatId, inputId, revision, intent: "steer", pauseReason: null });
      yield* this.manualDispatch(chatId, inputId, snapshot.revision);
    }.bind(this),
  );
  stop = Effect.fn("Scheduler.stop")(
    function* (this: Scheduler, chatId: string) {
      this.manualStops.add(chatId);
      // Stop is always available, even when durable policy persistence fails.
      const savePause = Effect.gen(
        function* (this: Scheduler) {
          if (this.available === false) return;
          let snapshot = yield* this.read(chatId),
            state = policyOf(snapshot);
          const first = yield* this.write(snapshot, {
            ...state,
            pause: "Stopped by you",
            wait: state.wait ? { ...state.wait, status: "paused", reason: "Stopped by you" } : undefined,
          }).pipe(Effect.result);
          if (first._tag === "Success") return;
          const latest = yield* this.read(chatId);
          // Cancellation itself changes the revision. Retry metadata only against
          // the observed new snapshot, never replay a delivery operation.
          if (latest.revision === snapshot.revision) return yield* Effect.fail(first.failure);
          snapshot = latest;
          state = policyOf(snapshot);
          yield* this.write(snapshot, {
            ...state,
            pause: "Stopped by you",
            wait: state.wait ? { ...state.wait, status: "paused", reason: "Stopped by you" } : undefined,
          });
        }.bind(this),
      ).pipe(Effect.catch((error) => (unknownMethod(error) ? Effect.void : this.report(chatId, error))));
      yield* Effect.all([savePause, this.controller.host.call("host/chats.cancel", { chatId })], {
        concurrency: "unbounded",
      });
    }.bind(this),
  );
  setAutoSend = Effect.fn("Scheduler.setAutoSend")(
    function* (this: Scheduler, agentId: string, enabled: boolean) {
      const preferences = yield* parse(
        "Chat automation preferences",
        preferencesSchema,
        this.controller.api.settings.get(),
      );
      yield* promise("settings.set", () =>
        this.controller.api.settings.set("autoSendAfterLimit", {
          ...preferences.autoSendAfterLimit,
          [agentId]: enabled,
        }),
      );
      yield* this.updatePreferences();
    }.bind(this),
  );
  updatePreferences = Effect.fn("Scheduler.preferences")(
    function* (this: Scheduler) {
      this.preferences = yield* parse(
        "Chat automation preferences",
        preferencesSchema,
        this.controller.api.settings.get(),
      );
      this.controller.settingsView?.update();
      for (const snapshot of Object.values(this.snapshots)) this.controller.refresh(snapshot.chatId);
      yield* this.evaluate();
    }.bind(this),
  );
  restore = Effect.fn("Scheduler.restore")(
    function* (this: Scheduler) {
      yield* this.updatePreferences();
      const { chats } = yield* this.controller.host.call("host/chats.list", {});
      // Restore every chat, not just the three retained panes.
      for (const chat of chats) {
        if (chat.archived) continue;
        yield* this.load(chat.id);
        const snapshot = this.snapshots[chat.id];
        if (snapshot && policyOf(snapshot).wait) this.restored.add(chat.id);
      }
      this.ready = true;
      yield* this.evaluate();
    }.bind(this),
  );
  dispatch = Effect.fn("Scheduler.dispatch")(
    function* (this: Scheduler, snapshot: InputSnapshot, inputId: string, scope?: string) {
      const { token } = yield* this.controller.host.call("host/inputs.policy", {});
      if (!token) return;
      const result = yield* this.controller.host.call("host/inputs.dispatch", {
        chatId: snapshot.chatId,
        inputId,
        revision: snapshot.revision,
        policy: token,
        scope,
      });
      this.receive(result);
      delete this.errors[snapshot.chatId];
    }.bind(this),
  );
  evaluate = Effect.fn("Scheduler.evaluate")(() => this.evaluation.withPermit(this.evaluateLocked()));
  private evaluateLocked = Effect.fn("Scheduler.evaluateLocked")(
    function* (this: Scheduler) {
      if (!this.ready || this.available !== true) return;
      for (const snapshot of Object.values(this.snapshots)) {
        if (this.controller.state.chats[snapshot.chatId]?.archived) continue;
        yield* this.reconcile(snapshot).pipe(Effect.catch((error) => this.report(snapshot.chatId, error)));
      }
      for (const snapshot of orderedWaits(Object.values(this.snapshots)))
        yield* this.recover(snapshot).pipe(Effect.catch((error) => this.report(snapshot.chatId, error)));
    }.bind(this),
  );
  private reconcile = Effect.fn("Scheduler.reconcile")(
    function* (this: Scheduler, original: InputSnapshot) {
      let snapshot = original,
        state = policyOf(snapshot);
      if (this.manualStops.has(snapshot.chatId) && !state.pause)
        snapshot = yield* this.write(snapshot, {
          ...state,
          pause: "Stopped by you",
          wait: state.wait ? { ...state.wait, status: "paused" } : undefined,
        });
      state = policyOf(snapshot);
      const previousWait = state.wait;
      const next = settlePolicy(snapshot, state);
      if (JSON.stringify(next) !== JSON.stringify(state)) {
        snapshot = yield* this.write(snapshot, next);
        state = next;
        if (
          previousWait?.status === "attempting" &&
          next.wait?.runId !== previousWait.runId &&
          snapshot.lastRun?.usageLimit
        )
          yield* this.pausePeers(snapshot, previousWait.resetAt);
      }
      if (!this.preferences.generateContinue) {
        for (const input of pending(snapshot).filter(
          (input) => input.origin === "generated_continue" && editable(input),
        )) {
          snapshot = yield* this.controller.host.call("host/inputs.update", {
            chatId: snapshot.chatId,
            revision: snapshot.revision,
            inputId: input.id,
            remove: true,
          });
          this.receive(snapshot);
        }
      } else if (
        state.wait &&
        !state.pause &&
        !state.wait.generated &&
        !pending(snapshot).length &&
        !snapshot.running &&
        !snapshot.settling
      ) {
        // Record generation before enqueue; its deterministic id makes a crash
        // after enqueue safe. User enqueue wins atomically in the host.
        state = { ...state, wait: { ...state.wait, generated: true } };
        snapshot = yield* this.write(snapshot, state);
        snapshot = yield* this.controller.host.call("host/inputs.enqueue", {
          chatId: snapshot.chatId,
          id: `continue:${snapshot.generation}:${state.wait?.runId}`,
          blocks: [{ type: "text", text: "Continue" }],
          intent: "queue",
          origin: "generated_continue",
        });
        this.receive(snapshot);
      }
      const input = fifoCandidate(snapshot, state);
      if (!input) return;
      snapshot = yield* this.write(snapshot, { ...state, releasedRun: snapshot.lastRun?.id });
      const delivery = yield* this.dispatch(snapshot, input.id).pipe(Effect.result);
      if (delivery._tag === "Failure") {
        const latest = yield* this.read(snapshot.chatId);
        const current = latest.inputs.find((entry) => entry.id === input.id),
          currentPolicy = policyOf(latest);
        // A stale claim must not consume the successful run's FIFO boundary.
        // Only an unchanged host-held attempt proves no admission occurred;
        // rejected/unknown/admitted states keep their original fencing intact.
        if (
          latest.generation === snapshot.generation &&
          latest.lastRun?.id === snapshot.lastRun?.id &&
          current?.state === "held" &&
          current.attemptId === input.attemptId &&
          currentPolicy.releasedRun === snapshot.lastRun?.id
        )
          yield* this.write(latest, { ...currentPolicy, releasedRun: state.releasedRun });
        return yield* Effect.fail(delivery.failure);
      }
    }.bind(this),
  );
  private pausePeers = Effect.fn("Scheduler.pausePeers")(
    function* (this: Scheduler, snapshot: InputSnapshot, resetAt?: string) {
      const currentWait = policyOf(snapshot).wait;
      if (!currentWait) return;
      const scope = scopeOf(snapshot, currentWait);
      if (!scope) return;
      for (const peer of Object.values(this.snapshots)) {
        const state = policyOf(peer),
          wait = state.wait;
        if (peer.chatId === snapshot.chatId || !wait || scopeOf(peer, wait) !== scope || wait.resetAt !== resetAt)
          continue;
        yield* this.write(peer, {
          ...state,
          wait: { ...wait, status: "paused", reason: "Shared quota was exhausted again; retry manually" },
        });
      }
    }.bind(this),
  );
  private recover = Effect.fn("Scheduler.recover")(
    function* (this: Scheduler, original: InputSnapshot) {
      const chatId = original.chatId;
      if (this.controller.state.chats[chatId]?.archived) return;
      let snapshot = original,
        state = policyOf(snapshot),
        wait = state.wait;
      if (
        !wait ||
        state.pause ||
        this.manualStops.has(chatId) ||
        snapshot.running ||
        snapshot.settling ||
        uncertain(snapshot) ||
        !autoSend(this.preferences, snapshot.agentId)
      )
        return;
      const input = firstQueued(snapshot);
      if (!input) return;
      const now = Date.now(),
        key = waitKey(wait);
      const signal = this.controller.state.agents[snapshot.agentId]?.usageLimits?.recovery;
      const explicitlyAllowed =
        signal?.availability === "allowed" && matching(wait, signal) && fresh(signal, wait, now);
      if (
        (wait.checkedKey === key && !explicitlyAllowed) ||
        (wait.resetAt && Date.parse(wait.resetAt) > now && !explicitlyAllowed)
      )
        return;
      const scope = scopeOf(snapshot, wait);
      if (
        scope &&
        Object.values(this.snapshots).some((peer) => {
          const peerWait = policyOf(peer).wait;
          return (
            peerWait?.status === "attempting" &&
            scopeOf(peer, peerWait) === scope &&
            (peer.running || peer.settling || uncertain(peer))
          );
        })
      )
        return;
      this.checking.add(chatId);
      this.controller.refresh(chatId);
      const result = yield* this.controller.host.call("host/agents.refresh_usage", { agentId: snapshot.agentId }).pipe(
        Effect.catch((error) => (unknownMethod(error) ? Effect.succeed({ limits: null }) : Effect.fail(error))),
        Effect.ensuring(
          Effect.sync(() => {
            this.checking.delete(chatId);
            this.controller.refresh(chatId);
          }),
        ),
        Effect.result,
      );
      // Re-read after the external check: user input, native retries, Stop,
      // attachment changes and another policy can all have won while it waited.
      snapshot = yield* this.read(chatId);
      state = policyOf(snapshot);
      this.preferences = yield* parse(
        "Chat automation preferences",
        preferencesSchema,
        this.controller.api.settings.get(),
      );
      if (
        snapshot.generation !== original.generation ||
        !state.wait ||
        waitKey(state.wait) !== key ||
        snapshot.running ||
        snapshot.settling ||
        uncertain(snapshot) ||
        state.pause ||
        this.manualStops.has(chatId) ||
        state.wait.status !== "armed" ||
        !autoSend(this.preferences, snapshot.agentId)
      )
        return;
      wait = state.wait;
      if (result._tag === "Failure") {
        yield* this.write(snapshot, {
          ...state,
          wait: { ...wait, status: "failed", reason: "Availability check failed; retry manually", checkedKey: key },
        });
        return yield* this.report(chatId, result.failure);
      }
      const checkedAt = Date.now();
      let evidence = result.success.limits?.recovery;
      const latest = this.controller.state.agents[snapshot.agentId]?.usageLimits?.recovery;
      // A refresh may be unsupported/null or race a native push. Preserve
      // affirmative native evidence and newer account changes; blocked wins
      // over sparse/equally old facts, not a newer matching affirmative allowed.
      if (latest && fresh(latest, wait, checkedAt)) {
        const latestAt = Date.parse(latest.observedAt),
          refreshedAt = evidence ? Date.parse(evidence.observedAt) : 0;
        if (
          !evidence?.identity ||
          (matching(wait, evidence) &&
            (!fresh(evidence, wait, checkedAt) ||
              (latest.identity && !matching(wait, latest) && latestAt >= refreshedAt) ||
              (matching(wait, latest) &&
                ((latest.availability === "blocked" && latestAt >= refreshedAt) ||
                  (latest.availability === "allowed" &&
                    (evidence.availability !== "blocked" || latestAt > refreshedAt))))))
        )
          evidence = latest;
      }
      const decision = recoveryDecision(wait, evidence, checkedAt, this.restored.has(chatId));
      if (evidence && matching(wait, evidence) && fresh(evidence, wait, checkedAt)) this.restored.delete(chatId);
      if (decision.kind !== "send") {
        const resetAt = decision.kind === "hold" ? decision.resetAt : undefined;
        const future = resetAt && Date.parse(resetAt) > checkedAt;
        yield* this.write(snapshot, {
          ...state,
          wait: {
            ...wait,
            status: decision.kind === "invalidate" || !wait.identity || this.restored.has(chatId) ? "paused" : "armed",
            reason: decision.reason,
            resetAt,
            checkedKey: future ? undefined : key,
          },
        });
        return;
      }
      this.restored.delete(chatId);
      const selected = firstQueued(snapshot);
      if (!selected || selected.id !== input.id) return;
      wait = {
        ...wait,
        status: "attempting",
        attemptKey: key,
        inputId: selected.id,
        reason: decision.speculative
          ? "Trying once after the reported reset"
          : "Availability confirmed; sending follow-up",
      };
      snapshot = yield* this.write(snapshot, { ...state, wait });
      yield* this.dispatch(snapshot, selected.id, scope).pipe(Effect.catch((error) => this.report(chatId, error)));
    }.bind(this),
  );
}

Versions

VersionPublishedPlugin APISizePermissionsStatus
0.2.0latestOct 5, 2026>=2 <3105.1 KB5 permissionsListed

Reviews and comments

0 threads · 0 reviews

No comments yet.