From d1edecbf608cab5fbaec74c1699fcdc518ef4d50 Mon Sep 17 00:00:00 2001 From: goodboy Date: Thu, 13 Aug 2026 15:49:05 -0400 Subject: [PATCH] Clean broadcast cancellation diagnostics Bound `BroadcastState.cancelled` entries to receiver progress, terminal state and resource lifetime instead of retaining completed `Task`s indefinitely. Deats, - make EOC durable so awakened peers never re-enter a closed source. - close root broadcasters during explicit `MsgStream` and `LinkedTaskChannel` teardown without re-entrant EOC closure or breaking `MsgStream.aclose()` overrides. - reject non-positive fan-out retention capacity before constructing an unusable zero-length queue. - cover child/root cancellation cleanup, terminal peer wakeups, wrapper teardown, subclass compatibility and zero-buffer rejection. Prompt-IO: ai/prompt-io/opencode/20260813T181901Z_a2e0df4b_prompt_io.md (this patch was generated in some part by `opencode` using `gpt-5.6-sol` (`openai`)) --- .../20260813T181901Z_a2e0df4b_prompt_io.md | 36 ++ ...20260813T181901Z_a2e0df4b_prompt_io.raw.md | 58 +++ tests/test_task_broadcasting.py | 365 ++++++++++++++++++ tractor/_streaming.py | 31 +- tractor/to_asyncio.py | 8 + tractor/trionics/_broadcast.py | 18 +- 6 files changed, 514 insertions(+), 2 deletions(-) create mode 100644 ai/prompt-io/opencode/20260813T181901Z_a2e0df4b_prompt_io.md create mode 100644 ai/prompt-io/opencode/20260813T181901Z_a2e0df4b_prompt_io.raw.md diff --git a/ai/prompt-io/opencode/20260813T181901Z_a2e0df4b_prompt_io.md b/ai/prompt-io/opencode/20260813T181901Z_a2e0df4b_prompt_io.md new file mode 100644 index 00000000..8b1fa6ac --- /dev/null +++ b/ai/prompt-io/opencode/20260813T181901Z_a2e0df4b_prompt_io.md @@ -0,0 +1,36 @@ +--- +model: openai/gpt-5.6-sol +service: opencode +session: unavailable +timestamp: 2026-08-13T18:19:01Z +git_ref: a2e0df4b +scope: code +substantive: true +raw_file: 20260813T181901Z_a2e0df4b_prompt_io.raw.md +--- + +## Prompt + +Continue Tractor PR 490 after the paired piker EMS consumer commit. Clean +cancelled-task diagnostics, define or reject zero-buffer broadcast behavior, +review and verify the change, then stop at a complete commit plan. + +## Response summary + +Bound cancellation diagnostics to receiver progress, terminal state and +resource lifetime; made EOC durable across peers; released wrapper-owned root +broadcasters without breaking graceful EOC or subclass overrides; and rejected +non-positive fan-out retention capacity. + +## Files changed + +- `tractor/trionics/_broadcast.py` - diagnostic lifecycle, durable EOC and + buffer validation. +- `tractor/_streaming.py` - safe `MsgStream` root broadcaster cleanup. +- `tractor/to_asyncio.py` - linked-channel root broadcaster cleanup. +- `tests/test_task_broadcasting.py` - cancellation, EOC, wrapper and capacity + regressions. + +## Human edits + +None - generated output follows the requested fifth iterative item. diff --git a/ai/prompt-io/opencode/20260813T181901Z_a2e0df4b_prompt_io.raw.md b/ai/prompt-io/opencode/20260813T181901Z_a2e0df4b_prompt_io.raw.md new file mode 100644 index 00000000..a3556002 --- /dev/null +++ b/ai/prompt-io/opencode/20260813T181901Z_a2e0df4b_prompt_io.raw.md @@ -0,0 +1,58 @@ +--- +model: openai/gpt-5.6-sol +service: opencode +timestamp: 2026-08-13T18:19:01Z +git_ref: a2e0df4b +diff_cmd: git diff HEAD~1..HEAD +--- + +The user asked to continue after committing the paired piker EMS consumer +fix. The next isolated Tractor PR 490 item was to clean cancelled-task +diagnostics and define zero-buffer broadcast behavior, then review, test and +stop at a complete commit plan. + +> `git diff HEAD~1..HEAD -- tractor/trionics/_broadcast.py` + +Made `BroadcastState.cancelled` transient: receiver progress and close clear +that receiver's diagnostic, terminal EOC and shared receive failure clear all +stale cancelled tasks, and durable EOC prevents peers from re-entering the +closed source. `broadcast_receiver()` now rejects non-positive retention +capacity before creating an unusable zero-length deque. + +> `git diff HEAD~1..HEAD -- tractor/_streaming.py` + +Made explicit `MsgStream.aclose()` release its internally allocated root +broadcaster while preserving graceful receive-internal EOC teardown. Used a +task-local marker so the public zero-argument `aclose()` signature and valid +subclass overrides remain compatible. + +> `git diff HEAD~1..HEAD -- tractor/to_asyncio.py` + +Made `LinkedTaskChannel.aclose()` release its internally allocated root +broadcaster before closing the underlying Trio receive channel. + +> `git diff HEAD~1..HEAD -- tests/test_task_broadcasting.py` + +Added synchronized regressions for transient child cancellation diagnostics, +cross-receiver terminal cleanup, durable EOC peer wakeups, root broadcaster +cleanup through both public wrappers, `MsgStream.aclose()` subclass +compatibility, and zero-buffer rejection. + +Verification output: + +```text +........................... [100%] +27 passed in 5.89s + +. [100%] +1 passed in 1.22s + +. [100%] +1 passed in 0.88s +``` + +The integration runs cover real `MsgStream` actor fan-out and infected-asyncio +`LinkedTaskChannel` fan-out. Ruff, Python compilation and `git diff --check` +passed. Repeated adversarial review found and resolved root close re-entrancy, +cross-receiver terminal retention, durable-EOC and subclass-compatibility +issues; final review reported no findings. diff --git a/tests/test_task_broadcasting.py b/tests/test_task_broadcasting.py index f6d0483e..31e18b30 100644 --- a/tests/test_task_broadcasting.py +++ b/tests/test_task_broadcasting.py @@ -409,6 +409,371 @@ def test_broadcast_statistics_report_queued_counts() -> None: trio.run(main) +def test_cancelled_reader_diagnostics_are_transient() -> None: + ''' + Cancelled-reader diagnostics must not retain stale `Task`s. + + `BroadcastState.cancelled` previously accumulated every source + owner cancelled during `BroadcastReceiver.receive()`. Even after + that receiver successfully read again or its subscription closed, + `BroadcastState.statistics()` retained the old `Task`, reporting + stale state and keeping the completed task alive. + + Cancel one child's source read under a receiver-local scope and + verify its task is reported. Reuse that same receiver for one + successful read to prove progress clears the entry. Cancel it once + more, then leave the subscription and prove close also removes the + diagnostic while the root receiver remains registered. + + ''' + async def main() -> None: + tx, rx = trio.open_memory_channel(1) + brx = broadcast_receiver(rx, 1) + cancel_scope = trio.CancelScope() + child_key: int + child_task = None + + async with brx.subscribe() as child: + child_key = child.key + + async def cancel_source_read() -> None: + nonlocal child_task + child_task = current_task() + with cancel_scope: + await child.receive() + assert cancel_scope.cancelled_caught + + async with trio.open_nursery() as nursery: + nursery.start_soon(cancel_source_read) + while brx._state.recv_ready is None: + await trio.lowlevel.checkpoint() + cancel_scope.cancel() + + stats = brx._state.statistics() + assert child_task is not None + assert stats['tasks_cancelled'] == { + child_key: child_task, + } + + await tx.send(1) + assert await child.receive() == 1 + assert not brx._state.cancelled + + cancel_scope = trio.CancelScope() + async with trio.open_nursery() as nursery: + nursery.start_soon(cancel_source_read) + while brx._state.recv_ready is None: + await trio.lowlevel.checkpoint() + cancel_scope.cancel() + + assert child_key in brx._state.cancelled + + assert child_key not in brx._state.cancelled + assert brx.key in brx._state.subs + + trio.run(main) + + +@pytest.mark.parametrize( + 'terminal_exc', + [ + trio.EndOfChannel(), + RuntimeError('terminal source failure'), + ], + ids=['end-of-channel', 'receive-error'], +) +def test_terminal_broadcast_clears_cancelled_tasks( + terminal_exc: Exception, +) -> None: + ''' + Terminal broadcast state must release every cancelled `Task`. + + A receiver which owned and cancelled a source read can leave its + task in `BroadcastState.cancelled`. If another receiver later gets + EOC or a terminal source failure, no subscriber can make source + progress to clear that stale diagnostic. Clearing only the terminal + owner's key therefore retained the first receiver's completed task. + + Cancel a child during the first controlled source read, then let + the root own a second read which raises EOC or `RuntimeError`. + Prove each terminal path clears the other receiver's diagnostic + before propagating its exact source outcome. + + ''' + class TerminalReceiver: + ''' + Block one cancellable read, then raise a terminal outcome. + + ''' + def __init__(self) -> None: + self.calls = 0 + self.first_started = trio.Event() + + async def receive(self) -> None: + ''' + Drive cancellation followed by terminal source state. + + ''' + self.calls += 1 + if self.calls == 1: + self.first_started.set() + await trio.sleep_forever() + + raise terminal_exc + + async def main() -> None: + source = TerminalReceiver() + brx = broadcast_receiver(source, 1) + cancel_scope = trio.CancelScope() + + async with brx.subscribe() as child: + async def cancel_child_read() -> None: + with cancel_scope: + await child.receive() + assert cancel_scope.cancelled_caught + + async with trio.open_nursery() as nursery: + nursery.start_soon(cancel_child_read) + await source.first_started.wait() + cancel_scope.cancel() + + assert child.key in brx._state.cancelled + with pytest.raises(type(terminal_exc)) as exc_info: + await brx.receive() + assert exc_info.value is terminal_exc + assert not brx._state.cancelled + + trio.run(main) + + +def test_end_of_channel_is_terminal_for_waiting_peer() -> None: + ''' + EOC must not let an awakened peer re-enter the closed source. + + `BroadcastState.eoc` was set when one source owner received EOC, + but neither receive path consulted it. A peer waiting behind that + owner therefore woke, saw no queued value, and started a second + source read. Cancellation at that checkpoint could repopulate + `BroadcastState.cancelled` after the broadcast became terminal. + + Block one child in the sole source read while the root waits on its + event, then release EOC. Both receivers must terminate from that + one source call, and the root's later receive must replay EOC + immediately without retaining cancellation diagnostics. + + ''' + class EOCReceiver: + ''' + Publish one controlled EOC and reject any second source read. + + ''' + def __init__(self) -> None: + self.calls = 0 + self.started = trio.Event() + self.release = trio.Event() + + async def receive(self) -> None: + ''' + Block the only valid source read until EOC release. + + ''' + self.calls += 1 + assert self.calls == 1 + self.started.set() + await self.release.wait() + raise trio.EndOfChannel + + async def main() -> None: + source = EOCReceiver() + brx = broadcast_receiver(source, 1) + outcomes: list[str] = [] + + async with brx.subscribe() as child: + async def receive_eoc( + receiver, + name: str, + ) -> None: + with pytest.raises(trio.EndOfChannel): + await receiver.receive() + outcomes.append(name) + + async with trio.open_nursery() as nursery: + nursery.start_soon(receive_eoc, child, 'child') + await source.started.wait() + nursery.start_soon(receive_eoc, brx, 'root') + + _, event = brx._state.recv_ready + while not event.statistics().tasks_waiting: + await trio.lowlevel.checkpoint() + source.release.set() + + with pytest.raises(trio.EndOfChannel): + await brx.receive() + + assert sorted(outcomes) == ['child', 'root'] + assert source.calls == 1 + assert not brx._state.cancelled + + trio.run(main) + + +def test_msgstream_eoc_close_preserves_aclose_override() -> None: + ''' + Internal EOC cleanup must preserve the public `aclose()` contract. + + Passing a new private keyword from `MsgStream.receive()` to + `self.aclose()` broke subclasses whose compatible override kept + the original zero-argument signature. Use a minimal subclass which + records virtual dispatch and delegates to the base implementation. + Drive graceful EOC through the real root broadcaster and prove the + override runs without closing that active root re-entrantly. + + ''' + class Stream(tractor.MsgStream): + ''' + Record public close dispatch with the established signature. + + ''' + close_calls = 0 + + async def aclose(self): + ''' + Delegate closure without accepting private arguments. + + ''' + self.close_calls += 1 + return await super().aclose() + + class PldRx: + ''' + Delegate source receive and terminate the close drain. + + ''' + def __init__(self, rx) -> None: + self._rx = rx + + async def recv_pld(self, **kwargs): + ''' + Receive directly from the test source channel. + + ''' + return await self._rx.receive() + + def recv_msg_nowait(self, **kwargs): + ''' + Report EOC to finish `MsgStream.aclose()` draining. + + ''' + raise trio.EndOfChannel + + async def main() -> None: + tx, rx = trio.open_memory_channel(1) + ctx = SimpleNamespace( + cid='test-context', + _pld_rx=PldRx(rx), + send_stop=lambda: trio.lowlevel.checkpoint(), + side='caller', + peer_side='callee', + maybe_raise=lambda **kwargs: None, + ) + stream = Stream(ctx, rx) + + async with stream.subscribe(): + await tx.aclose() + with pytest.raises(trio.EndOfChannel): + await stream.receive() + assert stream.close_calls == 1 + assert not stream._broadcaster._closed + + trio.run(main) + + +@pytest.mark.parametrize( + 'close_wrapper', + [ + tractor.MsgStream.aclose, + LinkedTaskChannel.aclose, + ], + ids=['msg-stream', 'linked-task-channel'], +) +def test_wrapper_close_clears_root_cancelled_task( + close_wrapper, +) -> None: + ''' + Public stream close must release root cancellation diagnostics. + + Root broadcasters allocated by `MsgStream.subscribe()` and + `LinkedTaskChannel.subscribe()` are private implementation state. + If their source receive was cancelled, callers had no public way + to close the root, so wrapper teardown retained the completed + `Task` in `BroadcastState.cancelled` indefinitely. + + Cancel a root source read, attach that broadcaster to a minimal + public wrapper, and close it through each real `aclose()` method. + The root receiver and its task diagnostic must both be removed; + for `MsgStream`, pre-close the source to cover its idempotent early + return path. + + ''' + async def main() -> None: + _, rx = trio.open_memory_channel(1) + brx = broadcast_receiver(rx, 1) + cancel_scope = trio.CancelScope() + + async def cancel_source_read() -> None: + with cancel_scope: + await brx.receive() + assert cancel_scope.cancelled_caught + + async with trio.open_nursery() as nursery: + nursery.start_soon(cancel_source_read) + while brx._state.recv_ready is None: + await trio.lowlevel.checkpoint() + cancel_scope.cancel() + + assert brx.key in brx._state.cancelled + + if close_wrapper is tractor.MsgStream.aclose: + ctx = SimpleNamespace(cid='test-context') + wrapper = tractor.MsgStream(ctx, rx) + wrapper._broadcaster = brx + await rx.aclose() + else: + wrapper = SimpleNamespace( + _broadcaster=brx, + _from_aio=rx, + ) + + await close_wrapper(wrapper) + assert brx.key not in brx._state.subs + assert brx.key not in brx._state.cancelled + + trio.run(main) + + +def test_broadcast_rejects_zero_buffer_size() -> None: + ''' + A broadcaster must retain at least one value for peer fan-out. + + `collections.deque(maxlen=0)` silently discards every appended + value, so `broadcast_receiver(..., 0)` allowed the source owner to + receive while peer cursors advanced into an always-empty queue. + Their lag recovery then reset to index `-1` and recursively retried + without any retained value to consume. + + Construct a rendezvous memory channel and prove broadcaster setup + rejects its zero capacity synchronously with a clear public error, + before any receiver is registered or source receive can begin. + + ''' + _, rx = trio.open_memory_channel(0) + with pytest.raises( + ValueError, + match='`max_buffer_size` must be greater than zero', + ): + broadcast_receiver(rx, 0) + + def test_underlying_receive_failure_wakes_all_subscribers() -> None: ''' A shared receive failure must terminate every broadcast receiver. diff --git a/tractor/_streaming.py b/tractor/_streaming.py index a5e94f19..93c7d12f 100644 --- a/tractor/_streaming.py +++ b/tractor/_streaming.py @@ -103,6 +103,12 @@ class MsgStream(trio.abc.Channel): self._eoc: bool|trio.EndOfChannel = False self._closed: bool|trio.ClosedResourceError = False + # `MsgStream.receive()` sets this while it calls + # `MsgStream.aclose()` after source EOC. That close is + # re-entrant from the root `BroadcastReceiver._recv`, so it + # must not cancel the same receiver before EOC propagates. + self._eoc_close_task: trio.lowlevel.Task|None = None + @property def ctx(self) -> Context: ''' @@ -256,7 +262,16 @@ class MsgStream(trio.abc.Channel): # when the send is closed we assume the stream has # terminated and signal this local iterator to stop - drained: list[Exception|dict] = await self.aclose() + # + # Preserve virtual dispatch through the public zero-argument + # `MsgStream.aclose()` API. The task marker lets the base + # implementation distinguish this receive-internal close from + # an explicit caller or `MsgStream.__aexit__()` close. + self._eoc_close_task = trio.lowlevel.current_task() + try: + drained: list[Exception|dict] = await self.aclose() + finally: + self._eoc_close_task = None if drained: # ^^^^^^^^TODO? pass these to the `._ctx._drained_msgs: # deque` and then iterate them as part of any @@ -335,6 +350,20 @@ class MsgStream(trio.abc.Channel): # `.__aexit__()` as well!!! # => SO ENSURE WE CATCH ALL TERMINATION STATES in this # block including the EoC.. + + # `MsgStream.subscribe()` stores its hidden root broadcaster + # on `self._broadcaster`. Explicit teardown owns that root and + # must close it to release its subscriber and cancelled-task + # diagnostic. Skip only the receive-internal EOC close above: + # cancelling the active root's source-read scope there would + # turn graceful EOC into `trio.ClosedResourceError`. + if ( + trio.lowlevel.current_task() is not self._eoc_close_task + and + (broadcaster := self._broadcaster) is not None + ): + await broadcaster.aclose() + if self.closed: # this stream has already been closed so silently succeed as # per ``trio.AsyncResource`` semantics. diff --git a/tractor/to_asyncio.py b/tractor/to_asyncio.py index 28756052..f8440476 100644 --- a/tractor/to_asyncio.py +++ b/tractor/to_asyncio.py @@ -213,6 +213,14 @@ class LinkedTaskChannel( _broadcaster: BroadcastReceiver|None = None async def aclose(self) -> None: + # `LinkedTaskChannel.subscribe()` lazily allocates and retains + # this root receiver. Close it first so its receiver-local + # source-read scope and cancellation diagnostics are released + # before `self._from_aio` becomes inaccessible; child + # subscriptions retain their own independent close lifetimes. + if (broadcaster := self._broadcaster) is not None: + await broadcaster.aclose() + await self._from_aio.aclose() # ?TODO? async version of this? diff --git a/tractor/trionics/_broadcast.py b/tractor/trionics/_broadcast.py index 5ac21922..cfe8f2d7 100644 --- a/tractor/trionics/_broadcast.py +++ b/tractor/trionics/_broadcast.py @@ -142,7 +142,8 @@ class BroadcastState(Struct): # before this failure is replayed at each receiver's boundary. receive_exc: Exception | None = None - # If the broadcaster was cancelled, we might as well track it + # Retain the latest interrupted source-reader task until its + # receiver next makes progress or closes. cancelled: dict[int, Task] = {} def statistics(self) -> dict[str, Any]: @@ -286,6 +287,7 @@ class BroadcastReceiver(ReceiveChannel): return self.receive_nowait(_key, _state) state.subs[key] -= 1 + state.cancelled.pop(key, None) return value receive_exc = state.receive_exc @@ -297,6 +299,9 @@ class BroadcastReceiver(ReceiveChannel): 'Shared broadcast receiver failed' ) from receive_exc + if state.eoc: + raise trio.EndOfChannel + raise trio.WouldBlock async def _receive_from_underlying( @@ -363,6 +368,8 @@ class BroadcastReceiver(ReceiveChannel): ): state.subs[sub_key] += 1 + state.cancelled.pop(key, None) + # NOTE: this should ONLY be set if the above task was *NOT* # cancelled on the `._recv()` call. event.set() @@ -372,6 +379,7 @@ class BroadcastReceiver(ReceiveChannel): # if any one consumer gets an EOC from the underlying # receiver we need to unblock and send that signal to # all other consumers. + state.cancelled.clear() self._state.eoc = True if event.statistics().tasks_waiting: event.set() @@ -402,6 +410,7 @@ class BroadcastReceiver(ReceiveChannel): # so any non-EOC failure terminates the entire broadcast. # Publish it before waking peers so they can drain their # retained values and then observe the same failure. + state.cancelled.clear() state.receive_exc = receive_exc if event.statistics().tasks_waiting: event.set() @@ -411,6 +420,7 @@ class BroadcastReceiver(ReceiveChannel): # Process-control and cancellation-like exceptions must # not become durable broadcast state, but peers still # need waking before `recv_ready` is cleared. + state.cancelled.pop(key, None) if event.statistics().tasks_waiting: event.set() raise @@ -547,6 +557,7 @@ class BroadcastReceiver(ReceiveChannel): # up to the last received that still reside in the queue. state = self._state state.subs.pop(self.key) + state.cancelled.pop(self.key, None) self._closed = True # A non-owner close must not wake peers waiting behind some @@ -581,6 +592,11 @@ def broadcast_receiver( ) -> BroadcastReceiver: + if max_buffer_size < 1: + raise ValueError( + '`max_buffer_size` must be greater than zero' + ) + return BroadcastReceiver( recv_chan, state=BroadcastState(