diff --git a/apps/server/src/coil/autoResume/Reactor.test.ts b/apps/server/src/coil/autoResume/Reactor.test.ts index c4f48943f65f..c72076948c1c 100644 --- a/apps/server/src/coil/autoResume/Reactor.test.ts +++ b/apps/server/src/coil/autoResume/Reactor.test.ts @@ -15,6 +15,7 @@ import * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; import * as FileSystem from "effect/FileSystem"; import * as Layer from "effect/Layer"; +import * as PubSub from "effect/PubSub"; import * as Ref from "effect/Ref"; import * as Stream from "effect/Stream"; import * as TestClock from "effect/testing/TestClock"; @@ -23,14 +24,22 @@ import { OrchestrationEngineService } from "../../orchestration/Services/Orchest import { ProjectionSnapshotQuery } from "../../orchestration/Services/ProjectionSnapshotQuery.ts"; import { ProviderService } from "../../provider/Services/ProviderService.ts"; import { AutoResumeReactorLive } from "./Reactor.ts"; +import { + type AutoResumeReactorReceipt, + AutoResumeReactorReceipts, + AutoResumeReactorReceiptsLive, +} from "./receipts.ts"; import { AutoResumeStore, makeAutoResumeStore } from "./state.ts"; // Defaults: safetyMargin 60s, pollMs 30s. With resetsAt=100s the resume is due at // 100_000 + 60_000 = 160_000ms, so advancing past that (with 30s wake ticks) fires it. -// A one-Claude-thread read model. Cast because building every branded field is noise for -// this test — the reactor only reads the fields set here. +// A Claude-thread read model, one row per id (one row unless a scenario asks for more). +// Cast because building every branded field is noise for this test — the reactor only reads +// the fields set here. const readModel = (o: { + /** Ids to build rows for, in order. Every row is otherwise identical. */ + threadIds?: ReadonlyArray; messages?: Array<{ id: string; role: string }>; status?: string; latestTurnId?: string; @@ -44,33 +53,35 @@ const readModel = (o: { snapshotSequence: 1, updatedAt: "2026-01-01T00:00:00.000Z", projects: [{ id: "project-1", workspaceRoot: "/tmp/coil-nonexistent-workspace" }], - threads: [ - { - id: "thread-1", - projectId: "project-1", - runtimeMode: "full-access", - interactionMode: "default", - worktreePath: null, - deletedAt: null, - archivedAt: null, - settledOverride: null, - messages: o.messages ?? [{ id: "u1", role: "user" }], - activities: [], - latestTurn: - o.latestTurn !== undefined - ? o.latestTurn - : { turnId: o.latestTurnId ?? "turn-1", state: "completed" }, - session: { status: o.status ?? "ready", providerName: "claudeAgent" }, - }, - ], + threads: (o.threadIds ?? ["thread-1"]).map((id) => ({ + id, + projectId: "project-1", + runtimeMode: "full-access", + interactionMode: "default", + worktreePath: null, + deletedAt: null, + archivedAt: null, + settledOverride: null, + messages: o.messages ?? [{ id: "u1", role: "user" }], + activities: [], + latestTurn: + o.latestTurn !== undefined + ? o.latestTurn + : { turnId: o.latestTurnId ?? "turn-1", state: "completed" }, + session: { status: o.status ?? "ready", providerName: "claudeAgent" }, + })), }) as unknown as OrchestrationReadModel; -const rejectedEvent = (resetsAtSeconds: number): ProviderRuntimeEvent => +const rejectedEvent = ( + resetsAtSeconds: number, + threadId = "thread-1", + eventId = "evt-1", +): ProviderRuntimeEvent => ({ type: "runtime.warning", - eventId: "evt-1", + eventId, provider: "claudeAgent", - threadId: "thread-1", + threadId, createdAt: "2026-01-01T00:00:00.000Z", payload: { message: "Claude usage limit reached", @@ -110,9 +121,19 @@ const harness = (initialModel: OrchestrationReadModel, events: ProviderRuntimeEv } as unknown as typeof OrchestrationEngineService.Service); const snapshotCalls = yield* Ref.make(0); + // One-shot: the NEXT snapshot read fails and the flag clears itself. `getSnapshot` has a + // typed `ProjectionRepositoryError` channel in production; the shape of the error does + // not matter here, only that the read fails where the reactor expects it can. + const failNextSnapshot = yield* Ref.make(false); const SnapshotStub = Layer.succeed(ProjectionSnapshotQuery, { getSnapshot: () => - Ref.update(snapshotCalls, (n) => n + 1).pipe(Effect.andThen(Ref.get(modelRef))), + Effect.gen(function* () { + yield* Ref.update(snapshotCalls, (n) => n + 1); + if (yield* Ref.getAndSet(failNextSnapshot, false)) { + return yield* Effect.fail(new Error("simulated snapshot read failure")); + } + return yield* Ref.get(modelRef); + }), } as unknown as typeof ProjectionSnapshotQuery.Service); const ProviderStub = Layer.succeed(ProviderService, { @@ -140,104 +161,134 @@ const harness = (initialModel: OrchestrationReadModel, events: ProviderRuntimeEv }), ); - const deps = Layer.mergeAll(EngineStub, SnapshotStub, ProviderStub, CryptoStub, StoreLive); - return { dispatched, modelRef, deps, store, snapshotCalls, failTurnStart }; + const deps = Layer.mergeAll( + EngineStub, + SnapshotStub, + ProviderStub, + CryptoStub, + StoreLive, + // Test-only, and the whole reason the waits below are exact. Production never provides + // it, so the reactor's emitter resolves to a no-op there. See `receipts.ts`. + AutoResumeReactorReceiptsLive, + ); + return { dispatched, modelRef, deps, store, snapshotCalls, failTurnStart, failNextSnapshot }; }); const types = (commands: ReadonlyArray) => commands.map((c) => c.type); -// Conditions for settleUntil/advanceUntil. Each is re-evaluated on every pump, so it always -// reads fresh state rather than a value captured before the loop. -const scheduledOne = (store: { readonly listPending: Effect.Effect> }) => - store.listPending.pipe(Effect.map((pending) => pending.length === 1)); - -const dispatchedIncludes = ( - dispatched: Ref.Ref, - type: OrchestrationCommand["type"], -) => Ref.get(dispatched).pipe(Effect.map((commands) => types(commands).includes(type))); - -const firedCount = ( - store: { readonly countFiredSince: (threadId: string, since: number) => Effect.Effect }, - expected: number, -) => store.countFiredSince("thread-1", 0).pipe(Effect.map((n) => n === expected)); - -// A real event-loop tick. The store persists via writeFileStringAtomically — real -// filesystem I/O whose completion callback fires on the Node event loop, NOT on -// TestClock — and schedule() gates its in-memory ref update behind that write. So -// yieldNow alone (which only pumps the Effect fiber scheduler) never observes a just- -// scheduled resume; we must also let real I/O drain. +// --- waiting ----------------------------------------------------------------------- +// +// Every wait below is an await on a receipt the reactor published (`receipts.ts`). Nothing +// here counts scheduler turns, and nothing here moves the clock except the polls a scenario +// asks for. +// +// It used to count turns: `advancePastResume` was eight clock steps each chased by a fixed +// ten-pump spin, a budget rather than a signal. The store persists through +// `writeFileStringAtomically` — real filesystem I/O whose completion fires on the Node event +// loop and NOT on TestClock — so on a loaded runner a turn buys less real progress and a +// budget that is generous on a laptop runs out. That is exactly what radroid/t3code#134 was: +// `pending must be cleared: expected 1 to equal +0`, an assertion that looked before the +// cancellation had landed, reported as a defect in the product. + +/** The reactor's poll cadence (`config.pollMs`), and the size of one clock step here. */ +const POLL_MS = 30_000; + +/** A real event-loop tick. */ const realTick = Effect.promise(() => new Promise((resolve) => setImmediate(resolve))); -// Give forked fibers scheduling turns to subscribe/process (they are message-blocked on -// the provider stream, so TestClock.adjust alone does not run them), interleaved with -// real ticks so store persistence completes deterministically. -/** One pump of both schedulers: the Effect fiber scheduler and the real Node event loop. */ -const pump = Effect.gen(function* () { - yield* realTick; - for (let j = 0; j < 5; j++) yield* Effect.yieldNow; -}); +type ReceiptMatcher = (receipt: AutoResumeReactorReceipt) => boolean; -/** - * A bounded spin, for asserting that something does NOT happen. - * - * You cannot wait for the absence of an event, so these sites keep a fixed number of turns. - * Anything waiting for something to APPEAR must use `settleUntil` — see the note there. - */ -const settleQuiet = Effect.gen(function* () { - for (let i = 0; i < 10; i++) yield* pump; -}); +/** An arm was scheduled — optionally, the one due at `resumeAtMs`. */ +const scheduled = + (resumeAtMs?: number): ReceiptMatcher => + (receipt) => + receipt.type === "resume.scheduled" && + (resumeAtMs === undefined || receipt.resumeAtMs === resumeAtMs); + +/** A resume attempt is over. Published after the dispatch, which may itself have failed. */ +const fired: ReceiptMatcher = (receipt) => receipt.type === "resume.fired"; + +const firedFor = + (threadId: string): ReceiptMatcher => + (receipt) => + receipt.type === "resume.fired" && receipt.threadId === threadId; + +const scheduledFor = + (threadId: string): ReceiptMatcher => + (receipt) => + receipt.type === "resume.scheduled" && receipt.threadId === threadId; + +const cancelledFor = + (reason: string): ReceiptMatcher => + (receipt) => + receipt.type === "resume.cancelled" && receipt.reason === reason; + +const skippedFor = + (reason: string): ReceiptMatcher => + (receipt) => + receipt.type === "resume.skipped" && receipt.reason === reason; -const MAX_SETTLE_PUMPS = 500; +/** The next receipt matching `matches`. Unbounded on purpose: the test timeout is the bound. */ +const untilReceipt = (matches: ReceiptMatcher) => + Effect.gen(function* () { + const { log } = yield* AutoResumeReactorReceipts; + while (true) { + const receipt = yield* PubSub.take(log); + if (matches(receipt)) return receipt; + } + }); /** - * Wait until `condition` holds, pumping both schedulers. + * Advance one poll and wait out the wake pass it triggers, returning what that pass published. * - * Condition-based rather than a fixed spin, because no tick count is correct on every machine: - * `schedule()` gates its in-memory ref update behind `writeFileStringAtomically` — real - * filesystem I/O that completes on the Node event loop, NOT on TestClock — and how many turns - * that takes depends on the disk and the load. + * `tick.completed` is emitted after every due arm has been decided and every write is durable, + * so when this returns the world is exactly one whole pass further on — no more and no less, + * on any machine. * - * The fixed 10-pump spin this replaces failed about 1 run in 13 locally and took a main CI run - * red on 2026-08-07 with `expected +0 to equal 1`: the assertion simply looked before the write - * landed. Exiting as soon as the condition holds also makes the common case FASTER than the old - * spin, which always paid for all ten. + * The `realTick` is a handoff, not a budget: a pass publishes its receipt just before it + * re-arms its own `Effect.sleep`, so this gives it the turn it needs to register that sleep + * before the clock steps over it. */ -const settleUntil = (condition: Effect.Effect, description: string) => +const advanceOneTick = (log: PubSub.Subscription) => Effect.gen(function* () { - for (let i = 0; i < MAX_SETTLE_PUMPS; i++) { - if (yield* condition) return; - yield* pump; + yield* realTick; + yield* TestClock.adjust(Duration.millis(POLL_MS)); + const seen: AutoResumeReactorReceipt[] = []; + while (true) { + const receipt = yield* PubSub.take(log); + if (receipt.type === "tick.completed") return seen; + seen.push(receipt); } - return yield* Effect.die( - new Error(`timed out waiting for ${description} after ${MAX_SETTLE_PUMPS} pumps`), - ); }); -// Advance the test clock in wake-poll-sized steps, letting the wake fiber run each tick. -// A single large adjust does not reliably drive a recurring delay+forever loop. -const advancePastResume = Effect.gen(function* () { - for (let i = 0; i < 8; i++) { - yield* TestClock.adjust(Duration.millis(30_000)); - yield* settleQuiet; - } -}); +/** + * Advance by whole polls, waiting out each wake pass. + * + * This is how a scenario asserts that something does NOT happen: N passes demonstrably ran + * and none of them did it, rather than N clock steps and a hope that the reactor kept up. + */ +const advanceTicks = (ticks: number) => + Effect.gen(function* () { + const { log } = yield* AutoResumeReactorReceipts; + for (let i = 0; i < ticks; i++) yield* advanceOneTick(log); + }); /** - * Advance the clock until `condition` holds, for the cases that expect the wake fiber to fire. + * Advance in poll-sized steps until a wake pass publishes a matching receipt. * - * Same reasoning as `settleUntil`, one layer out: the wake fiber's work also ends in a store - * write, so "advance eight times and look" has the same race. Runs a generous number of steps - * because each one is a virtual 30s and costs only scheduler turns. + * `maxPolls` bounds SIMULATED time — how long the scenario is willing to wait — and says + * nothing about the machine, so it can never expire early under load. */ -const advanceUntil = (condition: Effect.Effect, description: string) => +const advanceUntilReceipt = (matches: ReceiptMatcher, description: string, maxPolls = 60) => Effect.gen(function* () { - for (let i = 0; i < 40; i++) { - if (yield* condition) return; - yield* TestClock.adjust(Duration.millis(30_000)); - yield* settleQuiet; + const { log } = yield* AutoResumeReactorReceipts; + for (let poll = 0; poll < maxPolls; poll++) { + const hit = (yield* advanceOneTick(log)).find(matches); + if (hit !== undefined) return hit; } - if (yield* condition) return; - return yield* Effect.die(new Error(`timed out waiting for ${description} past the resume`)); + return yield* Effect.die( + new Error(`timed out waiting for ${description} after ${maxPolls} polls`), + ); }); describe("AutoResumeReactor (integration)", () => { @@ -248,7 +299,7 @@ describe("AutoResumeReactor (integration)", () => { ]); yield* Effect.gen(function* () { - yield* settleUntil(scheduledOne(store), "detection to schedule a pending resume"); + yield* untilReceipt(scheduled()); const calls = yield* Ref.get(snapshotCalls); assert.isAbove( @@ -263,7 +314,7 @@ describe("AutoResumeReactor (integration)", () => { assert.include(types(afterSchedule), "thread.activity.append"); assert.notInclude(types(afterSchedule), "thread.turn.start"); - yield* advanceUntil(dispatchedIncludes(dispatched, "thread.turn.start"), "the resume turn"); + yield* advanceUntilReceipt(fired, "the resume turn"); const afterWake = yield* Ref.get(dispatched); const turnStarts = afterWake.filter((c) => c.type === "thread.turn.start"); @@ -281,25 +332,19 @@ describe("AutoResumeReactor (integration)", () => { // snapshot reports latestTurn: null. That must NOT read as "thread-advanced". it.effect("resumes when the limited turn has settled away by wake time (latestTurn null)", () => Effect.gen(function* () { - const { dispatched, modelRef, deps, store } = yield* harness( + const { dispatched, modelRef, deps } = yield* harness( readModel({ status: "running", latestTurn: { turnId: "turn-1", state: "running" } }), [rejectedEvent(100)], ); yield* Effect.gen(function* () { - yield* settleUntil(scheduledOne(store), "detection to schedule"); // baseline.latestTurnId === "turn-1" + yield* untilReceipt(scheduled()); // baseline.latestTurnId === "turn-1" // The limited turn settles and the session stops during the wait — the // projection's latest_turn_id empties out, so the snapshot's latestTurn is null. yield* Ref.set(modelRef, readModel({ status: "stopped", latestTurn: null })); - // Waiting for something to APPEAR, so it must be condition-based — see the note - // on `settleUntil`. The fixed spin this replaces was the file's last "expect a - // resume, then look" site and went red under load. - yield* advanceUntil( - dispatchedIncludes(dispatched, "thread.turn.start"), - "the settled-away thread to resume", - ); + yield* advanceUntilReceipt(fired, "the settled-away thread to resume"); const commands = yield* Ref.get(dispatched); const turnStarts = commands.filter((c) => c.type === "thread.turn.start"); @@ -325,13 +370,13 @@ describe("AutoResumeReactor (integration)", () => { it.effect("resumes when the new turn died inside the still-closed window", () => Effect.gen(function* () { - const { dispatched, modelRef, deps, store } = yield* harness( + const { dispatched, modelRef, deps } = yield* harness( readModel({ latestTurn: { turnId: "turn-old", state: "completed" } }), [rejectedEvent(100)], ); yield* Effect.gen(function* () { - yield* settleUntil(scheduledOne(store), "detection to schedule"); // baseline: "turn-old" + yield* untilReceipt(scheduled()); // baseline: "turn-old" // A turn is requested and dies half a second later, at t=150s — ten seconds // before the window reopens. It produced nothing; it only became the latest turn. @@ -346,10 +391,7 @@ describe("AutoResumeReactor (integration)", () => { }), ); - yield* advanceUntil( - dispatchedIncludes(dispatched, "thread.turn.start"), - "the doomed turn to be ignored and the resume to fire", - ); + yield* advanceUntilReceipt(fired, "the doomed turn to be ignored and the resume to fire"); const commands = yield* Ref.get(dispatched); assert.strictEqual( @@ -370,13 +412,13 @@ describe("AutoResumeReactor (integration)", () => { it.effect("still cancels when the new turn outlived the window (genuine advancement)", () => Effect.gen(function* () { - const { dispatched, modelRef, deps, store } = yield* harness( + const { dispatched, modelRef, deps } = yield* harness( readModel({ latestTurn: { turnId: "turn-old", state: "completed" } }), [rejectedEvent(100)], ); yield* Effect.gen(function* () { - yield* settleUntil(scheduledOne(store), "detection to schedule"); // baseline: "turn-old" + yield* untilReceipt(scheduled()); // baseline: "turn-old" // Same fixture, one field moved across the boundary: this turn finished at t=170s, // after the window reopened at 160s. The thread really did move on. @@ -391,11 +433,11 @@ describe("AutoResumeReactor (integration)", () => { }), ); - // The verdict has landed once the arm is gone, so this waits for something to - // APPEAR rather than spinning a fixed number of turns and looking. - yield* advanceUntil( - store.listPending.pipe(Effect.map((pending) => pending.length === 0)), - "the wake tick to reach a verdict on the arm", + // The guard names its own verdict, so the wait is that verdict rather than an + // inference from the arm having gone away. + yield* advanceUntilReceipt( + cancelledFor("thread-advanced"), + "the wake pass to reach a verdict on the arm", ); const commands = yield* Ref.get(dispatched); @@ -437,16 +479,13 @@ describe("AutoResumeReactor (integration)", () => { // fixed cadence the ticks landed on 150_000 then 180_000, so nothing fired at 160_000. it.effect("fires at the armed time, not at the next poll boundary", () => Effect.gen(function* () { - const { dispatched, deps, store } = yield* harness(readModel({}), [rejectedEvent(100)]); + const { dispatched, deps } = yield* harness(readModel({}), [rejectedEvent(100)]); yield* Effect.gen(function* () { - yield* settleUntil(scheduledOne(store), "detection to schedule"); // due at 160_000 + yield* untilReceipt(scheduled(160_000)); - // t = 150_000: still ten seconds short of the window reopening. - for (let i = 0; i < 5; i++) { - yield* TestClock.adjust(Duration.millis(30_000)); - yield* settleQuiet; - } + // t = 150_000: five whole wake passes, still ten seconds short of the reopening. + yield* advanceTicks(5); assert.notInclude( types(yield* Ref.get(dispatched)), "thread.turn.start", @@ -456,9 +495,11 @@ describe("AutoResumeReactor (integration)", () => { // t = 160_000 exactly: the armed moment, and a moment the old fixed cadence // never visited. yield* TestClock.adjust(Duration.millis(10_000)); - yield* settleUntil( - dispatchedIncludes(dispatched, "thread.turn.start"), - "the resume to fire at the armed time rather than the next poll", + yield* untilReceipt(fired); + assert.include( + types(yield* Ref.get(dispatched)), + "thread.turn.start", + "the resume must fire at the armed time rather than the next poll", ); }).pipe(Effect.provide(AutoResumeReactorLive.pipe(Layer.provideMerge(deps)))); }).pipe(Effect.scoped, Effect.provide(Layer.mergeAll(NodeServices.layer, TestClock.layer()))), @@ -469,12 +510,10 @@ describe("AutoResumeReactor (integration)", () => { // nothing, and the arm used to be destroyed as `user-took-over`. It must survive. it.effect("still resumes when the user posted a message while the resume was pending", () => Effect.gen(function* () { - const { dispatched, modelRef, deps, store } = yield* harness(readModel({}), [ - rejectedEvent(100), - ]); + const { dispatched, modelRef, deps } = yield* harness(readModel({}), [rejectedEvent(100)]); yield* Effect.gen(function* () { - yield* settleUntil(scheduledOne(store), "detection to schedule from the pre-loaded event"); + yield* untilReceipt(scheduled()); // A new user message lands, and goes nowhere: the thread is still idle at wake time. yield* Ref.set( @@ -487,7 +526,7 @@ describe("AutoResumeReactor (integration)", () => { }), ); - yield* advanceUntil(dispatchedIncludes(dispatched, "thread.turn.start"), "the resume turn"); + yield* advanceUntilReceipt(fired, "the resume turn"); const commands = yield* Ref.get(dispatched); assert.strictEqual( @@ -516,10 +555,7 @@ describe("AutoResumeReactor (integration)", () => { ]); yield* Effect.gen(function* () { - yield* settleUntil( - store.listPending.pipe(Effect.map((p) => p[0]?.resumeAtMs === 1_060_000)), - "the second rejection to supersede the first arm", - ); + yield* untilReceipt(scheduled(1_060_000)); assert.strictEqual( (yield* store.listPending).length, @@ -536,10 +572,10 @@ describe("AutoResumeReactor (integration)", () => { ]); // The original 160_000 due time passes without firing: that window is still shut. - yield* advancePastResume; // 8 x 30s = 240_000ms + yield* advanceTicks(8); // 8 x 30s = 240_000ms assert.notInclude(types(yield* Ref.get(dispatched)), "thread.turn.start"); - yield* advanceUntil(dispatchedIncludes(dispatched, "thread.turn.start"), "the resume turn"); + yield* advanceUntilReceipt(fired, "the superseded resume to fire at its new time"); }).pipe(Effect.provide(AutoResumeReactorLive.pipe(Layer.provideMerge(deps)))); }).pipe(Effect.scoped, Effect.provide(Layer.mergeAll(NodeServices.layer, TestClock.layer()))), ); @@ -552,10 +588,10 @@ describe("AutoResumeReactor (integration)", () => { yield* Ref.set(failTurnStart, true); // make the resume's turn.start dispatch fail yield* Effect.gen(function* () { - yield* settleUntil(scheduledOne(store), "detection to schedule"); - // Fires once and the dispatch fails; the attempt is reserved either way, which is the - // observable signal that the wake actually ran. - yield* advanceUntil(firedCount(store, 1), "the attempt to be reserved"); + yield* untilReceipt(scheduled()); + // Fires once and the dispatch fails; the attempt is reserved either way, and the + // receipt is published after the failed dispatch for exactly that reason. + yield* advanceUntilReceipt(fired, "the attempt to be reserved"); // The attempt was reserved (pending cleared, one fire recorded) despite the failure. assert.strictEqual((yield* store.listPending).length, 0); @@ -565,7 +601,7 @@ describe("AutoResumeReactor (integration)", () => { const turnStartsBefore = (yield* Ref.get(dispatched)).filter( (c) => c.type === "thread.turn.start", ).length; - yield* advancePastResume; + yield* advanceTicks(8); const turnStartsAfter = (yield* Ref.get(dispatched)).filter( (c) => c.type === "thread.turn.start", ).length; @@ -581,9 +617,10 @@ describe("AutoResumeReactor (integration)", () => { yield* store.setEnabled("thread-1", false); yield* Effect.gen(function* () { - // A bounded spin, not settleUntil: a disabled thread is gated BEFORE getSnapshot, so - // there is no positive signal to wait for — the assertion is that nothing appears. - yield* settleQuiet; // give detection every chance to run against the pre-loaded rejection + // Detection announces the gate it hit, so even "nothing was scheduled" is an exact + // await: the assertions below run after detection has demonstrably finished with the + // pre-loaded rejection, not after a spin that hoped it had. + yield* untilReceipt(skippedFor("disabled")); assert.strictEqual( (yield* store.listPending).length, @@ -593,7 +630,7 @@ describe("AutoResumeReactor (integration)", () => { // Disabling is a deliberate user action, so it must not post timeline noise either. assert.notInclude(types(yield* Ref.get(dispatched)), "thread.activity.append"); - yield* advancePastResume; + yield* advanceTicks(8); assert.notInclude(types(yield* Ref.get(dispatched)), "thread.turn.start"); }).pipe(Effect.provide(AutoResumeReactorLive.pipe(Layer.provideMerge(deps)))); }).pipe(Effect.scoped, Effect.provide(Layer.mergeAll(NodeServices.layer, TestClock.layer()))), @@ -604,14 +641,16 @@ describe("AutoResumeReactor (integration)", () => { const { dispatched, deps, store } = yield* harness(readModel({}), [rejectedEvent(100)]); yield* Effect.gen(function* () { - yield* settleUntil(scheduledOne(store), "detection to schedule"); + yield* untilReceipt(scheduled()); assert.strictEqual((yield* store.listPending).length, 1, "precondition: it scheduled"); // The switch is flipped off *after* scheduling but *before* the window reopens. // fireOne must re-read the record rather than trust the scheduling-time value. yield* store.setEnabled("thread-1", false); - yield* advancePastResume; + // The wait that made this test radroid/t3code#134: it was eight clock steps and a + // turn budget, and under load the cancellation had not landed when the assertion ran. + yield* advanceUntilReceipt(cancelledFor("disabled"), "the arm to be cancelled at wake"); assert.notInclude(types(yield* Ref.get(dispatched)), "thread.turn.start"); assert.strictEqual((yield* store.listPending).length, 0, "pending must be cleared"); @@ -623,4 +662,61 @@ describe("AutoResumeReactor (integration)", () => { }).pipe(Effect.provide(AutoResumeReactorLive.pipe(Layer.provideMerge(deps)))); }).pipe(Effect.scoped, Effect.provide(Layer.mergeAll(NodeServices.layer, TestClock.layer()))), ); + + // `fireOne` reads a fresh snapshot per arm and `getSnapshot` has a typed + // `ProjectionRepositoryError` channel, so one arm failing mid-pass is an expected event, + // not a defect. Two things must survive it: the rest of the batch, and the end-of-pass + // receipt — a pass that dies before `tick.completed` leaves anything awaiting that receipt + // waiting for one that will never come. + it.effect("one failed fire neither aborts the batch nor suppresses the completed pass", () => + Effect.gen(function* () { + const { dispatched, deps, store, failNextSnapshot } = yield* harness( + readModel({ threadIds: ["thread-1", "thread-2"] }), + [rejectedEvent(100, "thread-1"), rejectedEvent(100, "thread-2", "evt-2")], + ); + + yield* Effect.gen(function* () { + yield* untilReceipt(scheduledFor("thread-1")); + yield* untilReceipt(scheduledFor("thread-2")); + // Both due at 160_000, and in this order: the pass walks `listPending`, so the + // failure below lands on thread-1 and thread-2 is the rest of the batch. + assert.deepStrictEqual( + (yield* store.listPending).map((p) => p.threadId), + ["thread-1", "thread-2"], + ); + + // The next snapshot read — the first `fireOne` of the coming pass — fails. + yield* Ref.set(failNextSnapshot, true); + + // Both assertions at once: thread-2 fired despite thread-1 failing, and the pass + // reached `tick.completed` at all (without it `advanceUntilReceipt` could not return, + // because it drains to the tick receipt on every step). + yield* advanceUntilReceipt(firedFor("thread-2"), "the rest of the batch to fire anyway"); + + // Existing behaviour, pinned rather than changed: the failed fire reserved nothing and + // cleared nothing, so the arm is untouched and the next pass tries it again. + assert.deepStrictEqual( + (yield* store.listPending).map((p) => p.threadId), + ["thread-1"], + "a failed fire must leave the arm alone, not drop it", + ); + assert.strictEqual( + yield* store.countFiredSince("thread-1", 0), + 0, + "a failed fire must not burn one of the 24h attempts", + ); + + yield* advanceUntilReceipt(firedFor("thread-1"), "the failed arm to be retried"); + + const turnStarts = (yield* Ref.get(dispatched)).filter( + (c) => c.type === "thread.turn.start", + ); + assert.deepStrictEqual( + turnStarts.map((c) => c.threadId), + ["thread-2", "thread-1"], + "both threads resume: the healthy one first, the failed one on the next pass", + ); + }).pipe(Effect.provide(AutoResumeReactorLive.pipe(Layer.provideMerge(deps)))); + }).pipe(Effect.scoped, Effect.provide(Layer.mergeAll(NodeServices.layer, TestClock.layer()))), + ); }); diff --git a/apps/server/src/coil/autoResume/Reactor.ts b/apps/server/src/coil/autoResume/Reactor.ts index b15c09587593..2e3d5312a02f 100644 --- a/apps/server/src/coil/autoResume/Reactor.ts +++ b/apps/server/src/coil/autoResume/Reactor.ts @@ -47,6 +47,7 @@ import { isClaudeThread, threadIsGone, } from "./guards.ts"; +import { receiptEmitter } from "./receipts.ts"; import { AutoResumeStore, type PendingResume } from "./state.ts"; const HOUR_MS = 60 * 60_000; @@ -62,6 +63,8 @@ const makeSupervisor = Effect.gen(function* () { const store = yield* AutoResumeStore; const crypto = yield* Crypto.Crypto; const config = resolveConfig(); + // Optional, and absent in every production graph. See `receipts.ts`. + const receipts = yield* receiptEmitter; const isoNow = DateTime.now.pipe(Effect.map(DateTime.formatIso)); @@ -151,7 +154,10 @@ const makeSupervisor = Effect.gen(function* () { const record = yield* store.getThread(threadId); // Per-thread switch. A thread the user turned off never schedules, and posts no // timeline note: they disabled it deliberately, so an activity row would be noise. - if (!record.enabled) return; + if (!record.enabled) { + yield* receipts.emit({ type: "resume.skipped", threadId, reason: "disabled" }); + return; + } const nowMs = yield* Clock.currentTimeMillis; const firedRecently = yield* store.countFiredSince(threadId, nowMs - BACKOFF_LOOKBACK_MS); const firedInCapWindow = yield* store.countFiredSince(threadId, nowMs - CAP_WINDOW_MS); @@ -166,11 +172,17 @@ const makeSupervisor = Effect.gen(function* () { }); // A capped-out thread stops scheduling here (and stops posting "scheduled" notes) // until its fires age out of the 24h window; all other skips are dedupe/no-op. - if (plan.kind === "skip") return; + if (plan.kind === "skip") { + yield* receipts.emit({ type: "resume.skipped", threadId, reason: plan.reason }); + return; + } const snapshot = yield* snapshotQuery.getSnapshot(); const thread = snapshot.threads.find((t) => t.id === threadId); - if (!thread || threadIsGone(thread) || !isClaudeThread(thread)) return; + if (!thread || threadIsGone(thread) || !isClaudeThread(thread)) { + yield* receipts.emit({ type: "resume.skipped", threadId, reason: "thread-ineligible" }); + return; + } // Replacing an existing arm (a later window superseded it — see decide.ts). The // fresh `captureBaseline` below is what re-baselines the resume onto whatever the @@ -195,6 +207,12 @@ const makeSupervisor = Effect.gen(function* () { ? `Usage limit window pushed back (${limitType}). Auto-resume rescheduled to ~${waitMinutes} min from now.` : `Usage limit reached (${limitType}). Auto-resume scheduled in ~${waitMinutes} min.`, ); + yield* receipts.emit({ + type: "resume.scheduled", + threadId, + resumeAtMs: plan.resumeAtMs, + superseded, + }); }); // --- wake ----------------------------------------------------------------- @@ -207,6 +225,11 @@ const makeSupervisor = Effect.gen(function* () { const thread = snapshot.threads.find((t) => t.id === pending.threadId); if (!thread) { yield* store.clearPending(pending.threadId); + yield* receipts.emit({ + type: "resume.cancelled", + threadId: pending.threadId, + reason: "thread-missing", + }); return; } @@ -222,6 +245,11 @@ const makeSupervisor = Effect.gen(function* () { "coil.auto-resume.cancelled", "Auto-resume cancelled: turned off for this thread.", ); + yield* receipts.emit({ + type: "resume.cancelled", + threadId: pending.threadId, + reason: "disabled", + }); return; } @@ -246,6 +274,7 @@ const makeSupervisor = Effect.gen(function* () { observedTurnCompletedAt: thread.latestTurn?.completedAt ?? null, }, ); + yield* receipts.emit({ type: "resume.cancelled", threadId: pending.threadId, reason }); return; } @@ -258,6 +287,11 @@ const makeSupervisor = Effect.gen(function* () { "coil.auto-resume.capped", `Auto-resume stopped after ${config.maxResumesPer24h} attempts in 24h.`, ); + yield* receipts.emit({ + type: "resume.cancelled", + threadId: pending.threadId, + reason: "capped", + }); return; } @@ -285,13 +319,43 @@ const makeSupervisor = Effect.gen(function* () { }), ), ); + // After the dispatch attempt, not before: a failed dispatch still burned the attempt, + // and the test that asserts exactly that needs to be told when the attempt is over. + yield* receipts.emit({ + type: "resume.fired", + threadId: pending.threadId, + attempt: firedIn24h + 1, + }); }); const processDue = Effect.gen(function* () { const nowMs = yield* Clock.currentTimeMillis; const pending = yield* store.listPending; const due = pending.filter((p) => p.resumeAtMs <= nowMs); - yield* Effect.forEach(due, (p) => fireOne(p, nowMs), { discard: true }); + yield* Effect.forEach( + due, + (p) => + fireOne(p, nowMs).pipe( + // Per arm, so one bad record cannot abort the rest of the batch — and cannot skip + // the end-of-pass receipt below. `fireOne` reads a fresh snapshot per item and + // `getSnapshot` has a typed `ProjectionRepositoryError` channel, so this is an + // expected failure, not just a defect. The arm is left as it was: nothing was + // reserved and nothing was cleared, so the next pass tries it again. + Effect.catchCause((cause) => + Effect.logWarning("coil auto-resume: fire failed", { + threadId: p.threadId, + cause: Cause.pretty(cause), + }), + ), + ), + { discard: true }, + ); + // The end-of-pass receipt, published on the empty path too: it is what makes "advance one + // poll and let the reactor finish" an exact await instead of a budget of scheduler turns. + // Assembled only when someone is listening — this runs every poll of every install. + if (receipts.enabled) { + yield* receipts.emit({ type: "tick.completed", dueCount: due.length, nowMs }); + } }); /** diff --git a/apps/server/src/coil/autoResume/receipts.ts b/apps/server/src/coil/autoResume/receipts.ts new file mode 100644 index 000000000000..cc207b7e64c9 --- /dev/null +++ b/apps/server/src/coil/autoResume/receipts.ts @@ -0,0 +1,118 @@ +/** + * The auto-resume supervisor's receipts. + * + * Every milestone this reactor reaches is asynchronous and most end in real filesystem I/O, + * so a test that wants to assert about one has two choices: infer it from a store read after + * spinning the scheduler, or be told. Inference is what this file replaces — `Reactor.test.ts` + * went red in CI (radroid/t3code#134) on `pending must be cleared: expected 1 to equal +0`, + * which was not a defect in the product but a wait that ran out of scheduler turns on a + * loaded runner. AGENTS.md is explicit about the remedy: wait on receipts and worker drains, + * never on sleeps or polling. + * + * So the reactor announces, exactly as the loop supervisor does (`coil/loop/receipts.ts`), + * and the tests await the announcement instead of guessing how many turns it should take. + * + * ## Nothing is paid for in production + * + * The service is **optional**. `receiptEmitter` resolves it with `Effect.serviceOption`, and + * when nobody has provided it — which is every production layer graph, since `coil/index.ts` + * does not mention this module — `emit` is a constant `Effect.void` and `enabled` is false. + * Callers use `enabled` to skip even assembling a receipt. + * + * The buffer is bounded and **dropping**: a subscriber that stops draining loses receipts + * rather than blocking the wake fiber. A supervisor that stalls because a test stopped + * listening would be a worse bug than the flake this replaces. + * + * @module coil/autoResume/receipts + */ + +import * as Context from "effect/Context"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; +import * as Option from "effect/Option"; +import * as PubSub from "effect/PubSub"; +import * as Scope from "effect/Scope"; + +/** + * One announced milestone. + * + * Every variant names something a test used to infer from the store or from the recorded + * command list. `tick.completed` is the important one: it is published at the **end** of + * every wake pass, including the one where nothing is due, so "advance one poll and let the + * reactor finish" is a single exact await rather than a budget of scheduler turns. + */ +export type AutoResumeReactorReceipt = + /** A whole wake pass is over: every due arm decided, every write durable. */ + | { readonly type: "tick.completed"; readonly dueCount: number; readonly nowMs: number } + /** Detection armed a resume. `superseded` is true when it replaced an existing arm. */ + | { + readonly type: "resume.scheduled"; + readonly threadId: string; + readonly resumeAtMs: number; + readonly superseded: boolean; + } + /** Detection saw a rejection and armed nothing; `reason` is the gate that stopped it. */ + | { readonly type: "resume.skipped"; readonly threadId: string; readonly reason: string } + /** An armed resume was dropped at wake time; `reason` is the guard that dropped it. */ + | { readonly type: "resume.cancelled"; readonly threadId: string; readonly reason: string } + /** The resume turn went out (or was attempted — a failed dispatch still burns the attempt). */ + | { readonly type: "resume.fired"; readonly threadId: string; readonly attempt: number }; + +export interface AutoResumeReactorReceiptsShape { + readonly publish: (receipt: AutoResumeReactorReceipt) => Effect.Effect; + /** + * Everything published for the life of the service, in order. + * + * Subscribed at construction rather than handed out per wait, because the service is built + * before the reactor is: detection can publish while the test body is still being + * assembled, and a subscription opened later would miss it and wait forever. + */ + readonly log: PubSub.Subscription; +} + +export class AutoResumeReactorReceipts extends Context.Service< + AutoResumeReactorReceipts, + AutoResumeReactorReceiptsShape +>()("t3/coil/autoResume/receipts/AutoResumeReactorReceipts") {} + +/** + * Room for every receipt a scenario can produce without draining. + * + * Scenarios run at most a few hundred simulated polls and drain the tick receipt on each one, + * so this is slack rather than a working limit — but it is a limit, because the alternative + * is a queue that grows with a stuck subscriber. + */ +const RECEIPT_BUFFER = 4096; + +export const makeAutoResumeReactorReceipts: Effect.Effect< + AutoResumeReactorReceiptsShape, + never, + Scope.Scope +> = Effect.gen(function* () { + const pubsub = yield* PubSub.dropping(RECEIPT_BUFFER); + const log = yield* PubSub.subscribe(pubsub); + return { + publish: (receipt) => PubSub.publish(pubsub, receipt).pipe(Effect.asVoid), + log, + }; +}); + +export const AutoResumeReactorReceiptsLive = Layer.effect( + AutoResumeReactorReceipts, + makeAutoResumeReactorReceipts, +); + +const noEmit = (_receipt: AutoResumeReactorReceipt): Effect.Effect => Effect.void; + +export interface AutoResumeReceiptEmitter { + /** Whether anyone is listening. Lets a caller skip assembling a receipt, not just sending. */ + readonly enabled: boolean; + readonly emit: (receipt: AutoResumeReactorReceipt) => Effect.Effect; +} + +/** Resolve the optional service into something the reactor can call unconditionally. */ +export const receiptEmitter: Effect.Effect = Effect.gen(function* () { + const service = yield* Effect.serviceOption(AutoResumeReactorReceipts); + if (Option.isNone(service)) return { enabled: false, emit: noEmit }; + return { enabled: true, emit: service.value.publish }; +}); diff --git a/apps/server/src/coil/autoResume/replay/incidentDoomedTurn.test.ts b/apps/server/src/coil/autoResume/replay/incidentDoomedTurn.test.ts index b1a9b89fd14e..7e90a9213b19 100644 --- a/apps/server/src/coil/autoResume/replay/incidentDoomedTurn.test.ts +++ b/apps/server/src/coil/autoResume/replay/incidentDoomedTurn.test.ts @@ -141,6 +141,16 @@ describe("incident 2026-08-18: doomed turn destroys the pending resume", () => { // 19:51:00 — the window has reopened and the resume is due. yield* advanceSteps(10); + // The clock is where it needs to be; now wait for the verdict rather than assuming + // ten quiet spins were enough for it to land (see `settleUntil`). + yield* settleUntil( + Ref.get(harness.dispatched).pipe( + Effect.map((commands) => + appendedActivityKinds(commands).includes("coil.auto-resume.resumed"), + ), + ), + "the reopened window to fire the resume", + ); const commands = yield* Ref.get(harness.dispatched); @@ -190,6 +200,16 @@ describe("incident 2026-08-18: doomed turn destroys the pending resume", () => { ]), ); yield* advanceSteps(10); + // Same rule as above: the cancellation is something that must APPEAR, so wait for it. + // A fixed spin looked too early on a loaded runner (PR #143's CI run 34734274701). + yield* settleUntil( + Ref.get(harness.dispatched).pipe( + Effect.map((commands) => + appendedActivityKinds(commands).includes("coil.auto-resume.cancelled"), + ), + ), + "the takeover to cancel the pending resume", + ); const commands = yield* Ref.get(harness.dispatched); expect(appendedActivitySummaries(commands)).toContain(