From e0ce98c893d573d9a211ff49299195ad3f4a7343 Mon Sep 17 00:00:00 2001 From: Abhinav Pola Date: Wed, 23 Sep 2026 01:42:52 +0000 Subject: [PATCH 1/3] fix(agent): close two teardown races found by TLA+ model checking ReusableReadableStream.createConsumer() during cancel()'s await of the source reader registered a consumer at a watermark the second backlog sweep then cleared, so its first read threw the buffer invariant error. cancel() now marks the stream cancelled synchronously and later consumers are done immediately. handleRunEndAsyncTasks('drain') only called dropSettledTasks() when work was still in flight, so a task that settled during the last drain turn was never harvested and stayed persisted as working. The final drop now runs unconditionally. Adds specs/tla with the models and TLC configs that found both. --- .changeset/quiet-streams-drain.md | 5 + packages/agent/src/lib/model-result.ts | 3 +- packages/agent/src/lib/reusable-stream.ts | 26 +- .../tests/unit/async-tool-background.test.ts | 83 ++++ .../agent/tests/unit/reusable-stream.test.ts | 29 ++ specs/tla/AsyncRunEndDrain.cfg | 8 + specs/tla/AsyncRunEndDrain.tla | 156 ++++++++ specs/tla/AsyncRunEndDrain_timeouts.cfg | 8 + specs/tla/MultiConsumerStream.tla | 371 ++++++++++++++++++ specs/tla/MultiConsumerStream_active.cfg | 14 + .../MultiConsumerStream_active_latecreate.cfg | 14 + ...MultiConsumerStream_broadcaster_active.cfg | 14 + .../MultiConsumerStream_broadcaster_full.cfg | 14 + specs/tla/MultiConsumerStream_full.cfg | 14 + specs/tla/README.md | 20 + 15 files changed, 766 insertions(+), 13 deletions(-) create mode 100644 .changeset/quiet-streams-drain.md create mode 100644 specs/tla/AsyncRunEndDrain.cfg create mode 100644 specs/tla/AsyncRunEndDrain.tla create mode 100644 specs/tla/AsyncRunEndDrain_timeouts.cfg create mode 100644 specs/tla/MultiConsumerStream.tla create mode 100644 specs/tla/MultiConsumerStream_active.cfg create mode 100644 specs/tla/MultiConsumerStream_active_latecreate.cfg create mode 100644 specs/tla/MultiConsumerStream_broadcaster_active.cfg create mode 100644 specs/tla/MultiConsumerStream_broadcaster_full.cfg create mode 100644 specs/tla/MultiConsumerStream_full.cfg create mode 100644 specs/tla/README.md diff --git a/.changeset/quiet-streams-drain.md b/.changeset/quiet-streams-drain.md new file mode 100644 index 00000000..764bd088 --- /dev/null +++ b/.changeset/quiet-streams-drain.md @@ -0,0 +1,5 @@ +--- +'@openrouter/agent': patch +--- + +Fix two run-teardown races found by TLA+ model checking. `ReusableReadableStream.createConsumer()` called while `cancel()` is awaiting the source reader now yields an already-done consumer instead of one that reads a cleared buffer slot and throws. `onRunEnd: 'drain'` now persists and broadcasts (`delivery: 'dropped'`) a background task settlement that lands during the final drain turn, where it was previously left unharvested with the persisted task still `working`. diff --git a/packages/agent/src/lib/model-result.ts b/packages/agent/src/lib/model-result.ts index bbc4f047..a9a510fe 100644 --- a/packages/agent/src/lib/model-result.ts +++ b/packages/agent/src/lib/model-result.ts @@ -4983,11 +4983,10 @@ export class ModelResult< } } - // Drain budget exhausted with work still in flight: cut it loose. if (registry.hasInFlight()) { registry.abortAll('Async tool drain budget exhausted at run end'); - await this.dropSettledTasks(); } + await this.dropSettledTasks(); return response; } diff --git a/packages/agent/src/lib/reusable-stream.ts b/packages/agent/src/lib/reusable-stream.ts index fdd4ba41..ea0606b1 100644 --- a/packages/agent/src/lib/reusable-stream.ts +++ b/packages/agent/src/lib/reusable-stream.ts @@ -28,6 +28,7 @@ export class ReusableReadableStream { private sourceComplete = false; private sourceError: Error | null = null; private pumpStarted = false; + private cancelled = false; private sourceCancelPromise: Promise | null = null; private readonly streamReplay: StreamReplay; private readonly onValue: ((value: T) => void) | undefined; @@ -73,20 +74,22 @@ export class ReusableReadableStream { * Create a new consumer that can independently iterate over the stream. * Full-replay consumers start at position 0. Active-consumer replay starts * at the current trim watermark. Multiple attached consumers advance - * independently in either mode. + * independently in either mode. Consumers created after `cancel()` are + * already done. */ createConsumer(): AsyncIterableIterator { const consumerId = this.nextConsumerId++; - const state: ConsumerState = { - position: this.trimOffset, - waitingPromise: null, - cancelled: false, - }; - this.consumers.set(consumerId, state); - - // Start pumping the source stream if not already started - if (!this.pumpStarted) { - this.startPump(); + if (!this.cancelled) { + this.consumers.set(consumerId, { + position: this.trimOffset, + waitingPromise: null, + cancelled: false, + }); + + // Start pumping the source stream if not already started + if (!this.pumpStarted) { + this.startPump(); + } } // eslint-disable-next-line @typescript-eslint/no-this-alias @@ -343,6 +346,7 @@ export class ReusableReadableStream { * Cancel the source stream and all consumers */ async cancel(): Promise { + this.cancelled = true; // Cancel all consumers for (const consumer of this.consumers.values()) { consumer.cancelled = true; diff --git a/packages/agent/tests/unit/async-tool-background.test.ts b/packages/agent/tests/unit/async-tool-background.test.ts index 8a95b4e4..d989addc 100644 --- a/packages/agent/tests/unit/async-tool-background.test.ts +++ b/packages/agent/tests/unit/async-tool-background.test.ts @@ -480,6 +480,89 @@ describe('tool.background — placeholder & delivery', () => { expect(sawAbort).toBe(true); }); + it("onRunEnd: 'drain' reports a task that settles during the last drain turn as dropped", async () => { + const first = makeControlledBackgroundTool('render_a'); + const second = makeControlledBackgroundTool('render_b'); + + mockBetaResponsesSend + .mockResolvedValueOnce({ + ok: true, + value: makeResponse('resp_1', [ + functionCallItem('call_a', 'render_a', '{"script":"a"}'), + functionCallItem('call_b', 'render_b', '{"script":"b"}'), + ]), + }) + .mockImplementationOnce(async () => { + first.release({ + url: 'https://cdn/a.mp4', + }); + return { + ok: true, + value: makeResponse('resp_2', [ + messageItem('msg_1', 'both started'), + ]), + }; + }) + // The only drain turn (maxDrainTurns: 1) carries call_a; call_b + // settles while the model is still answering, after the loop has + // spent its last turn. + .mockImplementationOnce(async () => { + second.release({ + url: 'https://cdn/b.mp4', + }); + await new Promise((resolve) => setTimeout(resolve, 10)); + return { + ok: true, + value: makeResponse('resp_3', [ + messageItem('msg_2', 'a is done'), + ]), + }; + }); + + const result = callModel(client, { + model: 'test-model', + input: 'render', + tools: [ + first.tool, + second.tool, + ] as const, + asyncTools: { + onRunEnd: 'drain', + drainTimeoutMs: 5_000, + maxDrainTurns: 1, + }, + }); + + const settled: Array<{ + toolCallId: string; + delivery: string; + }> = []; + for await (const event of result.getFullResponsesStream()) { + if (isToolAsyncSettledEvent(event)) { + settled.push({ + toolCallId: event.toolCallId, + delivery: event.delivery, + }); + } + } + + expect(mockBetaResponsesSend).toHaveBeenCalledTimes(3); + expect(settled).toEqual([ + { + toolCallId: 'call_a', + delivery: 'injected', + }, + { + toolCallId: 'call_b', + delivery: 'dropped', + }, + ]); + expect(result.getAsyncTasks().map((t) => t.status)).toEqual([ + 'completed', + 'completed', + ]); + }); + it('cancelTask(taskId) cancels a working background task', async () => { const controlled = makeControlledBackgroundTool('render_video'); diff --git a/packages/agent/tests/unit/reusable-stream.test.ts b/packages/agent/tests/unit/reusable-stream.test.ts index bd27dd40..161e8402 100644 --- a/packages/agent/tests/unit/reusable-stream.test.ts +++ b/packages/agent/tests/unit/reusable-stream.test.ts @@ -209,4 +209,33 @@ describe('ReusableReadableStream', () => { expect((await fresh.next()).done).toBe(true); // No source.close(): cancel() already terminated the source stream. }); + + it('active-consumers: a consumer created while cancel() awaits the source reader is done', async () => { + const source = controlledStream(); + const stream = new ReusableReadableStream(source.stream, { + streamReplay: 'active-consumers', + }); + const first = stream.createConsumer(); + await Promise.resolve(); + source.push(1); + + /* + * The pump has already read a chunk that lands in the buffer while + * cancel() awaits sourceReader.cancel(), so the late consumer is created + * between the two backlog sweeps. + */ + const cancelPromise = stream.cancel(); + const late = stream.createConsumer(); + await cancelPromise; + + expect(await first.next()).toEqual({ + done: true, + value: undefined, + }); + expect(await late.next()).toEqual({ + done: true, + value: undefined, + }); + expect(stream.findLastBuffered(() => true)).toBeUndefined(); + }); }); diff --git a/specs/tla/AsyncRunEndDrain.cfg b/specs/tla/AsyncRunEndDrain.cfg new file mode 100644 index 00000000..f782979d --- /dev/null +++ b/specs/tla/AsyncRunEndDrain.cfg @@ -0,0 +1,8 @@ +SPECIFICATION Spec +CONSTANTS + Tasks = {t1, t2, t3} + MaxDrainTurns = 2 + DrainTimeouts = FALSE + FixDropLeftover = TRUE +INVARIANT Safety +PROPERTY EventuallyDone diff --git a/specs/tla/AsyncRunEndDrain.tla b/specs/tla/AsyncRunEndDrain.tla new file mode 100644 index 00000000..8f2c0d20 --- /dev/null +++ b/specs/tla/AsyncRunEndDrain.tla @@ -0,0 +1,156 @@ +---------------------------- MODULE AsyncRunEndDrain ---------------------------- +(* + * Model of ModelResult.handleRunEndAsyncTasks (onRunEnd: 'drain') together + * with AsyncToolRegistry settlement, for background tool tasks that outlived + * their grace window and hold a pending placeholder in the transcript. + * + * Every JavaScript synchronous section is one atomic step. Background work + * settles at any await boundary (Settle action), which is exactly what the + * registry's `.then()` callbacks do at runtime. + * + * Property checked: when the run ends, every task that settled has been + * accounted for, either delivered in a drain turn or dropped and persisted by + * dropSettledTasks. The registry's settled queue is empty at run end. + *) +EXTENDS Naturals, FiniteSets + +CONSTANTS + Tasks, \* set of background tasks with a pending placeholder + MaxDrainTurns, \* asyncTools.maxDrainTurns + DrainTimeouts, \* TRUE: drainTimeoutMs may expire (registry.drain returns false) + FixDropLeftover \* TRUE: run end drops unharvested settlements even with nothing in flight + +VARIABLES + status, \* task -> "working" | "settled" + queued, \* task -> BOOLEAN, settlement sits in registry.settledQueue + outcome, \* task -> "none" | "delivered" | "dropped" + pc, \* program counter of handleRunEndAsyncTasks + drainTurn, \* loop counter + harvested, \* tasks taken by the current flushAsyncToolDeliveries + deadlinePassed + +vars == <> + +InFlight == \E t \in Tasks : status[t] = "working" +Unharvested == \E t \in Tasks : queued[t] + +TypeOK == + /\ status \in [Tasks -> {"working", "settled"}] + /\ queued \in [Tasks -> BOOLEAN] + /\ outcome \in [Tasks -> {"none", "delivered", "dropped"}] + /\ pc \in {"start", "loopTop", "awaitDrain", "harvest", "modelTurn", "afterTurn", + "finalize", "done"} + /\ drainTurn \in 0..MaxDrainTurns + /\ harvested \subseteq Tasks + /\ deadlinePassed \in BOOLEAN + +Init == + /\ status = [t \in Tasks |-> "working"] + /\ queued = [t \in Tasks |-> FALSE] + /\ outcome = [t \in Tasks |-> "none"] + /\ pc = "start" + /\ drainTurn = 0 + /\ harvested = {} + /\ deadlinePassed = FALSE + +\* Background work resolves or rejects: registry.settle pushes onto settledQueue. +Settle(t) == + /\ status[t] = "working" + /\ pc # "done" + /\ status' = [status EXCEPT ![t] = "settled"] + /\ queued' = [queued EXCEPT ![t] = TRUE] + /\ UNCHANGED <> + +\* Wall clock: drainTimeoutMs elapses at some point. +Deadline == + /\ DrainTimeouts + /\ ~deadlinePassed + /\ deadlinePassed' = TRUE + /\ UNCHANGED <> + +\* Early return: nothing in flight and nothing unharvested. +Start == + /\ pc = "start" + /\ pc' = IF ~InFlight /\ ~Unharvested THEN "done" ELSE "loopTop" + /\ UNCHANGED <> + +LoopTop == + /\ pc = "loopTop" + /\ pc' = IF drainTurn >= MaxDrainTurns THEN "finalize" + ELSE IF ~Unharvested /\ InFlight + THEN IF deadlinePassed THEN "finalize" ELSE "awaitDrain" + ELSE "harvest" + /\ UNCHANGED <> + +\* registry.drain(remaining) resumes on a settle event or on timeout. +AwaitDrain == + /\ pc = "awaitDrain" + /\ (Unharvested \/ deadlinePassed \/ ~InFlight) + /\ pc' = "harvest" + /\ UNCHANGED <> + +\* flushAsyncToolDeliveries: takeSettled(), inject envelope. Nothing settled: stop. +Harvest == + /\ pc = "harvest" + /\ LET taken == {t \in Tasks : queued[t]} + IN /\ harvested' = taken + /\ queued' = [t \in Tasks |-> FALSE] + /\ outcome' = [t \in Tasks |-> IF t \in taken THEN "delivered" ELSE outcome[t]] + /\ pc' = IF taken = {} THEN "finalize" ELSE "modelTurn" + /\ UNCHANGED <> + +\* retryCurrentRequest + saveResponseToState (awaits, tasks may settle meanwhile). +ModelTurn == + /\ pc = "modelTurn" + /\ pc' = "afterTurn" + /\ UNCHANGED <> + +AfterTurn == + /\ pc = "afterTurn" + /\ drainTurn' = drainTurn + 1 + /\ pc' = IF ~InFlight /\ ~Unharvested THEN "finalize" ELSE "loopTop" + /\ UNCHANGED <> + +\* Drain budget exhausted: abortAll settles working tasks as cancelled, then +\* dropSettledTasks persists every queued settlement. In the unfixed code the +\* drop only runs when something was still in flight. +Finalize == + /\ pc = "finalize" + /\ LET drops == InFlight \/ FixDropLeftover + settledNow == [t \in Tasks |-> "settled"] + queuedNow == [t \in Tasks |-> queued[t] \/ status[t] = "working"] + IN IF InFlight + THEN /\ status' = settledNow + /\ queued' = [t \in Tasks |-> FALSE] + /\ outcome' = [t \in Tasks |-> IF queuedNow[t] THEN "dropped" ELSE outcome[t]] + ELSE IF drops + THEN /\ status' = status + /\ queued' = [t \in Tasks |-> FALSE] + /\ outcome' = [t \in Tasks |-> IF queued[t] THEN "dropped" ELSE outcome[t]] + ELSE UNCHANGED <> + /\ pc' = "done" + /\ UNCHANGED <> + +Terminating == pc = "done" /\ UNCHANGED vars + +Next == + \/ \E t \in Tasks : Settle(t) + \/ Deadline + \/ Start \/ LoopTop \/ AwaitDrain \/ Harvest \/ ModelTurn \/ AfterTurn \/ Finalize + \/ Terminating + +Spec == Init /\ [][Next]_vars /\ WF_vars(Next) + +\* Every settlement is delivered or persisted by the time the run returns. +NoLeakedSettlement == + pc = "done" => \A t \in Tasks : ~queued[t] /\ (status[t] = "settled" => outcome[t] # "none") + +\* Delivery happens at most once and only for settled tasks. +DeliveryConsistent == + \A t \in Tasks : outcome[t] # "none" => status[t] = "settled" + +Safety == TypeOK /\ NoLeakedSettlement /\ DeliveryConsistent + +EventuallyDone == <>(pc = "done") + +============================================================================= diff --git a/specs/tla/AsyncRunEndDrain_timeouts.cfg b/specs/tla/AsyncRunEndDrain_timeouts.cfg new file mode 100644 index 00000000..62ab9e76 --- /dev/null +++ b/specs/tla/AsyncRunEndDrain_timeouts.cfg @@ -0,0 +1,8 @@ +SPECIFICATION Spec +CONSTANTS + Tasks = {t1, t2, t3} + MaxDrainTurns = 2 + DrainTimeouts = TRUE + FixDropLeftover = TRUE +INVARIANT Safety +PROPERTY EventuallyDone diff --git a/specs/tla/MultiConsumerStream.tla b/specs/tla/MultiConsumerStream.tla new file mode 100644 index 00000000..a8a83ff9 --- /dev/null +++ b/specs/tla/MultiConsumerStream.tla @@ -0,0 +1,371 @@ +-------------------------- MODULE MultiConsumerStream -------------------------- +(* + * Model of the multi-consumer buffer shared by + * packages/agent/src/lib/reusable-stream.ts (ReusableReadableStream) + * packages/agent/src/lib/tool-event-broadcaster.ts (ToolEventBroadcaster) + * + * One producer ("pump") appends values to a buffer. Consumers hold an + * absolute position; the buffer index of a consumer is + * head + position - trim + * Consumers that find nothing to read install a waiting promise and resume + * later (a separate atomic step, since resumption is a microtask). In + * 'active-consumers' replay mode the buffer is trimmed to the slowest live + * consumer, and the retained backlog is dropped when the last consumer + * leaves. ReusableReadableStream.cancel() is a two-phase operation with an + * await in the middle; both phases are modelled. + * + * Atomicity follows JavaScript: every synchronous section between two + * awaits is one action. + *) +EXTENDS Naturals, Sequences, FiniteSets, TLC + +CONSTANTS + Variant, \* "stream" (ReusableReadableStream) | "broadcaster" (ToolEventBroadcaster) + Consumers, \* set of consumer ids + MaxValues, \* bound on values the pump produces + Mode, \* "full" | "active" + CompactionMinHead, \* BUFFER_COMPACTION_MIN_HEAD (1024 in code; small here) + AllowStreamCancel, \* model ReusableReadableStream.cancel() + AllowLateCreate, \* allow createConsumer() after cancel() started + CancelledFlag \* fix: createConsumer() after cancel() registers nothing + +ASSUME Mode \in {"full", "active"} +ASSUME Variant \in {"stream", "broadcaster"} +ASSUME Variant = "broadcaster" => ~AllowStreamCancel + +CLEARED == 0 \* `undefined` slot left by buffer.fill(undefined, ...) + +VARIABLES + buf, \* Seq of values (1..MaxValues) or CLEARED + head, \* bufferHead + trim, \* trimOffset + produced, \* number of values pushed so far + complete, \* sourceComplete / isComplete + err, \* sourceError / completionError set + nextId, \* nextConsumerId (only used for the `=== 0` check) + pump, \* "running" | "done" + cons, \* consumer records + pc, \* consumer program counters + recv, \* values received by each consumer + joinPos, \* position at createConsumer() + joinedLate, \* createConsumer() happened after completion + cancelPhase, \* "none" | "awaitingSource" | "done" + cleanupPending \* broadcaster: queueMicrotask(() => this.cleanup()) outstanding + +vars == <> + +ConsumerRecord == [inMap: BOOLEAN, pos: Nat, waiting: BOOLEAN, + resolved: BOOLEAN, rejected: BOOLEAN] + +PCs == {"notStarted", "idle", "waiting", "done", "errored", "cancelled"} + +Live(c) == cons[c].inMap + +LiveConsumers == {c \in Consumers : Live(c)} + +Idx(c) == head + cons[c].pos - trim \* JS bufferIndex (0-based) + +HasUnread(c) == Idx(c) < Len(buf) \* JS: bufferIndex < buffer.length + +------------------------------------------------------------------------------ +(* Helpers mirroring the private methods *) + +\* dropUnreadBacklog(): returns [buf, head, trim] +DropUnreadBacklog(b, h, t, nid) == + IF nid = 0 \/ Len(b) - h = 0 + THEN <> + ELSE <<<<>>, 0, t + (Len(b) - h)>> + +Min(S) == CHOOSE m \in S : \A x \in S : m <= x + +\* trimConsumed() given the consumer map `cs`; returns [buf, head, trim] +TrimConsumed(b, h, t, cs, nid) == + IF Mode = "full" THEN <> + ELSE + LET live == {c \in Consumers : cs[c].inMap} + IN IF live = {} THEN DropUnreadBacklog(b, h, t, nid) + ELSE + LET m == Min({cs[c].pos : c \in live}) + nextHead == h + m - t + IN IF nextHead <= h THEN <> + ELSE IF nextHead = Len(b) THEN <<<<>>, 0, m>> + ELSE + LET filled == [i \in 1..Len(b) |-> + IF i > h /\ i <= nextHead THEN CLEARED ELSE b[i]] + IN IF nextHead >= CompactionMinHead /\ nextHead * 2 >= Len(b) + THEN <> + ELSE <> + +\* ToolEventBroadcaster.cleanup(): drop the buffer once complete with no +\* consumers left (active-consumers mode only). Note: trimOffset is NOT +\* advanced here, unlike dropUnreadBacklog(). +Cleanup(b, h, cs, isComplete) == + IF Variant = "broadcaster" /\ Mode = "active" /\ isComplete + /\ {c \in Consumers : cs[c].inMap} = {} + THEN <<<<>>, 0>> + ELSE <> + +\* notifyAllConsumers(): resolve or reject every waiting promise +Notify(cs, isErr) == + [c \in Consumers |-> + IF cs[c].inMap /\ cs[c].waiting + THEN [cs[c] EXCEPT !.waiting = FALSE, + !.resolved = ~isErr, + !.rejected = isErr] + ELSE cs[c]] + +------------------------------------------------------------------------------ +Init == + /\ buf = <<>> + /\ head = 0 + /\ trim = 0 + /\ produced = 0 + /\ complete = FALSE + /\ err = FALSE + /\ nextId = 0 + /\ pump = "running" + /\ cons = [c \in Consumers |-> [inMap |-> FALSE, pos |-> 0, waiting |-> FALSE, + resolved |-> FALSE, rejected |-> FALSE]] + /\ pc = [c \in Consumers |-> "notStarted"] + /\ recv = [c \in Consumers |-> <<>>] + /\ joinPos = [c \in Consumers |-> 0] + /\ joinedLate = [c \in Consumers |-> FALSE] + /\ cancelPhase = "none" + /\ cleanupPending = FALSE + +------------------------------------------------------------------------------ +(* Pump *) + +\* the stream's pump only starts with the first createConsumer(); the +\* broadcaster is pushed into regardless of consumers. +PumpStarted == Variant = "broadcaster" \/ nextId > 0 + +Push == + /\ pump = "running" + /\ PumpStarted + /\ ~complete /\ ~err + /\ produced < MaxValues + /\ buf' = Append(buf, produced + 1) + /\ produced' = produced + 1 + /\ cons' = Notify(cons, FALSE) + /\ UNCHANGED <> + +Complete == + /\ pump = "running" + /\ PumpStarted + /\ complete' = TRUE + /\ pump' = "done" + /\ cons' = Notify(cons, FALSE) + /\ cleanupPending' = (Variant = "broadcaster") + /\ UNCHANGED <> + +\* stream: pump error. broadcaster: complete(error). +Fail == + /\ pump = "running" + /\ PumpStarted + /\ cancelPhase = "none" + /\ err' = TRUE + /\ complete' = (Variant = "broadcaster") + /\ pump' = "done" + /\ cons' = Notify(cons, TRUE) + /\ cleanupPending' = (Variant = "broadcaster") + /\ UNCHANGED <> + +CleanupMicrotask == + /\ cleanupPending + /\ LET cl == Cleanup(buf, head, cons, complete) + IN buf' = cl[1] /\ head' = cl[2] + /\ cleanupPending' = FALSE + /\ UNCHANGED <> + +------------------------------------------------------------------------------ +(* Consumers *) + +Create(c) == + /\ pc[c] = "notStarted" + /\ (cancelPhase = "none" \/ AllowLateCreate) + /\ LET registers == ~(CancelledFlag /\ cancelPhase # "none") + IN cons' = [cons EXCEPT ![c] = [inMap |-> registers, pos |-> trim, waiting |-> FALSE, + resolved |-> FALSE, rejected |-> FALSE]] + /\ joinPos' = [joinPos EXCEPT ![c] = trim] + /\ joinedLate' = [joinedLate EXCEPT ![c] = complete] + /\ nextId' = nextId + 1 + /\ pc' = [pc EXCEPT ![c] = "idle"] + /\ UNCHANGED <> + +\* next(): synchronous section +Next(c) == + /\ pc[c] = "idle" + /\ IF ~Live(c) THEN + \* consumers.get(consumerId) === undefined -> { done: true } + /\ pc' = [pc EXCEPT ![c] = "done"] + /\ UNCHANGED <> + ELSE IF HasUnread(c) THEN + \* read the slot; a CLEARED slot is the thrown invariant error + LET v == IF Idx(c) >= 0 THEN buf[Idx(c) + 1] ELSE CLEARED + cs == [cons EXCEPT ![c].pos = @ + 1] + tr == TrimConsumed(buf, head, trim, cs, nextId) + IN + /\ v # CLEARED + /\ recv' = [recv EXCEPT ![c] = Append(@, v)] + /\ cons' = cs + /\ buf' = tr[1] /\ head' = tr[2] /\ trim' = tr[3] + /\ UNCHANGED pc + ELSE IF complete THEN + LET cs == [cons EXCEPT ![c].inMap = FALSE] + cl == Cleanup(buf, head, cs, TRUE) + IN + /\ cons' = cs + /\ buf' = cl[1] /\ head' = cl[2] + /\ pc' = [pc EXCEPT ![c] = IF err THEN "errored" ELSE "done"] + /\ UNCHANGED <> + ELSE IF err THEN + /\ cons' = [cons EXCEPT ![c].inMap = FALSE] + /\ pc' = [pc EXCEPT ![c] = "errored"] + /\ UNCHANGED <> + ELSE + \* install waitingPromise; the re-check inside the executor cannot + \* fire here because the same conditions were just evaluated + /\ cons' = [cons EXCEPT ![c].waiting = TRUE, ![c].resolved = FALSE, + ![c].rejected = FALSE] + /\ pc' = [pc EXCEPT ![c] = "waiting"] + /\ UNCHANGED <> + /\ UNCHANGED <> + +\* resumption after `await waitPromise` +Wake(c) == + /\ pc[c] = "waiting" + /\ cons[c].resolved \/ cons[c].rejected + /\ IF cons[c].rejected + THEN pc' = [pc EXCEPT ![c] = "errored"] \* consumer stays in the map + ELSE pc' = [pc EXCEPT ![c] = "idle"] \* `return this.next()` + /\ cons' = [cons EXCEPT ![c].resolved = FALSE, ![c].rejected = FALSE] + /\ UNCHANGED <> + +\* iterator.return(): consumer-side cancellation +Return(c) == + /\ pc[c] \in {"idle", "waiting"} + /\ Live(c) + /\ LET cs == [cons EXCEPT ![c].inMap = FALSE, ![c].waiting = FALSE] + tr == TrimConsumed(buf, head, trim, cs, nextId) + cl == Cleanup(tr[1], tr[2], cs, complete) + IN /\ cons' = cs + /\ buf' = cl[1] /\ head' = cl[2] /\ trim' = tr[3] + /\ pc' = [pc EXCEPT ![c] = "cancelled"] + /\ UNCHANGED <> + +------------------------------------------------------------------------------ +(* ReusableReadableStream.cancel(): synchronous part, then an await on the + source reader's cancel(), then a second dropUnreadBacklog(). *) + +CancelSync == + /\ AllowStreamCancel + /\ cancelPhase = "none" + /\ LET cs == [c \in Consumers |-> [cons[c] EXCEPT !.inMap = FALSE, !.waiting = FALSE]] + dr == DropUnreadBacklog(buf, head, trim, nextId) + IN /\ cons' = cs + /\ buf' = dr[1] /\ head' = dr[2] /\ trim' = dr[3] + \* every live consumer's pending/next call now returns { done: true } + /\ pc' = [c \in Consumers |-> IF pc[c] \in {"idle", "waiting"} THEN "cancelled" ELSE pc[c]] + \* without a source reader (pump never started) there is nothing to await + /\ cancelPhase' = IF PumpStarted THEN "awaitingSource" ELSE "done" + /\ UNCHANGED <> + +CancelResume == + /\ cancelPhase = "awaitingSource" + /\ pump = "done" \* reader.cancel() resolves after the pump observed done + /\ LET dr == DropUnreadBacklog(buf, head, trim, nextId) + IN buf' = dr[1] /\ head' = dr[2] /\ trim' = dr[3] + /\ cancelPhase' = "done" + /\ UNCHANGED <> + +------------------------------------------------------------------------------ +Terminating == + /\ (pump = "done" \/ ~PumpStarted) + /\ \A c \in Consumers : pc[c] \in {"done", "errored", "cancelled", "notStarted"} + /\ cancelPhase # "awaitingSource" + /\ ~cleanupPending + /\ UNCHANGED vars + +NextState == + \/ Push \/ Complete \/ Fail \/ CleanupMicrotask + \/ \E c \in Consumers : Create(c) \/ Next(c) \/ Wake(c) \/ Return(c) + \/ CancelSync \/ CancelResume + \/ Terminating + +Fairness == + /\ WF_vars(Complete) + /\ WF_vars(CancelResume) + /\ WF_vars(CleanupMicrotask) + /\ \A c \in Consumers : WF_vars(Next(c)) /\ WF_vars(Wake(c)) + +Spec == Init /\ [][NextState]_vars /\ Fairness + +------------------------------------------------------------------------------ +(* Invariants *) + +TypeOK == + /\ head \in Nat /\ trim \in Nat /\ produced \in Nat + /\ \A i \in 1..Len(buf) : buf[i] \in 0..MaxValues + /\ \A c \in Consumers : pc[c] \in PCs + +\* produced == trim + (buffered values past head). The broadcaster's +\* cleanup() clears the buffer without advancing trimOffset, so only the +\* stream variant keeps exact accounting. +BufferAccounting == + IF Variant = "stream" THEN produced = trim + Len(buf) - head + ELSE produced >= trim + Len(buf) - head + +\* A live consumer about to read never hits a negative index or a cleared +\* slot (the code throws "buffer invariant violated" in that case). +NoInvalidRead == + \A c \in Consumers : + (Live(c) /\ pc[c] \in {"idle", "waiting"} /\ HasUnread(c)) => + (Idx(c) >= 0 /\ buf[Idx(c) + 1] # CLEARED) + +\* A live consumer's position is never behind the trim watermark. +PositionNotBehindTrim == + \A c \in Consumers : Live(c) => cons[c].pos >= trim + +\* Each consumer receives values in order, contiguous from its join point. +InOrderNoGaps == + \A c \in Consumers : \A i \in 1..Len(recv[c]) : recv[c][i] = joinPos[c] + i + +\* A consumer that observed completion has seen every value from its join +\* point onward. Exception (documented trade-off): the broadcaster's +\* cleanup() drops the backlog on completion in active-consumers mode, so a +\* consumer created after completion may see nothing. +CompleteMeansAll == + \A c \in Consumers : + (pc[c] = "done" /\ cancelPhase = "none" + /\ ~(Variant = "broadcaster" /\ Mode = "active" /\ joinedLate[c])) => + joinPos[c] + Len(recv[c]) = produced + +\* No consumer wait is lost: a waiting consumer with something to observe +\* has been (or will be, in the same step) resolved. +NoLostWakeup == + \A c \in Consumers : + (Live(c) /\ pc[c] = "waiting" /\ cons[c].waiting) => + ~(HasUnread(c) \/ complete \/ err) + +Safety == TypeOK /\ BufferAccounting /\ NoInvalidRead /\ PositionNotBehindTrim + /\ InOrderNoGaps /\ CompleteMeansAll /\ NoLostWakeup + +------------------------------------------------------------------------------ +(* Liveness *) + +\* Every consumer that starts eventually stops waiting once the pump is done. +EventuallyDone == + \A c \in Consumers : []((pc[c] # "notStarted") => <>(pc[c] \in {"done", "errored", "cancelled"})) + +============================================================================== diff --git a/specs/tla/MultiConsumerStream_active.cfg b/specs/tla/MultiConsumerStream_active.cfg new file mode 100644 index 00000000..d799510b --- /dev/null +++ b/specs/tla/MultiConsumerStream_active.cfg @@ -0,0 +1,14 @@ +SPECIFICATION Spec +CONSTANTS + Variant = "stream" + Consumers = {c1, c2} + MaxValues = 3 + Mode = "active" + CompactionMinHead = 2 + AllowStreamCancel = TRUE + CancelledFlag = TRUE + AllowLateCreate = FALSE +INVARIANTS + Safety +PROPERTIES + EventuallyDone diff --git a/specs/tla/MultiConsumerStream_active_latecreate.cfg b/specs/tla/MultiConsumerStream_active_latecreate.cfg new file mode 100644 index 00000000..0b5be931 --- /dev/null +++ b/specs/tla/MultiConsumerStream_active_latecreate.cfg @@ -0,0 +1,14 @@ +SPECIFICATION Spec +CONSTANTS + Variant = "stream" + Consumers = {c1, c2} + MaxValues = 3 + Mode = "active" + CompactionMinHead = 2 + AllowStreamCancel = TRUE + CancelledFlag = TRUE + AllowLateCreate = TRUE +INVARIANTS + Safety +PROPERTIES + EventuallyDone diff --git a/specs/tla/MultiConsumerStream_broadcaster_active.cfg b/specs/tla/MultiConsumerStream_broadcaster_active.cfg new file mode 100644 index 00000000..673a85d6 --- /dev/null +++ b/specs/tla/MultiConsumerStream_broadcaster_active.cfg @@ -0,0 +1,14 @@ +SPECIFICATION Spec +CONSTANTS + Variant = "broadcaster" + Consumers = {c1, c2} + MaxValues = 3 + Mode = "active" + CompactionMinHead = 2 + AllowStreamCancel = FALSE + CancelledFlag = TRUE + AllowLateCreate = FALSE +INVARIANTS + Safety +PROPERTIES + EventuallyDone diff --git a/specs/tla/MultiConsumerStream_broadcaster_full.cfg b/specs/tla/MultiConsumerStream_broadcaster_full.cfg new file mode 100644 index 00000000..fc19b7d9 --- /dev/null +++ b/specs/tla/MultiConsumerStream_broadcaster_full.cfg @@ -0,0 +1,14 @@ +SPECIFICATION Spec +CONSTANTS + Variant = "broadcaster" + Consumers = {c1, c2} + MaxValues = 3 + Mode = "full" + CompactionMinHead = 2 + AllowStreamCancel = FALSE + CancelledFlag = TRUE + AllowLateCreate = FALSE +INVARIANTS + Safety +PROPERTIES + EventuallyDone diff --git a/specs/tla/MultiConsumerStream_full.cfg b/specs/tla/MultiConsumerStream_full.cfg new file mode 100644 index 00000000..f7ed3f7e --- /dev/null +++ b/specs/tla/MultiConsumerStream_full.cfg @@ -0,0 +1,14 @@ +SPECIFICATION Spec +CONSTANTS + Variant = "stream" + Consumers = {c1, c2} + MaxValues = 3 + Mode = "full" + CompactionMinHead = 2 + AllowStreamCancel = TRUE + CancelledFlag = TRUE + AllowLateCreate = FALSE +INVARIANTS + Safety +PROPERTIES + EventuallyDone diff --git a/specs/tla/README.md b/specs/tla/README.md new file mode 100644 index 00000000..425e81f0 --- /dev/null +++ b/specs/tla/README.md @@ -0,0 +1,20 @@ +# TLA+ models + +TLA+ specifications of the concurrency surfaces in `@openrouter/agent`, model checked with TLC. Each JavaScript synchronous section is one atomic step and each `await` is a point where other actions may interleave, so the models explore the microtask interleavings the runtime can produce without encoding the event loop itself. + +## Running + +Requires Java 17+ and `tla2tools.jar` from the [TLA+ releases](https://github.com/tlaplus/tlaplus/releases). + +```bash +cd specs/tla +java -cp /path/to/tla2tools.jar tlc2.TLC -workers auto -config MultiConsumerStream_active.cfg MultiConsumerStream.tla +``` + +Every `.cfg` in this directory is expected to pass. The `Fix*` and `CancelledFlag` constants switch a modeled fix on or off, so flipping one to `FALSE` reproduces the counterexample that motivated the change. + +## Models + +`MultiConsumerStream.tla` covers `ReusableReadableStream` and `ToolEventBroadcaster`: consumer creation, reads, waiting promises, `return()`, active-consumer trimming and compaction, source completion and failure, the two-phase `cancel()`, and the broadcaster's completion cleanup microtask. Invariants check buffer accounting, no reads from cleared or trimmed slots, in-order gap-free delivery, no lost wakeups, and eventual consumer termination. With `CancelledFlag = FALSE` and `AllowLateCreate = TRUE`, TLC finds a consumer created during `cancel()`'s await of the source reader that reads a slot the second backlog sweep cleared. `packages/agent/tests/unit/reusable-stream.test.ts` reproduces that trace against the implementation. + +`AsyncRunEndDrain.tla` covers `ModelResult.handleRunEndAsyncTasks` under `onRunEnd: 'drain'` together with `AsyncToolRegistry` settlement. The invariant requires every settled task to be delivered or persisted when the run returns. With `FixDropLeftover = FALSE`, TLC finds a settlement that lands during the last permitted drain turn and is never harvested. `packages/agent/tests/unit/async-tool-background.test.ts` reproduces that trace against the implementation. From 1062ab019447f8a0296f693431eacf1640ac462b Mon Sep 17 00:00:00 2001 From: Abhinav Pola Date: Wed, 23 Sep 2026 01:49:00 +0000 Subject: [PATCH 2/3] fix(agent): cancel an unstarted source in ReusableReadableStream.cancel() cancel() before the first consumer never acquired a reader, and the cancelled flag now stops later consumers from starting the pump, so the source stream would stay open. Cancel it directly in that case. Also extract registerConsumer() so createConsumer() stays under the structural gate's complexity limit. --- packages/agent/src/lib/reusable-stream.ts | 38 +++++++++++++------ .../agent/tests/unit/reusable-stream.test.ts | 21 ++++++++++ 2 files changed, 47 insertions(+), 12 deletions(-) diff --git a/packages/agent/src/lib/reusable-stream.ts b/packages/agent/src/lib/reusable-stream.ts index ea0606b1..1e4908fd 100644 --- a/packages/agent/src/lib/reusable-stream.ts +++ b/packages/agent/src/lib/reusable-stream.ts @@ -79,18 +79,7 @@ export class ReusableReadableStream { */ createConsumer(): AsyncIterableIterator { const consumerId = this.nextConsumerId++; - if (!this.cancelled) { - this.consumers.set(consumerId, { - position: this.trimOffset, - waitingPromise: null, - cancelled: false, - }); - - // Start pumping the source stream if not already started - if (!this.pumpStarted) { - this.startPump(); - } - } + this.registerConsumer(consumerId); // eslint-disable-next-line @typescript-eslint/no-this-alias const self = this; @@ -266,6 +255,22 @@ export class ReusableReadableStream { this.bufferHead = 0; } + private registerConsumer(consumerId: number): void { + if (this.cancelled) { + return; + } + this.consumers.set(consumerId, { + position: this.trimOffset, + waitingPromise: null, + cancelled: false, + }); + + // Start pumping the source stream if not already started + if (!this.pumpStarted) { + this.startPump(); + } + } + /** * Start pumping data from the source stream into the buffer */ @@ -326,6 +331,13 @@ export class ReusableReadableStream { return this.sourceCancelPromise; } + private cancelUnstartedSource(): Promise { + if (!this.sourceCancelPromise) { + this.sourceCancelPromise = this.sourceStream.cancel(); + } + return this.sourceCancelPromise; + } + /** * Notify all waiting consumers that new data is available */ @@ -362,6 +374,8 @@ export class ReusableReadableStream { // Cancel the source stream if (this.sourceReader) { await this.cancelSourceReader(this.sourceReader); + } else if (!this.pumpStarted) { + await this.cancelUnstartedSource(); } /* * The pump may have landed one in-flight chunk between the synchronous diff --git a/packages/agent/tests/unit/reusable-stream.test.ts b/packages/agent/tests/unit/reusable-stream.test.ts index 161e8402..144328c0 100644 --- a/packages/agent/tests/unit/reusable-stream.test.ts +++ b/packages/agent/tests/unit/reusable-stream.test.ts @@ -210,6 +210,27 @@ describe('ReusableReadableStream', () => { // No source.close(): cancel() already terminated the source stream. }); + it('cancel() before the first consumer cancels the source and later consumers are done', async () => { + let sourceCancelled = false; + const source = new ReadableStream({ + cancel(): void { + sourceCancelled = true; + }, + }); + const stream = new ReusableReadableStream(source); + + await stream.cancel(); + + expect(sourceCancelled).toBe(true); + expect(source.locked).toBe(false); + await expect(source.getReader().read()).resolves.toEqual({ + done: true, + value: undefined, + }); + expect((await stream.createConsumer().next()).done).toBe(true); + await stream.cancel(); + }); + it('active-consumers: a consumer created while cancel() awaits the source reader is done', async () => { const source = controlledStream(); const stream = new ReusableReadableStream(source.stream, { From 7c2660b62814d3fc81b45151d72b5e1cab2bf31f Mon Sep 17 00:00:00 2001 From: Abhinav Pola Date: Thu, 24 Sep 2026 17:37:42 +0000 Subject: [PATCH 3/3] chore: drop TLA+ specs from the PR, keep the fixes and unit tests --- .changeset/quiet-streams-drain.md | 2 +- specs/tla/AsyncRunEndDrain.cfg | 8 - specs/tla/AsyncRunEndDrain.tla | 156 -------- specs/tla/AsyncRunEndDrain_timeouts.cfg | 8 - specs/tla/MultiConsumerStream.tla | 371 ------------------ specs/tla/MultiConsumerStream_active.cfg | 14 - .../MultiConsumerStream_active_latecreate.cfg | 14 - ...MultiConsumerStream_broadcaster_active.cfg | 14 - .../MultiConsumerStream_broadcaster_full.cfg | 14 - specs/tla/MultiConsumerStream_full.cfg | 14 - specs/tla/README.md | 20 - 11 files changed, 1 insertion(+), 634 deletions(-) delete mode 100644 specs/tla/AsyncRunEndDrain.cfg delete mode 100644 specs/tla/AsyncRunEndDrain.tla delete mode 100644 specs/tla/AsyncRunEndDrain_timeouts.cfg delete mode 100644 specs/tla/MultiConsumerStream.tla delete mode 100644 specs/tla/MultiConsumerStream_active.cfg delete mode 100644 specs/tla/MultiConsumerStream_active_latecreate.cfg delete mode 100644 specs/tla/MultiConsumerStream_broadcaster_active.cfg delete mode 100644 specs/tla/MultiConsumerStream_broadcaster_full.cfg delete mode 100644 specs/tla/MultiConsumerStream_full.cfg delete mode 100644 specs/tla/README.md diff --git a/.changeset/quiet-streams-drain.md b/.changeset/quiet-streams-drain.md index 764bd088..93f292e3 100644 --- a/.changeset/quiet-streams-drain.md +++ b/.changeset/quiet-streams-drain.md @@ -2,4 +2,4 @@ '@openrouter/agent': patch --- -Fix two run-teardown races found by TLA+ model checking. `ReusableReadableStream.createConsumer()` called while `cancel()` is awaiting the source reader now yields an already-done consumer instead of one that reads a cleared buffer slot and throws. `onRunEnd: 'drain'` now persists and broadcasts (`delivery: 'dropped'`) a background task settlement that lands during the final drain turn, where it was previously left unharvested with the persisted task still `working`. +Fix two run-teardown races. `ReusableReadableStream.createConsumer()` called while `cancel()` is awaiting the source reader now yields an already-done consumer instead of one that reads a cleared buffer slot and throws. `onRunEnd: 'drain'` now persists and broadcasts (`delivery: 'dropped'`) a background task settlement that lands during the final drain turn, where it was previously left unharvested with the persisted task still `working`. diff --git a/specs/tla/AsyncRunEndDrain.cfg b/specs/tla/AsyncRunEndDrain.cfg deleted file mode 100644 index f782979d..00000000 --- a/specs/tla/AsyncRunEndDrain.cfg +++ /dev/null @@ -1,8 +0,0 @@ -SPECIFICATION Spec -CONSTANTS - Tasks = {t1, t2, t3} - MaxDrainTurns = 2 - DrainTimeouts = FALSE - FixDropLeftover = TRUE -INVARIANT Safety -PROPERTY EventuallyDone diff --git a/specs/tla/AsyncRunEndDrain.tla b/specs/tla/AsyncRunEndDrain.tla deleted file mode 100644 index 8f2c0d20..00000000 --- a/specs/tla/AsyncRunEndDrain.tla +++ /dev/null @@ -1,156 +0,0 @@ ----------------------------- MODULE AsyncRunEndDrain ---------------------------- -(* - * Model of ModelResult.handleRunEndAsyncTasks (onRunEnd: 'drain') together - * with AsyncToolRegistry settlement, for background tool tasks that outlived - * their grace window and hold a pending placeholder in the transcript. - * - * Every JavaScript synchronous section is one atomic step. Background work - * settles at any await boundary (Settle action), which is exactly what the - * registry's `.then()` callbacks do at runtime. - * - * Property checked: when the run ends, every task that settled has been - * accounted for, either delivered in a drain turn or dropped and persisted by - * dropSettledTasks. The registry's settled queue is empty at run end. - *) -EXTENDS Naturals, FiniteSets - -CONSTANTS - Tasks, \* set of background tasks with a pending placeholder - MaxDrainTurns, \* asyncTools.maxDrainTurns - DrainTimeouts, \* TRUE: drainTimeoutMs may expire (registry.drain returns false) - FixDropLeftover \* TRUE: run end drops unharvested settlements even with nothing in flight - -VARIABLES - status, \* task -> "working" | "settled" - queued, \* task -> BOOLEAN, settlement sits in registry.settledQueue - outcome, \* task -> "none" | "delivered" | "dropped" - pc, \* program counter of handleRunEndAsyncTasks - drainTurn, \* loop counter - harvested, \* tasks taken by the current flushAsyncToolDeliveries - deadlinePassed - -vars == <> - -InFlight == \E t \in Tasks : status[t] = "working" -Unharvested == \E t \in Tasks : queued[t] - -TypeOK == - /\ status \in [Tasks -> {"working", "settled"}] - /\ queued \in [Tasks -> BOOLEAN] - /\ outcome \in [Tasks -> {"none", "delivered", "dropped"}] - /\ pc \in {"start", "loopTop", "awaitDrain", "harvest", "modelTurn", "afterTurn", - "finalize", "done"} - /\ drainTurn \in 0..MaxDrainTurns - /\ harvested \subseteq Tasks - /\ deadlinePassed \in BOOLEAN - -Init == - /\ status = [t \in Tasks |-> "working"] - /\ queued = [t \in Tasks |-> FALSE] - /\ outcome = [t \in Tasks |-> "none"] - /\ pc = "start" - /\ drainTurn = 0 - /\ harvested = {} - /\ deadlinePassed = FALSE - -\* Background work resolves or rejects: registry.settle pushes onto settledQueue. -Settle(t) == - /\ status[t] = "working" - /\ pc # "done" - /\ status' = [status EXCEPT ![t] = "settled"] - /\ queued' = [queued EXCEPT ![t] = TRUE] - /\ UNCHANGED <> - -\* Wall clock: drainTimeoutMs elapses at some point. -Deadline == - /\ DrainTimeouts - /\ ~deadlinePassed - /\ deadlinePassed' = TRUE - /\ UNCHANGED <> - -\* Early return: nothing in flight and nothing unharvested. -Start == - /\ pc = "start" - /\ pc' = IF ~InFlight /\ ~Unharvested THEN "done" ELSE "loopTop" - /\ UNCHANGED <> - -LoopTop == - /\ pc = "loopTop" - /\ pc' = IF drainTurn >= MaxDrainTurns THEN "finalize" - ELSE IF ~Unharvested /\ InFlight - THEN IF deadlinePassed THEN "finalize" ELSE "awaitDrain" - ELSE "harvest" - /\ UNCHANGED <> - -\* registry.drain(remaining) resumes on a settle event or on timeout. -AwaitDrain == - /\ pc = "awaitDrain" - /\ (Unharvested \/ deadlinePassed \/ ~InFlight) - /\ pc' = "harvest" - /\ UNCHANGED <> - -\* flushAsyncToolDeliveries: takeSettled(), inject envelope. Nothing settled: stop. -Harvest == - /\ pc = "harvest" - /\ LET taken == {t \in Tasks : queued[t]} - IN /\ harvested' = taken - /\ queued' = [t \in Tasks |-> FALSE] - /\ outcome' = [t \in Tasks |-> IF t \in taken THEN "delivered" ELSE outcome[t]] - /\ pc' = IF taken = {} THEN "finalize" ELSE "modelTurn" - /\ UNCHANGED <> - -\* retryCurrentRequest + saveResponseToState (awaits, tasks may settle meanwhile). -ModelTurn == - /\ pc = "modelTurn" - /\ pc' = "afterTurn" - /\ UNCHANGED <> - -AfterTurn == - /\ pc = "afterTurn" - /\ drainTurn' = drainTurn + 1 - /\ pc' = IF ~InFlight /\ ~Unharvested THEN "finalize" ELSE "loopTop" - /\ UNCHANGED <> - -\* Drain budget exhausted: abortAll settles working tasks as cancelled, then -\* dropSettledTasks persists every queued settlement. In the unfixed code the -\* drop only runs when something was still in flight. -Finalize == - /\ pc = "finalize" - /\ LET drops == InFlight \/ FixDropLeftover - settledNow == [t \in Tasks |-> "settled"] - queuedNow == [t \in Tasks |-> queued[t] \/ status[t] = "working"] - IN IF InFlight - THEN /\ status' = settledNow - /\ queued' = [t \in Tasks |-> FALSE] - /\ outcome' = [t \in Tasks |-> IF queuedNow[t] THEN "dropped" ELSE outcome[t]] - ELSE IF drops - THEN /\ status' = status - /\ queued' = [t \in Tasks |-> FALSE] - /\ outcome' = [t \in Tasks |-> IF queued[t] THEN "dropped" ELSE outcome[t]] - ELSE UNCHANGED <> - /\ pc' = "done" - /\ UNCHANGED <> - -Terminating == pc = "done" /\ UNCHANGED vars - -Next == - \/ \E t \in Tasks : Settle(t) - \/ Deadline - \/ Start \/ LoopTop \/ AwaitDrain \/ Harvest \/ ModelTurn \/ AfterTurn \/ Finalize - \/ Terminating - -Spec == Init /\ [][Next]_vars /\ WF_vars(Next) - -\* Every settlement is delivered or persisted by the time the run returns. -NoLeakedSettlement == - pc = "done" => \A t \in Tasks : ~queued[t] /\ (status[t] = "settled" => outcome[t] # "none") - -\* Delivery happens at most once and only for settled tasks. -DeliveryConsistent == - \A t \in Tasks : outcome[t] # "none" => status[t] = "settled" - -Safety == TypeOK /\ NoLeakedSettlement /\ DeliveryConsistent - -EventuallyDone == <>(pc = "done") - -============================================================================= diff --git a/specs/tla/AsyncRunEndDrain_timeouts.cfg b/specs/tla/AsyncRunEndDrain_timeouts.cfg deleted file mode 100644 index 62ab9e76..00000000 --- a/specs/tla/AsyncRunEndDrain_timeouts.cfg +++ /dev/null @@ -1,8 +0,0 @@ -SPECIFICATION Spec -CONSTANTS - Tasks = {t1, t2, t3} - MaxDrainTurns = 2 - DrainTimeouts = TRUE - FixDropLeftover = TRUE -INVARIANT Safety -PROPERTY EventuallyDone diff --git a/specs/tla/MultiConsumerStream.tla b/specs/tla/MultiConsumerStream.tla deleted file mode 100644 index a8a83ff9..00000000 --- a/specs/tla/MultiConsumerStream.tla +++ /dev/null @@ -1,371 +0,0 @@ --------------------------- MODULE MultiConsumerStream -------------------------- -(* - * Model of the multi-consumer buffer shared by - * packages/agent/src/lib/reusable-stream.ts (ReusableReadableStream) - * packages/agent/src/lib/tool-event-broadcaster.ts (ToolEventBroadcaster) - * - * One producer ("pump") appends values to a buffer. Consumers hold an - * absolute position; the buffer index of a consumer is - * head + position - trim - * Consumers that find nothing to read install a waiting promise and resume - * later (a separate atomic step, since resumption is a microtask). In - * 'active-consumers' replay mode the buffer is trimmed to the slowest live - * consumer, and the retained backlog is dropped when the last consumer - * leaves. ReusableReadableStream.cancel() is a two-phase operation with an - * await in the middle; both phases are modelled. - * - * Atomicity follows JavaScript: every synchronous section between two - * awaits is one action. - *) -EXTENDS Naturals, Sequences, FiniteSets, TLC - -CONSTANTS - Variant, \* "stream" (ReusableReadableStream) | "broadcaster" (ToolEventBroadcaster) - Consumers, \* set of consumer ids - MaxValues, \* bound on values the pump produces - Mode, \* "full" | "active" - CompactionMinHead, \* BUFFER_COMPACTION_MIN_HEAD (1024 in code; small here) - AllowStreamCancel, \* model ReusableReadableStream.cancel() - AllowLateCreate, \* allow createConsumer() after cancel() started - CancelledFlag \* fix: createConsumer() after cancel() registers nothing - -ASSUME Mode \in {"full", "active"} -ASSUME Variant \in {"stream", "broadcaster"} -ASSUME Variant = "broadcaster" => ~AllowStreamCancel - -CLEARED == 0 \* `undefined` slot left by buffer.fill(undefined, ...) - -VARIABLES - buf, \* Seq of values (1..MaxValues) or CLEARED - head, \* bufferHead - trim, \* trimOffset - produced, \* number of values pushed so far - complete, \* sourceComplete / isComplete - err, \* sourceError / completionError set - nextId, \* nextConsumerId (only used for the `=== 0` check) - pump, \* "running" | "done" - cons, \* consumer records - pc, \* consumer program counters - recv, \* values received by each consumer - joinPos, \* position at createConsumer() - joinedLate, \* createConsumer() happened after completion - cancelPhase, \* "none" | "awaitingSource" | "done" - cleanupPending \* broadcaster: queueMicrotask(() => this.cleanup()) outstanding - -vars == <> - -ConsumerRecord == [inMap: BOOLEAN, pos: Nat, waiting: BOOLEAN, - resolved: BOOLEAN, rejected: BOOLEAN] - -PCs == {"notStarted", "idle", "waiting", "done", "errored", "cancelled"} - -Live(c) == cons[c].inMap - -LiveConsumers == {c \in Consumers : Live(c)} - -Idx(c) == head + cons[c].pos - trim \* JS bufferIndex (0-based) - -HasUnread(c) == Idx(c) < Len(buf) \* JS: bufferIndex < buffer.length - ------------------------------------------------------------------------------- -(* Helpers mirroring the private methods *) - -\* dropUnreadBacklog(): returns [buf, head, trim] -DropUnreadBacklog(b, h, t, nid) == - IF nid = 0 \/ Len(b) - h = 0 - THEN <> - ELSE <<<<>>, 0, t + (Len(b) - h)>> - -Min(S) == CHOOSE m \in S : \A x \in S : m <= x - -\* trimConsumed() given the consumer map `cs`; returns [buf, head, trim] -TrimConsumed(b, h, t, cs, nid) == - IF Mode = "full" THEN <> - ELSE - LET live == {c \in Consumers : cs[c].inMap} - IN IF live = {} THEN DropUnreadBacklog(b, h, t, nid) - ELSE - LET m == Min({cs[c].pos : c \in live}) - nextHead == h + m - t - IN IF nextHead <= h THEN <> - ELSE IF nextHead = Len(b) THEN <<<<>>, 0, m>> - ELSE - LET filled == [i \in 1..Len(b) |-> - IF i > h /\ i <= nextHead THEN CLEARED ELSE b[i]] - IN IF nextHead >= CompactionMinHead /\ nextHead * 2 >= Len(b) - THEN <> - ELSE <> - -\* ToolEventBroadcaster.cleanup(): drop the buffer once complete with no -\* consumers left (active-consumers mode only). Note: trimOffset is NOT -\* advanced here, unlike dropUnreadBacklog(). -Cleanup(b, h, cs, isComplete) == - IF Variant = "broadcaster" /\ Mode = "active" /\ isComplete - /\ {c \in Consumers : cs[c].inMap} = {} - THEN <<<<>>, 0>> - ELSE <> - -\* notifyAllConsumers(): resolve or reject every waiting promise -Notify(cs, isErr) == - [c \in Consumers |-> - IF cs[c].inMap /\ cs[c].waiting - THEN [cs[c] EXCEPT !.waiting = FALSE, - !.resolved = ~isErr, - !.rejected = isErr] - ELSE cs[c]] - ------------------------------------------------------------------------------- -Init == - /\ buf = <<>> - /\ head = 0 - /\ trim = 0 - /\ produced = 0 - /\ complete = FALSE - /\ err = FALSE - /\ nextId = 0 - /\ pump = "running" - /\ cons = [c \in Consumers |-> [inMap |-> FALSE, pos |-> 0, waiting |-> FALSE, - resolved |-> FALSE, rejected |-> FALSE]] - /\ pc = [c \in Consumers |-> "notStarted"] - /\ recv = [c \in Consumers |-> <<>>] - /\ joinPos = [c \in Consumers |-> 0] - /\ joinedLate = [c \in Consumers |-> FALSE] - /\ cancelPhase = "none" - /\ cleanupPending = FALSE - ------------------------------------------------------------------------------- -(* Pump *) - -\* the stream's pump only starts with the first createConsumer(); the -\* broadcaster is pushed into regardless of consumers. -PumpStarted == Variant = "broadcaster" \/ nextId > 0 - -Push == - /\ pump = "running" - /\ PumpStarted - /\ ~complete /\ ~err - /\ produced < MaxValues - /\ buf' = Append(buf, produced + 1) - /\ produced' = produced + 1 - /\ cons' = Notify(cons, FALSE) - /\ UNCHANGED <> - -Complete == - /\ pump = "running" - /\ PumpStarted - /\ complete' = TRUE - /\ pump' = "done" - /\ cons' = Notify(cons, FALSE) - /\ cleanupPending' = (Variant = "broadcaster") - /\ UNCHANGED <> - -\* stream: pump error. broadcaster: complete(error). -Fail == - /\ pump = "running" - /\ PumpStarted - /\ cancelPhase = "none" - /\ err' = TRUE - /\ complete' = (Variant = "broadcaster") - /\ pump' = "done" - /\ cons' = Notify(cons, TRUE) - /\ cleanupPending' = (Variant = "broadcaster") - /\ UNCHANGED <> - -CleanupMicrotask == - /\ cleanupPending - /\ LET cl == Cleanup(buf, head, cons, complete) - IN buf' = cl[1] /\ head' = cl[2] - /\ cleanupPending' = FALSE - /\ UNCHANGED <> - ------------------------------------------------------------------------------- -(* Consumers *) - -Create(c) == - /\ pc[c] = "notStarted" - /\ (cancelPhase = "none" \/ AllowLateCreate) - /\ LET registers == ~(CancelledFlag /\ cancelPhase # "none") - IN cons' = [cons EXCEPT ![c] = [inMap |-> registers, pos |-> trim, waiting |-> FALSE, - resolved |-> FALSE, rejected |-> FALSE]] - /\ joinPos' = [joinPos EXCEPT ![c] = trim] - /\ joinedLate' = [joinedLate EXCEPT ![c] = complete] - /\ nextId' = nextId + 1 - /\ pc' = [pc EXCEPT ![c] = "idle"] - /\ UNCHANGED <> - -\* next(): synchronous section -Next(c) == - /\ pc[c] = "idle" - /\ IF ~Live(c) THEN - \* consumers.get(consumerId) === undefined -> { done: true } - /\ pc' = [pc EXCEPT ![c] = "done"] - /\ UNCHANGED <> - ELSE IF HasUnread(c) THEN - \* read the slot; a CLEARED slot is the thrown invariant error - LET v == IF Idx(c) >= 0 THEN buf[Idx(c) + 1] ELSE CLEARED - cs == [cons EXCEPT ![c].pos = @ + 1] - tr == TrimConsumed(buf, head, trim, cs, nextId) - IN - /\ v # CLEARED - /\ recv' = [recv EXCEPT ![c] = Append(@, v)] - /\ cons' = cs - /\ buf' = tr[1] /\ head' = tr[2] /\ trim' = tr[3] - /\ UNCHANGED pc - ELSE IF complete THEN - LET cs == [cons EXCEPT ![c].inMap = FALSE] - cl == Cleanup(buf, head, cs, TRUE) - IN - /\ cons' = cs - /\ buf' = cl[1] /\ head' = cl[2] - /\ pc' = [pc EXCEPT ![c] = IF err THEN "errored" ELSE "done"] - /\ UNCHANGED <> - ELSE IF err THEN - /\ cons' = [cons EXCEPT ![c].inMap = FALSE] - /\ pc' = [pc EXCEPT ![c] = "errored"] - /\ UNCHANGED <> - ELSE - \* install waitingPromise; the re-check inside the executor cannot - \* fire here because the same conditions were just evaluated - /\ cons' = [cons EXCEPT ![c].waiting = TRUE, ![c].resolved = FALSE, - ![c].rejected = FALSE] - /\ pc' = [pc EXCEPT ![c] = "waiting"] - /\ UNCHANGED <> - /\ UNCHANGED <> - -\* resumption after `await waitPromise` -Wake(c) == - /\ pc[c] = "waiting" - /\ cons[c].resolved \/ cons[c].rejected - /\ IF cons[c].rejected - THEN pc' = [pc EXCEPT ![c] = "errored"] \* consumer stays in the map - ELSE pc' = [pc EXCEPT ![c] = "idle"] \* `return this.next()` - /\ cons' = [cons EXCEPT ![c].resolved = FALSE, ![c].rejected = FALSE] - /\ UNCHANGED <> - -\* iterator.return(): consumer-side cancellation -Return(c) == - /\ pc[c] \in {"idle", "waiting"} - /\ Live(c) - /\ LET cs == [cons EXCEPT ![c].inMap = FALSE, ![c].waiting = FALSE] - tr == TrimConsumed(buf, head, trim, cs, nextId) - cl == Cleanup(tr[1], tr[2], cs, complete) - IN /\ cons' = cs - /\ buf' = cl[1] /\ head' = cl[2] /\ trim' = tr[3] - /\ pc' = [pc EXCEPT ![c] = "cancelled"] - /\ UNCHANGED <> - ------------------------------------------------------------------------------- -(* ReusableReadableStream.cancel(): synchronous part, then an await on the - source reader's cancel(), then a second dropUnreadBacklog(). *) - -CancelSync == - /\ AllowStreamCancel - /\ cancelPhase = "none" - /\ LET cs == [c \in Consumers |-> [cons[c] EXCEPT !.inMap = FALSE, !.waiting = FALSE]] - dr == DropUnreadBacklog(buf, head, trim, nextId) - IN /\ cons' = cs - /\ buf' = dr[1] /\ head' = dr[2] /\ trim' = dr[3] - \* every live consumer's pending/next call now returns { done: true } - /\ pc' = [c \in Consumers |-> IF pc[c] \in {"idle", "waiting"} THEN "cancelled" ELSE pc[c]] - \* without a source reader (pump never started) there is nothing to await - /\ cancelPhase' = IF PumpStarted THEN "awaitingSource" ELSE "done" - /\ UNCHANGED <> - -CancelResume == - /\ cancelPhase = "awaitingSource" - /\ pump = "done" \* reader.cancel() resolves after the pump observed done - /\ LET dr == DropUnreadBacklog(buf, head, trim, nextId) - IN buf' = dr[1] /\ head' = dr[2] /\ trim' = dr[3] - /\ cancelPhase' = "done" - /\ UNCHANGED <> - ------------------------------------------------------------------------------- -Terminating == - /\ (pump = "done" \/ ~PumpStarted) - /\ \A c \in Consumers : pc[c] \in {"done", "errored", "cancelled", "notStarted"} - /\ cancelPhase # "awaitingSource" - /\ ~cleanupPending - /\ UNCHANGED vars - -NextState == - \/ Push \/ Complete \/ Fail \/ CleanupMicrotask - \/ \E c \in Consumers : Create(c) \/ Next(c) \/ Wake(c) \/ Return(c) - \/ CancelSync \/ CancelResume - \/ Terminating - -Fairness == - /\ WF_vars(Complete) - /\ WF_vars(CancelResume) - /\ WF_vars(CleanupMicrotask) - /\ \A c \in Consumers : WF_vars(Next(c)) /\ WF_vars(Wake(c)) - -Spec == Init /\ [][NextState]_vars /\ Fairness - ------------------------------------------------------------------------------- -(* Invariants *) - -TypeOK == - /\ head \in Nat /\ trim \in Nat /\ produced \in Nat - /\ \A i \in 1..Len(buf) : buf[i] \in 0..MaxValues - /\ \A c \in Consumers : pc[c] \in PCs - -\* produced == trim + (buffered values past head). The broadcaster's -\* cleanup() clears the buffer without advancing trimOffset, so only the -\* stream variant keeps exact accounting. -BufferAccounting == - IF Variant = "stream" THEN produced = trim + Len(buf) - head - ELSE produced >= trim + Len(buf) - head - -\* A live consumer about to read never hits a negative index or a cleared -\* slot (the code throws "buffer invariant violated" in that case). -NoInvalidRead == - \A c \in Consumers : - (Live(c) /\ pc[c] \in {"idle", "waiting"} /\ HasUnread(c)) => - (Idx(c) >= 0 /\ buf[Idx(c) + 1] # CLEARED) - -\* A live consumer's position is never behind the trim watermark. -PositionNotBehindTrim == - \A c \in Consumers : Live(c) => cons[c].pos >= trim - -\* Each consumer receives values in order, contiguous from its join point. -InOrderNoGaps == - \A c \in Consumers : \A i \in 1..Len(recv[c]) : recv[c][i] = joinPos[c] + i - -\* A consumer that observed completion has seen every value from its join -\* point onward. Exception (documented trade-off): the broadcaster's -\* cleanup() drops the backlog on completion in active-consumers mode, so a -\* consumer created after completion may see nothing. -CompleteMeansAll == - \A c \in Consumers : - (pc[c] = "done" /\ cancelPhase = "none" - /\ ~(Variant = "broadcaster" /\ Mode = "active" /\ joinedLate[c])) => - joinPos[c] + Len(recv[c]) = produced - -\* No consumer wait is lost: a waiting consumer with something to observe -\* has been (or will be, in the same step) resolved. -NoLostWakeup == - \A c \in Consumers : - (Live(c) /\ pc[c] = "waiting" /\ cons[c].waiting) => - ~(HasUnread(c) \/ complete \/ err) - -Safety == TypeOK /\ BufferAccounting /\ NoInvalidRead /\ PositionNotBehindTrim - /\ InOrderNoGaps /\ CompleteMeansAll /\ NoLostWakeup - ------------------------------------------------------------------------------- -(* Liveness *) - -\* Every consumer that starts eventually stops waiting once the pump is done. -EventuallyDone == - \A c \in Consumers : []((pc[c] # "notStarted") => <>(pc[c] \in {"done", "errored", "cancelled"})) - -============================================================================== diff --git a/specs/tla/MultiConsumerStream_active.cfg b/specs/tla/MultiConsumerStream_active.cfg deleted file mode 100644 index d799510b..00000000 --- a/specs/tla/MultiConsumerStream_active.cfg +++ /dev/null @@ -1,14 +0,0 @@ -SPECIFICATION Spec -CONSTANTS - Variant = "stream" - Consumers = {c1, c2} - MaxValues = 3 - Mode = "active" - CompactionMinHead = 2 - AllowStreamCancel = TRUE - CancelledFlag = TRUE - AllowLateCreate = FALSE -INVARIANTS - Safety -PROPERTIES - EventuallyDone diff --git a/specs/tla/MultiConsumerStream_active_latecreate.cfg b/specs/tla/MultiConsumerStream_active_latecreate.cfg deleted file mode 100644 index 0b5be931..00000000 --- a/specs/tla/MultiConsumerStream_active_latecreate.cfg +++ /dev/null @@ -1,14 +0,0 @@ -SPECIFICATION Spec -CONSTANTS - Variant = "stream" - Consumers = {c1, c2} - MaxValues = 3 - Mode = "active" - CompactionMinHead = 2 - AllowStreamCancel = TRUE - CancelledFlag = TRUE - AllowLateCreate = TRUE -INVARIANTS - Safety -PROPERTIES - EventuallyDone diff --git a/specs/tla/MultiConsumerStream_broadcaster_active.cfg b/specs/tla/MultiConsumerStream_broadcaster_active.cfg deleted file mode 100644 index 673a85d6..00000000 --- a/specs/tla/MultiConsumerStream_broadcaster_active.cfg +++ /dev/null @@ -1,14 +0,0 @@ -SPECIFICATION Spec -CONSTANTS - Variant = "broadcaster" - Consumers = {c1, c2} - MaxValues = 3 - Mode = "active" - CompactionMinHead = 2 - AllowStreamCancel = FALSE - CancelledFlag = TRUE - AllowLateCreate = FALSE -INVARIANTS - Safety -PROPERTIES - EventuallyDone diff --git a/specs/tla/MultiConsumerStream_broadcaster_full.cfg b/specs/tla/MultiConsumerStream_broadcaster_full.cfg deleted file mode 100644 index fc19b7d9..00000000 --- a/specs/tla/MultiConsumerStream_broadcaster_full.cfg +++ /dev/null @@ -1,14 +0,0 @@ -SPECIFICATION Spec -CONSTANTS - Variant = "broadcaster" - Consumers = {c1, c2} - MaxValues = 3 - Mode = "full" - CompactionMinHead = 2 - AllowStreamCancel = FALSE - CancelledFlag = TRUE - AllowLateCreate = FALSE -INVARIANTS - Safety -PROPERTIES - EventuallyDone diff --git a/specs/tla/MultiConsumerStream_full.cfg b/specs/tla/MultiConsumerStream_full.cfg deleted file mode 100644 index f7ed3f7e..00000000 --- a/specs/tla/MultiConsumerStream_full.cfg +++ /dev/null @@ -1,14 +0,0 @@ -SPECIFICATION Spec -CONSTANTS - Variant = "stream" - Consumers = {c1, c2} - MaxValues = 3 - Mode = "full" - CompactionMinHead = 2 - AllowStreamCancel = TRUE - CancelledFlag = TRUE - AllowLateCreate = FALSE -INVARIANTS - Safety -PROPERTIES - EventuallyDone diff --git a/specs/tla/README.md b/specs/tla/README.md deleted file mode 100644 index 425e81f0..00000000 --- a/specs/tla/README.md +++ /dev/null @@ -1,20 +0,0 @@ -# TLA+ models - -TLA+ specifications of the concurrency surfaces in `@openrouter/agent`, model checked with TLC. Each JavaScript synchronous section is one atomic step and each `await` is a point where other actions may interleave, so the models explore the microtask interleavings the runtime can produce without encoding the event loop itself. - -## Running - -Requires Java 17+ and `tla2tools.jar` from the [TLA+ releases](https://github.com/tlaplus/tlaplus/releases). - -```bash -cd specs/tla -java -cp /path/to/tla2tools.jar tlc2.TLC -workers auto -config MultiConsumerStream_active.cfg MultiConsumerStream.tla -``` - -Every `.cfg` in this directory is expected to pass. The `Fix*` and `CancelledFlag` constants switch a modeled fix on or off, so flipping one to `FALSE` reproduces the counterexample that motivated the change. - -## Models - -`MultiConsumerStream.tla` covers `ReusableReadableStream` and `ToolEventBroadcaster`: consumer creation, reads, waiting promises, `return()`, active-consumer trimming and compaction, source completion and failure, the two-phase `cancel()`, and the broadcaster's completion cleanup microtask. Invariants check buffer accounting, no reads from cleared or trimmed slots, in-order gap-free delivery, no lost wakeups, and eventual consumer termination. With `CancelledFlag = FALSE` and `AllowLateCreate = TRUE`, TLC finds a consumer created during `cancel()`'s await of the source reader that reads a slot the second backlog sweep cleared. `packages/agent/tests/unit/reusable-stream.test.ts` reproduces that trace against the implementation. - -`AsyncRunEndDrain.tla` covers `ModelResult.handleRunEndAsyncTasks` under `onRunEnd: 'drain'` together with `AsyncToolRegistry` settlement. The invariant requires every settled task to be delivered or persisted when the run returns. With `FixDropLeftover = FALSE`, TLC finds a settlement that lands during the last permitted drain turn and is never harvested. `packages/agent/tests/unit/async-tool-background.test.ts` reproduces that trace against the implementation.