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.test.ts30.2 KB
// Competing lifecycle observations exercise a stateful CAS input ledger, not
// echo mocks. Finishing a transcript and settling a host run are distinct.
import assert from "node:assert/strict";
import test from "node:test";
import * as Effect from "effect/Effect";
import * as Fiber from "effect/Fiber";
import { promise } from "convergence/effect";
import { InputLedger, held, settledRun, withController } from "./testing.ts";
import { onAgent } from "./events.ts";
import { policyOf, recoveryDecision, defaults, autoSend, scopeOf, waitKey } from "./scheduler-policy.ts";
import type { RecoveryWait } from "./scheduler-policy.ts";
import type { UsageRecovery, UsageLimits } from "convergence";
import { restoreQueues } from "./data.ts";
const stopped = "2026-10-04T12:00:00Z",
reset = "2026-10-04T12:05:00Z";
const evidence = (patch: Partial<UsageRecovery> = {}): UsageRecovery => ({
availability: "blocked",
source: "native",
observedAt: stopped,
identity: "account:A",
scope: "quota",
resetsAt: reset,
reason: "Subscription exhausted",
...patch,
});
const wait = (patch: Partial<RecoveryWait> = {}): RecoveryWait => ({
runId: "limit",
stoppedAt: stopped,
identity: "account:A",
scope: "quota",
resetAt: reset,
reason: "Subscription exhausted",
status: "armed",
...patch,
});
function arm(ledger: InputLedger, id = "c1", patch: Partial<RecoveryWait> = {}) {
ledger.mutate(id, (snapshot) => {
snapshot.inputs = [held(`next:${id}`, { chatId: id })];
snapshot.lastRun = settledRun("limit", { outcome: { status: "failed", message: "limit" }, usageLimit: evidence() });
snapshot.policyState = { version: 1, generation: 1, handledRun: "limit", wait: wait(patch) };
});
}
test("FIFO dispatches one after settlement, never run_finished, and uses the live policy token", async () => {
const ledger = new InputLedger();
await withController(ledger, (c) =>
Effect.gen(function* () {
const scheduler = c.scheduler;
yield* scheduler.enqueue("c1", "one", [{ type: "text", text: "One" }], "queue");
yield* scheduler.enqueue("c1", "two", [{ type: "text", text: "Two" }], "queue");
yield* scheduler.evaluate();
assert.equal(ledger.deliveries.length, 0, "idle Queue starts nothing");
ledger.mutate("c1", (snapshot) => {
snapshot.running = true;
snapshot.runId = "initial";
});
scheduler.receive(ledger.finalizing("c1"));
yield* onAgent(c, {
event: "agent",
chatId: "c1",
kind: { event: "run_finished", outcome: { status: "completed" } },
});
yield* scheduler.evaluate();
assert.equal(ledger.deliveries.length, 0, "checkpoint finalization has not settled");
scheduler.receive(ledger.settled("c1", settledRun("initial")));
yield* Effect.all([scheduler.evaluate(), scheduler.evaluate()], { concurrency: "unbounded" });
assert.deepEqual(
ledger.deliveries.map((delivery) => delivery.inputId),
["one"],
);
assert.equal(ledger.deliveries[0].policy, ledger.token);
scheduler.receive(ledger.consumed("c1", "one"));
scheduler.receive(ledger.settled("c1", settledRun("run:one")));
yield* scheduler.evaluate();
assert.deepEqual(
ledger.deliveries.map((delivery) => delivery.inputId),
["one", "two"],
);
}),
);
});
test("Stop, ordinary failure and simulated intermediate cancellation do not drain FIFO", async () => {
for (const run of [
settledRun("failed", { outcome: { status: "failed", message: "ordinary failure" } }),
settledRun("stop", { outcome: { status: "cancelled" } }),
settledRun("simulation", { outcome: { status: "cancelled" }, steering: true }),
]) {
const ledger = new InputLedger();
await withController(ledger, (c) =>
Effect.gen(function* () {
yield* c.scheduler.enqueue("c1", "q", [{ type: "text", text: "Follow-up" }], "queue");
c.scheduler.receive(ledger.settled("c1", run));
yield* c.scheduler.evaluate();
assert.equal(ledger.deliveries.length, 0);
assert.equal(ledger.snapshots.c1.inputs[0].state, "held");
assert.equal(
policyOf(ledger.snapshots.c1).pause,
run.steering ? undefined : run.outcome.status === "cancelled" ? "Stopped by you" : "Run failed",
);
}),
);
}
});
test("generation changes, foreign policy state and unknown delivery cannot automatically replay", async () => {
const ledger = new InputLedger();
arm(ledger);
ledger.mutate("c1", (snapshot) => {
snapshot.generation = 2;
});
await withController(ledger, (c) =>
Effect.gen(function* () {
yield* c.scheduler.evaluate();
assert.match(policyOf(ledger.snapshots.c1).pause ?? "", /Attachment changed/);
assert.equal(ledger.deliveries.length, 0);
}),
);
ledger.mutate("c1", (snapshot) => {
snapshot.generation = 1;
snapshot.inputs[0].state = "delivery_unknown";
snapshot.policyState = { version: 1, generation: 1, wait: wait(), handledRun: "limit" };
});
await withController(ledger, (c) =>
Effect.gen(function* () {
yield* c.scheduler.restore();
yield* c.scheduler.evaluate();
assert.equal(ledger.deliveries.length, 0);
}),
);
});
test("promote preserves identity, attachments and an unrelated draft; stale CAS never edits a newer entry", async () => {
const ledger = new InputLedger();
await withController(ledger, (c) =>
Effect.gen(function* () {
c.composer("c1").text = "Unrelated draft";
const snapshot = yield* c.scheduler.enqueue(
"c1",
"q",
[
{ type: "text", text: "Inspect" },
{ type: "image", mimeType: "image/png", data: "bytes" },
],
"queue",
);
yield* c.scheduler.change({
chatId: "c1",
inputId: "q",
revision: snapshot.revision,
blocks: [{ type: "text", text: "Newer" }, ...snapshot.inputs[0].blocks.slice(1)],
});
const stale = yield* c.scheduler
.change({ chatId: "c1", inputId: "q", revision: snapshot.revision, remove: true })
.pipe(Effect.result);
assert.equal(stale._tag, "Failure");
yield* c.scheduler.promote("c1", "q", ledger.snapshots.c1.revision);
assert.equal(c.composer("c1").text, "Unrelated draft");
assert.deepEqual(ledger.deliveries[0].blocks, [
{ type: "text", text: "Newer" },
{ type: "image", mimeType: "image/png", data: "bytes" },
]);
assert.equal(ledger.deliveries[0].inputId, "q");
}),
);
});
test("append uses host enqueue and keeps the complete last held steer in contribution order", async () => {
const ledger = new InputLedger();
await withController(ledger, (c) =>
Effect.gen(function* () {
yield* c.scheduler.enqueue("c1", "a", [{ type: "text", text: "Inspect alpha" }], "steer");
const snapshot = yield* c.scheduler.enqueue(
"c1",
"b",
[
{ type: "text", text: "Do not change beta" },
{ type: "image", mimeType: "image/png", data: "bytes" },
],
"steer",
);
assert.equal(snapshot.inputs.length, 1);
assert.deepEqual(snapshot.inputs[0].blocks, [
{ type: "text", text: "Inspect alpha\nDo not change beta" },
{ type: "image", mimeType: "image/png", data: "bytes" },
]);
assert.equal(snapshot.inputs[0].id, "a");
}),
);
});
test("legacy migration is idempotent across a crash and retains whole-chat queues", async () => {
const ledger = new InputLedger();
ledger.storage.queues = {
prompts: {
c1: [
"First",
{
text: "With image",
attachments: [],
blocks: [
{ type: "text", text: "With image" },
{ type: "image", mimeType: "image/png", data: "durable" },
],
},
],
},
chats: [{ chatId: "other", text: "Whole chat" }],
};
const original = structuredClone(ledger.storage.queues);
await withController(ledger, (c) => restoreQueues(c));
assert.equal(ledger.snapshots.c1.inputs.length, 2);
ledger.storage.queues = original; // source deletion was lost in the crash
await withController(ledger, (c) => restoreQueues(c));
assert.equal(ledger.snapshots.c1.inputs.length, 2);
assert.deepEqual(ledger.snapshots.c1.inputs[1].blocks[1], { type: "image", mimeType: "image/png", data: "durable" });
assert.deepEqual(ledger.storage.queues, { prompts: {}, chats: [{ chatId: "other" }] });
assert.equal(ledger.deliveries.length, 0);
});
test("recovery defaults and permissions never infer allowed from null, percentages or a stale reset", () => {
assert.equal(defaults.generateContinue, false);
assert.equal(autoSend(defaults, "agent"), true);
const now = Date.parse(reset);
assert.equal(
recoveryDecision(wait(), undefined, now, false).kind,
"send",
"one visibly speculative future-reset attempt",
);
assert.equal(recoveryDecision(wait(), undefined, now, true).kind, "hold", "restart must establish identity");
assert.equal(recoveryDecision(wait({ resetAt: stopped }), undefined, now, false).kind, "hold");
assert.equal(
recoveryDecision(wait({ resetAt: undefined }), evidence({ availability: "unknown", observedAt: reset }), now, false)
.kind,
"hold",
);
assert.equal(
recoveryDecision(wait(), evidence({ availability: "allowed", identity: "other", observedAt: reset }), now, false)
.kind,
"invalidate",
);
assert.equal(
recoveryDecision(wait(), evidence({ observedAt: reset }), now, false).kind,
"hold",
"affirmative rejection wins",
);
assert.equal(
recoveryDecision(wait({ resetAt: undefined }), evidence({ availability: "allowed", observedAt: reset }), now, true)
.kind,
"send",
);
});
test("speculative retry is bounded before, at and after reset for unknown/unsupported evidence", () => {
for (const availability of ["unknown", "unsupported"] as const) {
for (const offset of [-1, 0, 1]) {
const now = Date.parse(reset) + offset;
const signal = evidence({ availability, observedAt: new Date(now).toISOString() });
const decision = recoveryDecision(wait(), signal, now, false);
assert.equal(decision.kind, offset < 0 ? "hold" : "send");
if (decision.kind === "send") assert.equal(decision.speculative, true);
assert.equal(recoveryDecision(wait({ attemptKey: waitKey(wait()) }), signal, now, false).kind, "hold");
}
}
});
test("host settlement alone does not invent pickup or replay admitted native input", async () => {
const ledger = new InputLedger();
ledger.mutate("c1", (snapshot) => {
snapshot.inputs = [
held("admitted", { intent: "steer", state: "native_admitted", runId: "original" }),
held("queued"),
];
snapshot.policyState = { version: 1, generation: 1 };
});
await withController(ledger, (c) =>
Effect.gen(function* () {
c.scheduler.receive(ledger.settled("c1", settledRun("original")));
yield* c.scheduler.evaluate();
assert.equal(ledger.snapshots.c1.inputs[0].state, "native_admitted");
assert.equal(ledger.deliveries.length, 0);
c.scheduler.receive(ledger.consumed("c1", "admitted"));
yield* c.scheduler.evaluate();
assert.deepEqual(
ledger.deliveries.map((delivery) => delivery.inputId),
["queued"],
);
}),
);
});
test("finite retry at the deadline; renewed exhaustion pauses peers and old reset never loops", async (ctx) => {
let now = Date.parse(reset) - 1;
ctx.mock.method(Date, "now", () => now);
const ledger = new InputLedger(["newer", "older", "distinct"]);
arm(ledger, "newer", { stoppedAt: "2026-10-04T12:01:00Z" });
arm(ledger, "older");
arm(ledger, "distinct", { identity: "account:B", stoppedAt: "2026-10-04T12:02:00Z" });
await withController(ledger, (c) =>
Effect.gen(function* () {
yield* c.scheduler.evaluate();
assert.equal(ledger.deliveries.length, 0);
now++;
yield* c.scheduler.evaluate();
assert.deepEqual(
ledger.deliveries.map((delivery) => delivery.chatId),
["older", "distinct"],
);
assert.equal(ledger.deliveries[0].scope, scopeOf(ledger.snapshots.older, wait()));
c.scheduler.receive(
ledger.settled(
"older",
settledRun("run:next:older", {
stoppedAt: "2026-10-04T12:05:01Z",
outcome: { status: "failed", message: "quota again" },
usageLimit: evidence(),
}),
),
);
yield* c.scheduler.evaluate();
assert.equal(policyOf(ledger.snapshots.newer).wait?.status, "paused");
for (let i = 0; i < 20; i++) {
now += 1000;
yield* c.scheduler.evaluate();
}
assert.equal(ledger.deliveries.length, 2, "no speculative prompt polling or peer stampede");
}),
);
});
test("fresh blocked refresh postpones a wait; identity mismatch invalidates it", async (ctx) => {
let now = Date.parse(reset);
ctx.mock.method(Date, "now", () => now);
const ledger = new InputLedger();
arm(ledger);
let checks = 0;
ledger.availability = async () => {
checks++;
return {
limits: { recovery: evidence({ observedAt: new Date(now).toISOString(), resetsAt: "2026-10-04T12:10:00Z" }) },
};
};
await withController(ledger, (c) =>
Effect.gen(function* () {
yield* c.scheduler.evaluate();
assert.equal(policyOf(ledger.snapshots.c1).wait?.resetAt, "2026-10-04T12:10:00Z");
for (let i = 0; i < 10; i++) yield* c.scheduler.evaluate();
assert.equal(checks, 1);
now = Date.parse("2026-10-04T12:10:00Z");
ledger.availability = async () => ({
limits: {
recovery: evidence({
availability: "allowed",
identity: "account:other",
observedAt: new Date(now).toISOString(),
}),
},
});
yield* c.scheduler.evaluate();
assert.equal(policyOf(ledger.snapshots.c1).wait?.status, "paused");
assert.equal(ledger.deliveries.length, 0);
}),
);
});
test("null or stale refresh cannot override a fresh matching native rejection", async (ctx) => {
ctx.mock.method(Date, "now", () => Date.parse(reset));
for (const limits of [
null,
{ recovery: evidence({ availability: "allowed", observedAt: stopped }) },
{ recovery: evidence({ availability: "unknown", observedAt: reset, identity: undefined }) },
{ recovery: evidence({ availability: "unknown", observedAt: reset }) },
]) {
const ledger = new InputLedger();
arm(ledger);
ledger.availability = async () => ({ limits });
await withController(ledger, (c) =>
Effect.gen(function* () {
c.state.agents.agent.usageLimits = {
recovery: evidence({ observedAt: reset, resetsAt: "2026-10-04T12:10:00Z" }),
};
yield* c.scheduler.evaluate();
assert.equal(ledger.deliveries.length, 0, "elapsed reset is not permission over fresh blocked evidence");
assert.equal(policyOf(ledger.snapshots.c1).wait?.resetAt, "2026-10-04T12:10:00Z");
}),
);
}
});
test("fresh matching allowed supersedes an older rejection, but a changed identity does not", async (ctx) => {
ctx.mock.method(Date, "now", () => Date.parse(reset));
for (const identity of ["account:A", "account:B"]) {
const ledger = new InputLedger();
arm(ledger);
ledger.availability = async () => ({
limits: { recovery: evidence({ availability: "allowed", observedAt: reset, identity }) },
});
await withController(ledger, (c) =>
Effect.gen(function* () {
c.state.agents.agent.usageLimits = { recovery: evidence({ observedAt: "2026-10-04T12:04:59Z" }) };
yield* c.scheduler.evaluate();
assert.equal(ledger.deliveries.length, identity === "account:A" ? 1 : 0);
if (identity !== "account:A") assert.equal(policyOf(ledger.snapshots.c1).wait?.status, "paused");
}),
);
}
});
test("fresh native availability remains authoritative when refresh is unsupported or races identity change", async (ctx) => {
ctx.mock.method(Date, "now", () => Date.parse(reset));
for (const kind of ["allowed", "identity-change"] as const) {
const ledger = new InputLedger();
arm(ledger);
if (kind === "identity-change")
ledger.availability = async () => ({
limits: {
recovery: evidence({ availability: "allowed", observedAt: "2026-10-04T12:04:59Z" }),
},
});
await withController(ledger, (c) =>
Effect.gen(function* () {
c.state.agents.agent.usageLimits = {
recovery: evidence({
observedAt: reset,
availability: kind === "allowed" ? "allowed" : "blocked",
identity: kind === "allowed" ? "account:A" : "account:B",
}),
};
// Allowed native evidence needs no elapsed reset; a newer account change
// must invalidate even an earlier affirmative refresh of the old account.
if (kind === "allowed")
c.scheduler.receive(
ledger.mutate("c1", (snapshot) => {
snapshot.policyState = { ...policyOf(snapshot), wait: wait({ resetAt: undefined }) };
}),
);
yield* c.scheduler.evaluate();
assert.equal(ledger.deliveries.length, kind === "allowed" ? 1 : 0);
if (kind === "identity-change") assert.equal(policyOf(ledger.snapshots.c1).wait?.status, "paused");
}),
);
}
});
test("restart restores all chats with due identity-validated waits, but not manual pauses or unknown delivery", async (ctx) => {
const now = Date.parse(reset) + 60_000;
ctx.mock.method(Date, "now", () => now);
const ledger = new InputLedger(["hidden", "manual", "unknown"]);
arm(ledger, "hidden");
arm(ledger, "manual");
arm(ledger, "unknown");
ledger.mutate("manual", (snapshot) => {
snapshot.policyState = { ...policyOf(snapshot), pause: "Stopped by you" };
});
ledger.mutate("unknown", (snapshot) => {
snapshot.inputs[0].state = "delivery_unknown";
});
ledger.availability = async () => ({
limits: { recovery: evidence({ availability: "unknown", observedAt: new Date(now).toISOString() }) },
});
await withController(ledger, (c) =>
Effect.gen(function* () {
c.scheduler.ready = false;
yield* c.scheduler.restore();
yield* c.scheduler.evaluate();
assert.deepEqual(
ledger.deliveries.map((delivery) => delivery.chatId),
["hidden"],
);
assert.equal(policyOf(ledger.snapshots.hidden).wait?.attemptKey, waitKey(wait()));
}),
);
});
test("usage grace leaves work running and cannot generate Continue; generation is exact, once, and user input supersedes", async (ctx) => {
ctx.mock.method(Date, "now", () => Date.parse(stopped));
const ledger = new InputLedger();
ledger.settings.generateContinue = true;
await withController(ledger, (c) =>
Effect.gen(function* () {
yield* c.scheduler.enqueue("c1", "initial", [{ type: "text", text: "Work" }], "send");
yield* c.scheduler.manualDispatch("c1", "initial", ledger.snapshots.c1.revision);
yield* c.scheduler.updatePreferences();
yield* onAgent(c, { event: "agent", chatId: "c1", kind: { event: "usage_blocked", recovery: evidence() } });
yield* c.scheduler.evaluate();
assert.equal(ledger.snapshots.c1.running, true);
assert.equal(ledger.snapshots.c1.inputs.length, 1);
assert.ok(!ledger.calls.some((call) => call.method === "host/chats.cancel"));
c.scheduler.receive(ledger.consumed("c1", "initial"));
c.scheduler.receive(
ledger.settled(
"c1",
settledRun("run:initial", { outcome: { status: "failed", message: "actual limit" }, usageLimit: evidence() }),
),
);
yield* c.scheduler.evaluate();
yield* c.scheduler.evaluate();
assert.deepEqual(
ledger.snapshots.c1.inputs
.filter((input) => input.origin === "generated_continue")
.map((input) => input.blocks),
[[{ type: "text", text: "Continue" }]],
);
yield* c.scheduler.enqueue("c1", "user", [{ type: "text", text: "My instruction" }], "queue");
assert.equal(ledger.snapshots.c1.inputs.filter((input) => input.origin === "generated_continue").length, 0);
}),
);
});
test("per-agent preference changes affect every chat immediately, generation-off removes only held generated input", async () => {
const ledger = new InputLedger(["c1", "c2"]);
arm(ledger, "c1");
arm(ledger, "c2");
ledger.mutate("c1", (snapshot) => {
snapshot.inputs.push(
held("generated", { origin: "generated_continue" }),
held("admitted", { origin: "generated_continue", state: "native_admitted" }),
);
});
await withController(ledger, (c) =>
Effect.gen(function* () {
yield* c.scheduler.setAutoSend("agent", false);
assert.equal(autoSend(c.scheduler.preferences, "agent"), false);
assert.equal(autoSend(c.scheduler.preferences, "other"), true);
assert.equal(
ledger.snapshots.c1.inputs.some((input) => input.id === "generated"),
false,
);
assert.equal(
ledger.snapshots.c1.inputs.some((input) => input.id === "admitted"),
true,
);
assert.ok(ledger.snapshots.c1.inputs.some((input) => input.id === "next:c1"));
assert.ok(ledger.snapshots.c2.inputs.some((input) => input.id === "next:c2"));
assert.equal(ledger.deliveries.length, 0);
}),
);
});
test("availability failure is finite and manual Stop during an in-flight check wins", async (ctx) => {
ctx.mock.method(Date, "now", () => Date.parse(reset));
const failed = new InputLedger();
arm(failed);
failed.availability = async () => {
throw new Error("usage API unavailable");
};
await withController(failed, (c) =>
Effect.gen(function* () {
for (let i = 0; i < 10; i++) yield* c.scheduler.evaluate();
assert.equal(failed.calls.filter((call) => call.method === "host/agents.refresh_usage").length, 1);
assert.equal(policyOf(failed.snapshots.c1).wait?.status, "failed");
assert.equal(failed.deliveries.length, 0);
}),
);
const ledger = new InputLedger();
arm(ledger);
const pending = Promise.withResolvers<{ limits: UsageLimits | null }>();
ledger.availability = () => pending.promise;
await withController(ledger, (c) =>
Effect.gen(function* () {
yield* Effect.forkIn(c.scheduler.evaluate(), c.scope, { startImmediately: true });
yield* Effect.yieldNow;
yield* c.scheduler.stop("c1");
pending.resolve({ limits: { recovery: evidence({ availability: "allowed", observedAt: reset }) } });
yield* Effect.yieldNow;
assert.equal(ledger.deliveries.length, 0);
assert.equal(policyOf(ledger.snapshots.c1).pause, "Stopped by you");
}),
);
});
test("slow automatic admission and usage refresh cannot block durable enqueue or Stop", async (ctx) => {
ctx.mock.method(Date, "now", () => Date.parse(reset));
for (const operation of ["admission", "usage"] as const) {
const ledger = new InputLedger();
const started = Promise.withResolvers<void>(),
release = Promise.withResolvers<void>();
if (operation === "usage") {
arm(ledger);
ledger.availability = async () => {
started.resolve();
await release.promise;
return { limits: { recovery: evidence({ availability: "allowed", observedAt: reset }) } };
};
} else {
ledger.mutate("c1", (snapshot) => {
snapshot.inputs = [held("next:c1")];
snapshot.lastRun = settledRun("completed");
snapshot.policyState = { version: 1, generation: 1 };
});
ledger.afterClaim = async () => {
started.resolve();
await release.promise;
};
}
await withController(ledger, (c) =>
Effect.gen(function* () {
const evaluation = yield* Effect.forkIn(c.scheduler.evaluate(), c.scope, { startImmediately: true });
yield* promise("test.started", () => started.promise);
yield* Effect.gen(function* () {
yield* c.scheduler
.enqueue("c1", "new-user", [{ type: "text", text: "Retain this" }], "queue")
.pipe(Effect.timeout("1 second"));
assert.ok(
ledger.snapshots.c1.inputs.some((input) => input.id === "new-user"),
`${operation} must not hold the mutation semaphore`,
);
yield* c.scheduler.stop("c1").pipe(Effect.timeout("1 second"));
assert.equal(policyOf(ledger.snapshots.c1).pause, "Stopped by you");
}).pipe(Effect.ensuring(Effect.sync(() => release.resolve())));
yield* Fiber.join(evaluation);
assert.equal(ledger.deliveries.length, operation === "admission" ? 1 : 0);
c.scheduler.receive(ledger.settled("c1", settledRun("cancelled", { outcome: { status: "cancelled" } })));
yield* c.scheduler.evaluate();
assert.equal(ledger.snapshots.c1.inputs.find((input) => input.id === "new-user")?.state, "held");
}),
);
}
});
test("a positively rejected native steer stays recoverable until settlement and explicit idle dispatch reuses its identity", async () => {
const ledger = new InputLedger();
ledger.mutate("c1", (snapshot) => {
snapshot.running = true;
snapshot.runId = "original";
});
ledger.afterClaim = async (claimed) => {
ledger.finalizing("c1");
ledger.mutate("c1", (snapshot) => {
const input = snapshot.inputs.find((input) => input.id === "steer");
if (input) {
input.state = "rejected";
input.error = "Fixture positive non-admission: native turn already ended";
}
});
assert.equal(claimed.inputs[0].runId, "original");
throw new Error("native turn already ended");
};
await withController(ledger, (c) =>
Effect.gen(function* () {
const queued = yield* c.scheduler.enqueue(
"c1",
"steer",
[
{ type: "text", text: "New direction" },
{ type: "image", mimeType: "image/png", data: "bytes" },
],
"steer",
);
const first = yield* c.scheduler.manualDispatch("c1", "steer", queued.revision).pipe(Effect.result);
assert.equal(first._tag, "Failure");
yield* c.scheduler.load("c1");
const rejected = ledger.snapshots.c1.inputs[0],
attempt = rejected.attemptId;
assert.equal(rejected.state, "rejected");
yield* c.scheduler.evaluate();
const premature = yield* c.scheduler
.manualDispatch("c1", "steer", ledger.snapshots.c1.revision)
.pipe(Effect.result);
assert.equal(premature._tag, "Failure");
assert.equal(ledger.deliveries.length, 1, "native completion is not host settlement");
c.scheduler.receive(ledger.settled("c1", settledRun("original")));
ledger.afterClaim = undefined;
yield* c.scheduler.manualDispatch("c1", "steer", ledger.snapshots.c1.revision);
assert.deepEqual(
ledger.deliveries.map((delivery) => delivery.inputId),
["steer", "steer"],
);
assert.deepEqual(ledger.deliveries[1].blocks, rejected.blocks);
assert.notEqual(ledger.snapshots.c1.inputs[0].attemptId, attempt);
assert.equal(ledger.snapshots.c1.inputs[0].runId, "run:steer");
}),
);
});
test("Stop and service replacement while automatic authority is loading cannot send a stale claim", async () => {
for (const invalidate of ["stop", "replace"] as const) {
const ledger = new InputLedger();
ledger.mutate("c1", (snapshot) => {
snapshot.inputs = [held("next")];
snapshot.lastRun = settledRun("completed");
snapshot.policyState = { version: 1, generation: 1 };
});
const started = Promise.withResolvers<void>(),
release = Promise.withResolvers<void>();
ledger.beforePolicy = async () => {
started.resolve();
await release.promise;
};
await withController(ledger, (c) =>
Effect.gen(function* () {
const evaluation = yield* Effect.forkIn(c.scheduler.evaluate(), c.scope, { startImmediately: true });
yield* promise("test.started", () => started.promise);
if (invalidate === "stop") yield* c.scheduler.stop("c1");
else ledger.token = null;
release.resolve();
yield* Fiber.join(evaluation);
assert.equal(ledger.deliveries.length, 0);
assert.equal(ledger.snapshots.c1.inputs[0].state, "held");
}),
);
}
});
test("editing held FIFO input during authority lookup rejects stale CAS without losing the settlement release", async () => {
const ledger = new InputLedger();
ledger.mutate("c1", (snapshot) => {
snapshot.inputs = [held("next")];
snapshot.lastRun = settledRun("completed");
snapshot.policyState = { version: 1, generation: 1 };
});
const started = Promise.withResolvers<void>(),
release = Promise.withResolvers<void>();
ledger.beforePolicy = async () => {
started.resolve();
await release.promise;
};
await withController(ledger, (c) =>
Effect.gen(function* () {
const evaluation = yield* Effect.forkIn(c.scheduler.evaluate(), c.scope, { startImmediately: true });
yield* promise("test.started", () => started.promise);
yield* c.scheduler.change({
chatId: "c1",
inputId: "next",
revision: ledger.snapshots.c1.revision,
blocks: [{ type: "text", text: "Edited before claim" }],
});
release.resolve();
yield* Fiber.join(evaluation);
assert.equal(ledger.deliveries.length, 0, "stale CAS never reaches native delivery");
assert.equal(ledger.snapshots.c1.inputs[0].state, "held");
ledger.beforePolicy = undefined;
yield* c.scheduler.evaluate();
assert.deepEqual(
ledger.deliveries.map((delivery) => delivery.blocks),
[[{ type: "text", text: "Edited before claim" }]],
);
}),
);
});
test("ambiguous automatic delivery never rolls back its FIFO release or replays the input", async () => {
const ledger = new InputLedger();
ledger.rejectDispatch = true;
ledger.mutate("c1", (snapshot) => {
snapshot.inputs = [held("uncertain"), held("later")];
snapshot.lastRun = settledRun("completed");
snapshot.policyState = { version: 1, generation: 1 };
});
await withController(ledger, (c) =>
Effect.gen(function* () {
for (let i = 0; i < 5; i++) yield* c.scheduler.evaluate();
assert.equal(ledger.calls.filter((call) => call.method === "host/inputs.dispatch").length, 1);
assert.equal(ledger.snapshots.c1.inputs[0].state, "delivery_unknown");
assert.equal(ledger.snapshots.c1.inputs[1].state, "held");
assert.equal(policyOf(ledger.snapshots.c1).releasedRun, "completed");
assert.ok(c.scheduler.errors.c1);
}),
);
});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.