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
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
| Version | Published | Plugin API | Size | Permissions | Status |
|---|---|---|---|---|---|
| 0.2.0latest | Oct 5, 2026 | >=2 <3 | 105.1 KB | 5 permissions | Listed |
No comments yet.