Compare commits

...

6 Commits

Author SHA1 Message Date
Gud Boi d1edecbf60 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`))
2026-08-13 15:49:05 -04:00
Gud Boi a2e0df4ba1 Expose stream subscriber lag policy
`MsgStream.subscribe()` and `LinkedTaskChannel.subscribe()` omitted
`BroadcastReceiver.raise_on_lag`, forcing downstream consumers to
mutate a private receiver attribute when overruns were acceptable.

Add `raise_on_lag` to both public wrappers. The first subscription
sets the irreversible root broadcaster's policy, while every child
selects its own strict or warn/drop/resume behavior independently.

Document both fan-out APIs. Cover policy forwarding plus real IPC
and infected-asyncio paths.

Prompt-IO: ai/prompt-io/opencode/20260812T213117Z_51185487_prompt_io.md

(this patch was generated in some part by `opencode` using `gpt-5.6-sol` (`openai`))
2026-08-13 00:37:07 -04:00
Gud Boi 511854870b Isolate `BroadcastReceiver.aclose()` wakeups
Closing any subscriber set the shared `recv_ready` event, even when
another receiver owned the source read. Waiting peers then looped
until an idle source produced another value.

Give each receiver private source-read and peer-wait cancellation
scopes. Closing a waiting peer interrupts only that peer. Closing
the source owner discards post-close source outcomes and then wakes
peers for a clean ownership handoff.

Keep outer task cancellation as `trio.Cancelled`; only explicit
receiver close maps either private scope's cancellation to
`ClosedResourceError`. Assert that scope cancellation implies the
receiver is closed and document the source-owner key check.

Prompt-IO: ai/prompt-io/opencode/20260812T150027Z_c2a6ccef_prompt_io.md

(this patch was generated in some part by `opencode` using `gpt-5.6-sol` (`openai`))
2026-08-12 17:22:59 -04:00
Gud Boi c2a6ccefd0 Wake broadcast peers on shared receive failures
Only `EndOfChannel` and direct cancellation woke tasks waiting
behind the subscriber which owned the underlying receive. Any other
failure cleared `BroadcastState.recv_ready` while peers remained
blocked on its unreachable event.

Publish ordinary receive exceptions as terminal broadcast state.
The owner keeps the original failure while peers drain retained
values and then raise `BroadcastReceiveError` from that cause. Also
wake peers on process-control exits without retaining them as state.

Document the public owner/peer contract and cover current, late and
control-flow subscribers with deterministic bounded regressions.

Prompt-IO: ai/prompt-io/opencode/20260812T030608Z_1095e7f7_prompt_io.md

(this patch was generated in some part by `opencode` using `gpt-5.6-sol` (`openai`))
2026-08-12 00:51:22 -04:00
Gud Boi 1095e7f710 Fix `BroadcastState.statistics()` queue counts
`BroadcastState.subs` stores each receiver's next unread deque
index, but `.statistics()` exposed that index as a queue length. A
caught-up receiver looked correct by accident while every queued
count was one short.

Convert cursors to retained, receivable counts and clamp lagged
receivers to the current queue length. Also avoid deprecated
`trio.Event` truthiness when reporting waiter counts.

Cover caught-up, queued, lagged and real-event states using actual
broadcast sends and receives.

Prompt-IO: ai/prompt-io/opencode/20260812T012324Z_06c4af17_prompt_io.md

(this patch was generated in some part by `opencode` using `gpt-5.6-sol` (`openai`))
2026-08-11 22:49:23 -04:00
Gud Boi 06c4af17e4 Fix `BroadcastReceiver` lag counts
`BroadcastReceiver.receive_nowait()` treated `seq` as a deque index
but subtracted `BroadcastState.maxlen` without counting the first
invalid index. A one-slot queue thus claimed it dropped zero values
after its subscriber missed one.

Include that first displaced value in the count. Preserve the
existing Tokio-style reset to the oldest retained item.

Also, cover exact loss reporting and recovery for one- and
three-slot retention windows.

Prompt-IO: ai/prompt-io/opencode/20260811T233833Z_7cbd64ee_prompt_io.md

(this patch was generated in some part by `opencode` using `gpt-5.6-sol` (`openai`))
2026-08-11 19:45:53 -04:00
20 changed files with 1595 additions and 17 deletions

View File

@ -0,0 +1,30 @@
---
model: openai/gpt-5.6-sol
service: opencode
session: ses_0799212ebffe42arY96czXn89F
timestamp: 2026-08-11T23:38:33Z
git_ref: 7cbd64ee
scope: code
substantive: true
raw_file: 20260811T233833Z_7cbd64ee_prompt_io.raw.md
---
## Prompt
Open a new isolated worktree in the local Tractor repository and draft a fix
for `BroadcastReceiver` reporting that a lagged one-slot consumer dropped
zero values when one value had actually been displaced.
## Response summary
Corrected the off-by-one lag count while preserving cursor recovery and added
deterministic narrow- and wider-window regressions for exact loss reporting.
## Files changed
- `tractor/trionics/_broadcast.py` - exact broadcast overrun count.
- `tests/test_task_broadcasting.py` - lag count and recovery regression.
## Human edits
None - generated output follows the user's diagnosed edge case.

View File

@ -0,0 +1,50 @@
---
model: openai/gpt-5.6-sol
service: opencode
timestamp: 2026-08-11T23:38:33Z
git_ref: 7cbd64ee
diff_cmd: git diff HEAD~1..HEAD
---
The user asked for a Tractor fix in a new isolated worktree after a live piker
failure reported:
```text
tractor.trionics._broadcast.Lagged:
Task `piker.brokers.ib.broker.handle_order_requests` overrun and
dropped `0` values
```
Inspection showed the lag exception was valid but its count was off by one.
`BroadcastReceiver.receive_nowait()` treats `seq` as a deque index. With a
one-entry queue, index zero is the only retained value and `seq == 1` already
means one value was displaced. The old `seq - maxlen` calculation therefore
reported zero instead of one.
> `git diff HEAD~1..HEAD -- tractor/trionics/_broadcast.py`
Adjusted the lag count to `seq - maxlen + 1` and documented why the first
invalid deque index must be included. The existing Tokio-style cursor reset
remains unchanged.
> `git diff HEAD~1..HEAD -- tests/test_task_broadcasting.py`
Added a deterministic parameterized regression covering a one-slot queue
with one dropped value and a three-slot queue with two dropped values. The
test keeps the root receiver idle while a child subscriber synchronously
drains each produced value, asserts the exact `Lagged` message, and proves the
next receive resumes at the oldest retained item.
Verification output:
```text
.. [100%]
2 passed in 0.04s
.......... [100%]
10 passed in 4.65s
```
The targeted import resolved to the new Tractor worktree. Python compilation
and `git diff --check` passed. Adversarial review found no actionable issues;
zero-capacity channels remain a pre-existing untested edge outside this fix.

View File

@ -0,0 +1,31 @@
---
model: openai/gpt-5.6-sol
service: opencode
session: ses_0799212ebffe42arY96czXn89F
timestamp: 2026-08-12T01:23:24Z
git_ref: 06c4af17
scope: code
substantive: true
raw_file: 20260812T012324Z_06c4af17_prompt_io.raw.md
---
## Prompt
Iteratively refine Tractor PR 490. For item one, correct broadcast statistics
queue counts and Trio event checks, verify and review the exact change, then
return a complete commit plan before proceeding.
## Response summary
Converted subscriber cursor indexes into clamped retained queue counts,
removed deprecated event truthiness, and added deterministic state and
deprecation regressions.
## Files changed
- `tractor/trionics/_broadcast.py` - accurate queue and waiter statistics.
- `tests/test_task_broadcasting.py` - retained-count and event regression.
## Human edits
None - generated output follows the requested first iterative item.

View File

@ -0,0 +1,37 @@
---
model: openai/gpt-5.6-sol
service: opencode
timestamp: 2026-08-12T01:23:24Z
git_ref: 06c4af17
diff_cmd: git diff HEAD~1..HEAD
---
After opening draft Tractor PR 490, the user requested an iterative pass over
additional broadcast subsystem findings. The first item was to correct
`BroadcastState.statistics()` queued counts and its deprecated Trio event
truthiness check, then stop for a complete commit plan.
> `git diff HEAD~1..HEAD -- tractor/trionics/_broadcast.py`
Changed `queued_len_by_task` from raw deque cursor indexes to retained,
receivable counts. Caught-up `-1` reports zero, valid indexes report index plus
one, and lagged cursors clamp to the current retained queue length. Replaced
`trio.Event` truthiness with an explicit `is not None` branch.
> `git diff HEAD~1..HEAD -- tests/test_task_broadcasting.py`
Added a deterministic statistics regression using actual sends and receives.
It verifies caught-up and one-queued states, drives a root receiver beyond a
three-slot retention window to prove clamping, and installs a real
`trio.Event` while treating deprecations as errors.
Verification output:
```text
........... [100%]
11 passed in 5.76s
```
Python compilation and `git diff --check` passed. Initial adversarial review
caught unclamped lagged cursors and an ineffective event test; both were
fixed. Final review found no actionable issues.

View File

@ -0,0 +1,33 @@
---
model: openai/gpt-5.6-sol
service: opencode
session: ses_0799212ebffe42arY96czXn89F
timestamp: 2026-08-12T03:06:08Z
git_ref: 1095e7f7
scope: code
substantive: true
raw_file: 20260812T030608Z_1095e7f7_prompt_io.raw.md
---
## Prompt
For Tractor PR 490 item two, make shared underlying receive failures wake and
terminate every broadcast subscriber without losing retained values. Review,
verify and return a complete commit plan before continuing.
## Response summary
Published ordinary receive failures as terminal broadcast state, introduced
a public chained peer exception, kept control-flow exits transient while
waking peers, documented the contract, and added deterministic regressions.
## Files changed
- `tractor/trionics/_broadcast.py` - terminal failure and peer wake protocol.
- `tractor/trionics/__init__.py` - public peer exception export.
- `docs/api/trionics.rst` - failure-delivery API contract.
- `tests/test_task_broadcasting.py` - terminal and transient failure tests.
## Human edits
None - generated output follows the requested second iterative item.

View File

@ -0,0 +1,49 @@
---
model: openai/gpt-5.6-sol
service: opencode
timestamp: 2026-08-12T03:06:08Z
git_ref: 1095e7f7
diff_cmd: git diff HEAD~1..HEAD
---
The user requested the second iterative refinement for Tractor PR 490: ensure
non-EOC failures from a shared underlying broadcast receiver do not leave peer
subscribers blocked forever, then stop for a complete commit plan.
> `git diff HEAD~1..HEAD -- tractor/trionics/_broadcast.py`
Added shared terminal failure publication for ordinary `Exception` values.
The receive owner gets the original exception; peers may drain retained
values and then get a fresh `BroadcastReceiveError` chained from the original.
Late subscribers observe the same terminal state without retrying the failed
underlying receiver. Process-control and cancellation-like `BaseException`
values wake peers but are re-raised without becoming durable channel state.
> `git diff HEAD~1..HEAD -- tractor/trionics/__init__.py`
Exported `BroadcastReceiveError` as the public peer-delivery exception.
> `git diff HEAD~1..HEAD -- docs/api/trionics.rst`
Documented `BroadcastReceiveError` and the owner-versus-peer delivery
contract, including retained-value draining and late subscribers.
> `git diff HEAD~1..HEAD -- tests/test_task_broadcasting.py`
Added deterministic bounded regressions. One scripts a successful receive
followed by `RuntimeError`, proving the root drains retained data, all current
and late receivers observe terminal failure, and the source is not retried.
The second scripts a custom `BaseException`, proving peers wake and take over
the next source receive without retaining control-flow state.
Verification output:
```text
............. [100%]
13 passed in 5.62s
```
Python compilation and `git diff --check` passed. Iterative adversarial review
drove independent peer exception wrappers, ordinary-versus-control-flow
classification, bounded test completion, public docs, and the final catch-all
peer wake. Final review found no issues.

View File

@ -0,0 +1,33 @@
---
model: openai/gpt-5.6-sol
service: opencode
session: ses_0799212ebffe42arY96czXn89F
timestamp: 2026-08-12T15:00:27Z
git_ref: c2a6ccef
scope: code
substantive: true
raw_file: 20260812T150027Z_c2a6ccef_prompt_io.raw.md
---
## Prompt
For Tractor PR 490 item three, prevent subscriber closure from waking another
receiver's shared event or stranding peers. Preserve close-safe source-read
ownership handoff, then review, verify and return a complete commit plan.
## Response summary
Introduced receiver-local wait/source cancellation, owner-specific handoff,
and close precedence over shielded source values/errors/EOC. Clarified that
private scope cancellation means explicit close while outer task cancellation
remains `trio.Cancelled` for both source-read and peer-wait scopes, with
deterministic receiver-close regressions.
## Files changed
- `tractor/trionics/_broadcast.py` - receiver-local close and ownership scopes.
- `tests/test_task_broadcasting.py` - peer close and owner handoff regressions.
## Human edits
None - generated output follows the requested third iterative item.

View File

@ -0,0 +1,51 @@
---
model: openai/gpt-5.6-sol
service: opencode
timestamp: 2026-08-12T15:00:27Z
git_ref: c2a6ccef
diff_cmd: git diff HEAD~1..HEAD
---
The user requested the third iterative refinement for Tractor PR 490: closing
one broadcast subscriber must not set another receiver owner's shared event
and create a runnable hot loop, then stop for a complete commit plan.
> `git diff HEAD~1..HEAD -- tractor/trionics/_broadcast.py`
Added receiver-local wait cancellation and source-read ownership scopes.
Closing a non-owner waiting behind another source reader cancels only that
receiver's private wait and maps it to `ClosedResourceError`; the shared event
remains untouched. Closing the active source owner cancels only its private
source-read scope, wakes peers after cleanup, and lets one peer take ownership.
Source outcomes are captured inside the owner scope and classified only after
checking close/cancel state. A cancellation-shielding source therefore cannot
publish a returned value, ordinary error, or EOC after its owner was closed.
The private scope's `cancel_called` bit is asserted to imply receiver closure;
outer task cancellation remains `trio.Cancelled` and is not translated into
`ClosedResourceError`. Owner-key comments document that only the receiver
identified by `recv_ready[0]` may cancel the shared source-read scope. The
same explicit-close invariant is enforced symmetrically for private peer-wait
scope cancellation.
> `git diff HEAD~1..HEAD -- tests/test_task_broadcasting.py`
Added deterministic bounded regressions for both close positions. The
non-owner test places two peers behind an active source read, closes one and
proves only that peer exits while the shared event stays unset. The owner test
closes a source owner whose receive shields cancellation and parameterizes a
returned value, `RuntimeError`, and `EndOfChannel`; each discarded outcome
hands the next source receive to the waiting root without terminal-state or
EOC publication.
Verification output:
```text
................. [100%]
17 passed in 5.84s
```
Python compilation and `git diff --check` passed. Iterative adversarial review
caught a waiting non-owner hang, cancellation-shielded source returns, and
shielded source exceptions. All were fixed. Final review found no actionable
issues.

View File

@ -0,0 +1,34 @@
---
model: openai/gpt-5.6-sol
service: opencode
session: ses_0799212ebffe42arY96czXn89F
timestamp: 2026-08-12T21:31:17Z
git_ref: 51185487
scope: code
substantive: true
raw_file: 20260812T213117Z_51185487_prompt_io.raw.md
---
## Prompt
For Tractor PR 490 item four, expose `raise_on_lag` through the public IPC and
asyncio linked-channel subscription wrappers. Review, verify and return a
complete commit plan before applying the API downstream in piker.
## Response summary
Added public per-subscription lag policy to both wrappers, preserved first-call
root policy, documented the semantics, and covered forwarding plus real fan-out
paths.
## Files changed
- `tractor/_streaming.py` - `MsgStream` lag policy forwarding.
- `tractor/to_asyncio.py` - linked-channel lag policy forwarding.
- `docs/guide/streaming.rst` - IPC fan-out policy docs.
- `docs/guide/asyncio.rst` - linked-channel fan-out policy docs.
- `tests/test_task_broadcasting.py` - wrapper policy regression.
## Human edits
None - generated output follows the requested fourth iterative item.

View File

@ -0,0 +1,53 @@
---
model: openai/gpt-5.6-sol
service: opencode
timestamp: 2026-08-12T21:31:17Z
git_ref: 51185487
diff_cmd: git diff HEAD~1..HEAD
---
The user requested the fourth iterative refinement for Tractor PR 490: expose
subscriber lag policy through the public `MsgStream.subscribe()` and
`LinkedTaskChannel.subscribe()` wrappers, then stop for a complete commit
plan. This enables piker to replace private receiver mutation.
> `git diff HEAD~1..HEAD -- tractor/_streaming.py`
Added `raise_on_lag: bool = True` to `MsgStream.subscribe()`. The first call
passes the policy to the irreversibly allocated root broadcaster and its
child; later calls configure each child independently while retaining the
root's first-call policy.
> `git diff HEAD~1..HEAD -- tractor/to_asyncio.py`
Added equivalent lag-policy forwarding to `LinkedTaskChannel.subscribe()`.
> `git diff HEAD~1..HEAD -- docs/guide/streaming.rst`
> `git diff HEAD~1..HEAD -- docs/guide/asyncio.rst`
Documented strict versus warn/drop/resume behavior, independent child policy,
and first-call root policy for both wrapper types.
> `git diff HEAD~1..HEAD -- tests/test_task_broadcasting.py`
Added a parameterized wrapper-level regression using minimal receive-compatible
handles. It verifies a first non-raising subscription configures root and
child, then a later strict child does not mutate the sticky root policy.
Verification output:
```text
................... [100%]
19 passed in 5.71s
.................... [100%]
20 passed in 6.88s
. [100%]
1 passed in 0.86s
```
The second and third runs cover actual `MsgStream` and infected-asyncio
`LinkedTaskChannel` fan-out respectively. Python compilation and
`git diff --check` passed. Adversarial review found no actionable issues.

View File

@ -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.

View File

@ -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.

View File

@ -40,6 +40,9 @@ Broadcast fan-out
.. autoexception:: Lagged .. autoexception:: Lagged
:show-inheritance: :show-inheritance:
.. autoexception:: BroadcastReceiveError
:show-inheritance:
A single-producer, many-consumer broadcast layer over any A single-producer, many-consumer broadcast layer over any
``trio``-style receive channel: non-lossy for the *fastest* ``trio``-style receive channel: non-lossy for the *fastest*
consumer while slower consumers raise :class:`Lagged` (a consumer while slower consumers raise :class:`Lagged` (a
@ -48,6 +51,13 @@ internal ring. This is exactly the machinery behind
:meth:`tractor.MsgStream.subscribe` — see :meth:`tractor.MsgStream.subscribe` — see
``examples/streaming_broadcast_fanout.py``. ``examples/streaming_broadcast_fanout.py``.
If the shared underlying receiver raises an ordinary exception, the
subscriber which owned that receive gets the original failure.
Waiting peers drain their retained values and then raise
:class:`BroadcastReceiveError`, with the original failure available
as ``__cause__``. Later subscribers observe the same terminal state
without retrying the failed underlying receiver.
ExceptionGroup helpers ExceptionGroup helpers
---------------------- ----------------------

View File

@ -209,6 +209,11 @@ The underlying broadcast machinery is lazily allocated on first
use and is *not* reversible for the channel's remaining lifetime, use and is *not* reversible for the channel's remaining lifetime,
so only reach for it when you actually want the fan-out. so only reach for it when you actually want the fan-out.
As with ``MsgStream``, pass ``raise_on_lag=False`` for a consumer
which may warn, drop old values and resume from the retained window.
Each child chooses independently; the first subscription also fixes
the linked channel's root receive policy.
One-shot calls with ``run_task()`` One-shot calls with ``run_task()``
---------------------------------- ----------------------------------
When you just want a single ``asyncio`` result and no streaming When you just want a single ``asyncio`` result and no streaming

View File

@ -171,6 +171,12 @@ keeps pace with the *fastest* subscriber; a task falling more
than the buffered window behind has its next receive raise than the buffered window behind has its next receive raise
``tractor.trionics.Lagged`` to say it lost data. ``tractor.trionics.Lagged`` to say it lost data.
Pass ``raise_on_lag=False`` when a consumer may drop old values and
resume from the oldest retained item instead. The receiver logs the
overrun rather than raising. Each child subscription chooses its own
policy; the first call also fixes the policy of the stream's root
receive handle because broadcaster allocation is irreversible.
The broadcast handle stays duplex btw: it proxies ``send()`` The broadcast handle stays duplex btw: it proxies ``send()``
through to the underlying stream, so each subscriber task can through to the underlying stream, so each subscriber task can
keep talking upstream while consuming its fan-out copy. keep talking upstream while consuming its fan-out copy.

View File

@ -8,14 +8,18 @@ from contextlib import (
from functools import partial from functools import partial
from itertools import cycle from itertools import cycle
import time import time
from types import SimpleNamespace
from typing import Optional from typing import Optional
import warnings
import pytest import pytest
import trio import trio
from trio.lowlevel import current_task from trio.lowlevel import current_task
import tractor import tractor
from tractor.to_asyncio import LinkedTaskChannel
from tractor.trionics import ( from tractor.trionics import (
broadcast_receiver, broadcast_receiver,
BroadcastReceiveError,
Lagged, Lagged,
collapse_eg, collapse_eg,
) )
@ -307,6 +311,801 @@ def test_subscribe_errors_after_close():
trio.run(main) trio.run(main)
@pytest.mark.parametrize(
('size', 'sent', 'dropped'),
[
(1, 2, 1),
(3, 5, 2),
],
)
def test_lagged_reports_exact_drop_count(
size: int,
sent: int,
dropped: int,
) -> None:
'''
`Lagged` must report every value outside the retained window.
`BroadcastReceiver.receive_nowait()` previously subtracted the
queue length from an already-invalid deque index without counting
that first displaced value. A one-slot queue therefore claimed it
dropped zero values after two sends. Keep one root receiver idle
while a child subscriber drains every produced value, then prove
the lag error reports the exact overrun and positions the root at
the oldest value still retained by `BroadcastState.queue`.
'''
async def main() -> None:
tx, rx = trio.open_memory_channel(size)
brx = broadcast_receiver(rx, size)
async with brx.subscribe() as fast:
for value in range(sent):
await tx.send(value)
assert await fast.receive() == value
match = rf'dropped `{dropped}` values'
with pytest.raises(Lagged, match=match):
await brx.receive()
assert await brx.receive() == sent - size
trio.run(main)
def test_broadcast_statistics_report_queued_counts() -> None:
'''
`BroadcastState.statistics()` must report counts, not indexes.
Each `BroadcastState.subs` value is the deque index of a
receiver's next unread value, with `-1` meaning caught up. The
statistics API returned these indexes directly, so one queued
value appeared as zero and every positive count was one short.
Keep one root receiver idle while a child synchronously receives
four produced values. Prove the root count advances through one
and three retained values, then remains clamped to the three-slot
retention window after lagging.
Finally install an actual unwaited `trio.Event` in
`BroadcastState.recv_ready` while treating deprecations as errors.
This proves statistics checks `None` explicitly instead of using
deprecated `trio.Event` truthiness.
'''
async def main() -> None:
tx, rx = trio.open_memory_channel(3)
brx = broadcast_receiver(rx, 3)
async with brx.subscribe() as child:
state = brx._state
assert state.statistics()['queued_len_by_task'] == {
brx.key: 0,
child.key: 0,
}
await tx.send(0)
assert await child.receive() == 0
assert state.statistics()['queued_len_by_task'] == {
brx.key: 1,
child.key: 0,
}
for value in range(1, 4):
await tx.send(value)
assert await child.receive() == value
state.recv_ready = (child.key, trio.Event())
with warnings.catch_warnings():
warnings.simplefilter('error', DeprecationWarning)
stats = state.statistics()
assert stats['queued_len_by_task'] == {
brx.key: 3,
child.key: 0,
}
assert stats['tasks_waiting'] == 0
state.recv_ready = 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.
Previously, only `EndOfChannel` and receiver cancellation woke
peer tasks waiting on `BroadcastState.recv_ready`. If the shared
underlying receiver raised another error, its owner propagated
the failure and cleared the event while every peer remained
blocked forever.
Script one successful receive followed by a controlled
`RuntimeError`. Let a fast child own both underlying receives
while the root first drains its retained value and then waits on
the child's second receive. Release the failure only after both
tasks have reached those positions. Both exact errors prove the
peer was awakened without losing buffered data. A later
subscriber proves the terminal failure remains published for new
receivers instead of retrying the failed underlying channel.
'''
class FailingReceiver:
'''
Return one value, then fail after deterministic release.
'''
def __init__(self) -> None:
self.calls: int = 0
self.failure_started = trio.Event()
self.release_failure = trio.Event()
async def receive(self) -> int:
'''
Drive the scripted success-then-failure sequence.
'''
self.calls += 1
if self.calls == 1:
return 1
self.failure_started.set()
await self.release_failure.wait()
raise RuntimeError('underlying receive failed')
async def main() -> None:
source = FailingReceiver()
brx = broadcast_receiver(source, 3)
child_error: list[RuntimeError] = []
root_error: list[BroadcastReceiveError] = []
late_error: list[BroadcastReceiveError] = []
root_drained = trio.Event()
async def receive_child() -> None:
async with brx.subscribe() as child:
assert await child.receive() == 1
try:
await child.receive()
except RuntimeError as exc:
child_error.append(exc)
async def receive_root() -> None:
assert await brx.receive() == 1
root_drained.set()
try:
await brx.receive()
except BroadcastReceiveError as exc:
root_error.append(exc)
with trio.fail_after(1):
async with trio.open_nursery() as nursery:
nursery.start_soon(receive_child)
await source.failure_started.wait()
nursery.start_soon(receive_root)
await root_drained.wait()
source.release_failure.set()
assert source.calls == 2
assert [str(exc) for exc in child_error] == [
'underlying receive failed',
]
assert [str(exc) for exc in root_error] == [
'Shared broadcast receiver failed',
]
assert child_error[0] is not root_error[0]
assert root_error[0].__cause__ is child_error[0]
async with brx.subscribe() as late:
with pytest.raises(
BroadcastReceiveError,
match='Shared broadcast receiver failed',
) as exc_info:
await late.receive()
late_error.append(exc_info.value)
assert late_error[0] is not child_error[0]
assert late_error[0] is not root_error[0]
assert late_error[0].__cause__ is child_error[0]
assert source.calls == 2
trio.run(main)
def test_control_flow_exit_wakes_broadcast_peer() -> None:
'''
Non-terminal control flow must wake peers without being retained.
Process-control and cancellation-like `BaseException` values
should remain local to the task which receives them, but the old
owner still has to wake subscribers blocked on its shared event.
Make one child own a controlled `BaseException` receive while the
root waits behind it. After release, prove the child gets that
exact exit and the root takes ownership of the next underlying
receive instead of hanging or replaying the control-flow event.
'''
class ReceiveExit(BaseException):
'''
Model a non-terminal process-control receive exit.
'''
class ControlFlowReceiver:
'''
Raise one controlled exit, then return a value.
'''
def __init__(self) -> None:
self.calls: int = 0
self.exit_started = trio.Event()
self.release_exit = trio.Event()
async def receive(self) -> int:
'''
Drive the scripted control-flow-then-value sequence.
'''
self.calls += 1
if self.calls == 1:
self.exit_started.set()
await self.release_exit.wait()
raise ReceiveExit
return 2
async def main() -> None:
source = ControlFlowReceiver()
brx = broadcast_receiver(source, 3)
child_exit: list[ReceiveExit] = []
root_value: list[int] = []
async def receive_child() -> None:
async with brx.subscribe() as child:
try:
await child.receive()
except ReceiveExit as exc:
child_exit.append(exc)
async def receive_root() -> None:
root_value.append(await brx.receive())
with trio.fail_after(1):
async with trio.open_nursery() as nursery:
nursery.start_soon(receive_child)
await source.exit_started.wait()
nursery.start_soon(receive_root)
while True:
_, event = brx._state.recv_ready
if event.statistics().tasks_waiting:
break
await trio.lowlevel.checkpoint()
source.release_exit.set()
assert len(child_exit) == 1
assert root_value == [2]
assert source.calls == 2
assert brx._state.receive_exc is None
trio.run(main)
def test_closing_non_owner_preserves_source_wait() -> None:
'''
Closing one subscriber must not wake another receiver's peers.
`BroadcastReceiver.aclose()` previously set the one shared
`BroadcastState.recv_ready` event even when a different receiver
owned the source read. Waiting peers then repeatedly awaited an
already-set event until the source produced another value,
creating a runnable hot loop on idle streams.
Block one child in the source receive, then place both the root
and a closing child behind its event. Close only that waiting
child and prove it gets `ClosedResourceError` without setting the
shared event. Both remaining receivers must still get the same
value after the source is released.
'''
async def main() -> None:
tx, rx = trio.open_memory_channel(1)
brx = broadcast_receiver(rx, 3)
owner_value: list[int] = []
root_value: list[int] = []
closing_closed = trio.Event()
async with (
brx.subscribe() as owner,
brx.subscribe() as closing,
):
async def receive_owner() -> None:
owner_value.append(await owner.receive())
async def receive_root() -> None:
root_value.append(await brx.receive())
async def receive_closing() -> None:
with pytest.raises(trio.ClosedResourceError):
await closing.receive()
closing_closed.set()
with trio.fail_after(1):
async with trio.open_nursery() as nursery:
nursery.start_soon(receive_owner)
while brx._state.recv_ready is None:
await trio.lowlevel.checkpoint()
nursery.start_soon(receive_root)
nursery.start_soon(receive_closing)
_, event = brx._state.recv_ready
while event.statistics().tasks_waiting < 2:
await trio.lowlevel.checkpoint()
await closing.aclose()
await closing_closed.wait()
assert not event.is_set()
await tx.send(1)
assert owner_value == [1]
assert root_value == [1]
trio.run(main)
@pytest.mark.parametrize(
'first_outcome',
[
1,
RuntimeError('discarded source error'),
trio.EndOfChannel(),
],
)
def test_closing_source_owner_hands_read_to_peer(
first_outcome: int|Exception,
) -> None:
'''
Closing the source-read owner must transfer ownership to a peer.
Merely suppressing the old shared-event wake would leave peers
blocked behind an externally closed receiver that still owned an
idle source read. Script a first receive which blocks until its
private scope is cancelled and a second which returns immediately.
Close that owner only after the root is waiting behind it. Cover
a shielded value, ordinary error and EOC from the cancelled source
read. The owner must always get `ClosedResourceError`, while the
awakened root takes the second source read without publishing the
discarded source outcome.
'''
class HandoffReceiver:
'''
Block the first source read and satisfy the second.
'''
def __init__(self) -> None:
self.calls: int = 0
self.first_started = trio.Event()
self.release_first = trio.Event()
async def receive(self) -> int:
'''
Drive one cancelled read followed by one value.
'''
self.calls += 1
if self.calls == 1:
self.first_started.set()
with trio.CancelScope(shield=True):
await self.release_first.wait()
if isinstance(first_outcome, BaseException):
raise first_outcome
return first_outcome
return 2
async def main() -> None:
source = HandoffReceiver()
brx = broadcast_receiver(source, 3)
owner_closed = trio.Event()
root_value: list[int] = []
async with brx.subscribe() as owner:
async def receive_owner() -> None:
with pytest.raises(trio.ClosedResourceError):
await owner.receive()
owner_closed.set()
async def receive_root() -> None:
root_value.append(await brx.receive())
with trio.fail_after(1):
async with trio.open_nursery() as nursery:
nursery.start_soon(receive_owner)
await source.first_started.wait()
nursery.start_soon(receive_root)
_, event = brx._state.recv_ready
while not event.statistics().tasks_waiting:
await trio.lowlevel.checkpoint()
await owner.aclose()
source.release_first.set()
await owner_closed.wait()
assert source.calls == 2
assert root_value == [2]
assert brx._state.receive_exc is None
assert not brx._state.eoc
trio.run(main)
def test_ensure_slow_consumers_lag_out( def test_ensure_slow_consumers_lag_out(
reg_addr, reg_addr,
start_method, start_method,
@ -448,6 +1247,7 @@ def test_first_recver_is_cancelled():
async with brx.subscribe() as bc: async with brx.subscribe() as bc:
async for value in bc: async for value in bc:
print(value) print(value)
assert cs.cancelled_caught
async def cancel_and_send(): async def cancel_and_send():
await trio.sleep(0.2) await trio.sleep(0.2)
@ -519,3 +1319,74 @@ def test_no_raise_on_lag():
with pytest.raises(KeyboardInterrupt): with pytest.raises(KeyboardInterrupt):
trio.run(main) trio.run(main)
@pytest.mark.parametrize(
('subscribe', 'chan_attr'),
[
(tractor.MsgStream.subscribe, '_rx_chan'),
(LinkedTaskChannel.subscribe, '_from_aio'),
],
ids=['msg-stream', 'linked-task-channel'],
)
def test_stream_subscribe_forwards_lag_policy(
subscribe,
chan_attr: str,
) -> None:
'''
Stream wrappers must expose per-subscriber lag policy.
`MsgStream.subscribe()` and `LinkedTaskChannel.subscribe()`
previously omitted `BroadcastReceiver.raise_on_lag`, forcing
downstream users to mutate a private receiver attribute. Invoke
each public wrapper against a minimal receive-compatible handle.
Prove the first non-raising subscription configures both the
irreversible root broadcaster and its child, while a later strict
child selects its own policy without changing that root.
'''
class StreamHandle:
'''
Provide the wrapper fields needed for local fan-out.
'''
def __init__(self) -> None:
self._broadcaster = None
setattr(
self,
chan_attr,
SimpleNamespace(
_state=SimpleNamespace(max_buffer_size=1),
),
)
async def receive(self):
'''
Block if a regression unexpectedly enters source receive.
'''
await trio.sleep_forever()
async def send(self, value) -> None:
'''
Satisfy `MsgStream` duplex-handle patching.
'''
async def main() -> None:
stream = StreamHandle()
async with subscribe(
stream,
raise_on_lag=False,
) as first:
assert not stream._broadcaster._raise_on_lag
assert not first._raise_on_lag
async with subscribe(
stream,
raise_on_lag=True,
) as second:
assert not stream._broadcaster._raise_on_lag
assert second._raise_on_lag
trio.run(main)

View File

@ -103,6 +103,12 @@ class MsgStream(trio.abc.Channel):
self._eoc: bool|trio.EndOfChannel = False self._eoc: bool|trio.EndOfChannel = False
self._closed: bool|trio.ClosedResourceError = 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 @property
def ctx(self) -> Context: def ctx(self) -> Context:
''' '''
@ -256,7 +262,16 @@ class MsgStream(trio.abc.Channel):
# when the send is closed we assume the stream has # when the send is closed we assume the stream has
# terminated and signal this local iterator to stop # terminated and signal this local iterator to stop
#
# 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() drained: list[Exception|dict] = await self.aclose()
finally:
self._eoc_close_task = None
if drained: if drained:
# ^^^^^^^^TODO? pass these to the `._ctx._drained_msgs: # ^^^^^^^^TODO? pass these to the `._ctx._drained_msgs:
# deque` and then iterate them as part of any # deque` and then iterate them as part of any
@ -335,6 +350,20 @@ class MsgStream(trio.abc.Channel):
# `.__aexit__()` as well!!! # `.__aexit__()` as well!!!
# => SO ENSURE WE CATCH ALL TERMINATION STATES in this # => SO ENSURE WE CATCH ALL TERMINATION STATES in this
# block including the EoC.. # 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: if self.closed:
# this stream has already been closed so silently succeed as # this stream has already been closed so silently succeed as
# per ``trio.AsyncResource`` semantics. # per ``trio.AsyncResource`` semantics.
@ -512,6 +541,7 @@ class MsgStream(trio.abc.Channel):
@acm @acm
async def subscribe( async def subscribe(
self, self,
raise_on_lag: bool = True,
) -> AsyncIterator[BroadcastReceiver]: ) -> AsyncIterator[BroadcastReceiver]:
''' '''
@ -526,6 +556,11 @@ class MsgStream(trio.abc.Channel):
value from the far end via the internally created broudcast value from the far end via the internally created broudcast
receiver wrapper. receiver wrapper.
``raise_on_lag=False`` makes this subscription warn and resume
at the oldest retained value after an overrun. The first call
also sets that policy for this stream's root receive handle;
later child subscriptions choose their policy independently.
''' '''
# NOTE: This operation is indempotent and non-reversible, so be # NOTE: This operation is indempotent and non-reversible, so be
# sure you can deal with any (theoretical) overhead of the the # sure you can deal with any (theoretical) overhead of the the
@ -541,6 +576,7 @@ class MsgStream(trio.abc.Channel):
# TODO: can remove this kwarg right since # TODO: can remove this kwarg right since
# by default behaviour is to do this anyway? # by default behaviour is to do this anyway?
receive_afunc=self.receive, receive_afunc=self.receive,
raise_on_lag=raise_on_lag,
) )
# NOTE: we override the original stream instance's receive # NOTE: we override the original stream instance's receive
@ -552,7 +588,9 @@ class MsgStream(trio.abc.Channel):
# seems there's no graceful way to type this with ``mypy``? # seems there's no graceful way to type this with ``mypy``?
# https://github.com/python/mypy/issues/708 # https://github.com/python/mypy/issues/708
async with self._broadcaster.subscribe() as bstream: async with self._broadcaster.subscribe(
raise_on_lag=raise_on_lag,
) as bstream:
assert bstream.key != self._broadcaster.key assert bstream.key != self._broadcaster.key
assert bstream._recv == self._broadcaster._recv assert bstream._recv == self._broadcaster._recv

View File

@ -213,6 +213,14 @@ class LinkedTaskChannel(
_broadcaster: BroadcastReceiver|None = None _broadcaster: BroadcastReceiver|None = None
async def aclose(self) -> 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() await self._from_aio.aclose()
# ?TODO? async version of this? # ?TODO? async version of this?
@ -324,6 +332,7 @@ class LinkedTaskChannel(
@acm @acm
async def subscribe( async def subscribe(
self, self,
raise_on_lag: bool = True,
) -> AsyncIterator[BroadcastReceiver]: ) -> AsyncIterator[BroadcastReceiver]:
''' '''
@ -335,6 +344,11 @@ class LinkedTaskChannel(
See ``tractor._streaming.MsgStream.subscribe()`` for further See ``tractor._streaming.MsgStream.subscribe()`` for further
similar details. similar details.
``raise_on_lag=False`` makes this subscription warn and resume
at the oldest retained value after an overrun. The first call
also sets that policy for this channel's root receive handle;
later child subscriptions choose their policy independently.
''' '''
if self._broadcaster is None: if self._broadcaster is None:
@ -343,11 +357,14 @@ class LinkedTaskChannel(
# use memory channel size by default # use memory channel size by default
self._from_aio._state.max_buffer_size, # type: ignore self._from_aio._state.max_buffer_size, # type: ignore
receive_afunc=self.receive, receive_afunc=self.receive,
raise_on_lag=raise_on_lag,
) )
self.receive = bcast.receive # type: ignore self.receive = bcast.receive # type: ignore
async with self._broadcaster.subscribe() as bstream: async with self._broadcaster.subscribe(
raise_on_lag=raise_on_lag,
) as bstream:
assert bstream.key != self._broadcaster.key assert bstream.key != self._broadcaster.key
assert bstream._recv == self._broadcaster._recv assert bstream._recv == self._broadcaster._recv
yield bstream yield bstream

View File

@ -26,6 +26,7 @@ from ._mngrs import (
from ._broadcast import ( from ._broadcast import (
AsyncReceiver as AsyncReceiver, AsyncReceiver as AsyncReceiver,
broadcast_receiver as broadcast_receiver, broadcast_receiver as broadcast_receiver,
BroadcastReceiveError as BroadcastReceiveError,
BroadcastReceiver as BroadcastReceiver, BroadcastReceiver as BroadcastReceiver,
Lagged as Lagged, Lagged as Lagged,
) )

View File

@ -100,6 +100,20 @@ class Lagged(trio.TooSlowError):
''' '''
class BroadcastReceiveError(Exception):
'''
A shared underlying receiver failed in another subscriber task.
'''
class _BroadcastReceiverClosed(Exception):
'''
An active receiver was closed while owning the source read.
'''
class BroadcastState(Struct): class BroadcastState(Struct):
''' '''
Common state to all receivers of a broadcast. Common state to all receivers of a broadcast.
@ -115,6 +129,7 @@ class BroadcastState(Struct):
# broadcast event to wake up all sleeping consumer tasks # broadcast event to wake up all sleeping consumer tasks
# on a newly produced value from the sender. # on a newly produced value from the sender.
recv_ready: tuple[int, trio.Event]|None = None recv_ready: tuple[int, trio.Event]|None = None
recv_scope: trio.CancelScope|None = None
# if a ``trio.EndOfChannel`` is received on any # if a ``trio.EndOfChannel`` is received on any
# consumer all consumers should be placed in this state # consumer all consumers should be placed in this state
@ -122,7 +137,13 @@ class BroadcastState(Struct):
# For now, this is solely for testing/debugging purposes. # For now, this is solely for testing/debugging purposes.
eoc: bool = False eoc: bool = False
# If the broadcaster was cancelled, we might as well track it # Any non-EOC failure from the shared underlying receiver is
# terminal for every subscriber. Retained values remain readable
# before this failure is replayed at each receiver's boundary.
receive_exc: Exception | None = None
# Retain the latest interrupted source-reader task until its
# receiver next makes progress or closes.
cancelled: dict[int, Task] = {} cancelled: dict[int, Task] = {}
def statistics(self) -> dict[str, Any]: def statistics(self) -> dict[str, Any]:
@ -142,13 +163,20 @@ class BroadcastState(Struct):
qlens: dict[int, int] = {} qlens: dict[int, int] = {}
for tid, sz in subs.items(): for tid, sz in subs.items():
qlens[tid] = sz if sz != -1 else 0 qlens[tid] = min(
sz + 1,
len(self.queue),
)
return { return {
'open_consumers': len(subs), 'open_consumers': len(subs),
'queued_len_by_task': qlens, 'queued_len_by_task': qlens,
'max_buffer_size': self.maxlen, 'max_buffer_size': self.maxlen,
'tasks_waiting': ev.statistics().tasks_waiting if ev else 0, 'tasks_waiting': (
ev.statistics().tasks_waiting
if ev is not None
else 0
),
'tasks_cancelled': self.cancelled, 'tasks_cancelled': self.cancelled,
'next_value_receiver_id': key, 'next_value_receiver_id': key,
} }
@ -190,6 +218,7 @@ class BroadcastReceiver(ReceiveChannel):
self._recv = receive_afunc or rx_chan.receive self._recv = receive_afunc or rx_chan.receive
self._closed: bool = False self._closed: bool = False
self._raise_on_lag = raise_on_lag self._raise_on_lag = raise_on_lag
self._wait_scope: trio.CancelScope|None = None
def receive_nowait( def receive_nowait(
self, self,
@ -237,7 +266,10 @@ class BroadcastReceiver(ReceiveChannel):
# https://docs.rs/tokio/1.11.0/tokio/sync/broadcast/index.html#lagging # https://docs.rs/tokio/1.11.0/tokio/sync/broadcast/index.html#lagging
mxln = state.maxlen mxln = state.maxlen
lost = seq - mxln # `seq == mxln` is already one past the final
# valid deque index, so include that first
# displaced value in the loss count.
lost = seq - mxln + 1
# decrement to the last value and expect # decrement to the last value and expect
# consumer to either handle the ``Lagged`` and come back # consumer to either handle the ``Lagged`` and come back
@ -255,8 +287,21 @@ class BroadcastReceiver(ReceiveChannel):
return self.receive_nowait(_key, _state) return self.receive_nowait(_key, _state)
state.subs[key] -= 1 state.subs[key] -= 1
state.cancelled.pop(key, None)
return value return value
receive_exc = state.receive_exc
if receive_exc is not None:
# Re-raising one shared exception mutates its traceback on
# every delivery. Give each receiver a stable wrapper while
# retaining the original failure as its cause.
raise BroadcastReceiveError(
'Shared broadcast receiver failed'
) from receive_exc
if state.eoc:
raise trio.EndOfChannel
raise trio.WouldBlock raise trio.WouldBlock
async def _receive_from_underlying( async def _receive_from_underlying(
@ -270,14 +315,34 @@ class BroadcastReceiver(ReceiveChannel):
raise trio.ClosedResourceError raise trio.ClosedResourceError
event = trio.Event() event = trio.Event()
recv_scope = trio.CancelScope()
assert state.recv_ready is None assert state.recv_ready is None
assert state.recv_scope is None
state.recv_ready = key, event state.recv_ready = key, event
state.recv_scope = recv_scope
try: try:
# if we're cancelled here it should be # if we're cancelled here it should be
# fine to bail without affecting any other consumers # fine to bail without affecting any other consumers
# right? # right?
receive_exc: BaseException|None = None
with recv_scope:
try:
value = await self._recv() value = await self._recv()
except BaseException as exc:
receive_exc = exc
# Only this receiver's `aclose()` cancels its private
# source-read scope, and it marks the receiver closed
# first without a checkpoint. Outer task cancellation does
# not set `recv_scope.cancel_called`; it remains a real
# `trio.Cancelled` and follows the handler below.
if recv_scope.cancel_called:
assert self._closed
if self._closed:
raise _BroadcastReceiverClosed
if receive_exc is not None:
raise receive_exc
# items with lower indices are "newer" # items with lower indices are "newer"
# NOTE: ``collections.deque`` implicitly takes care of # NOTE: ``collections.deque`` implicitly takes care of
@ -303,6 +368,8 @@ class BroadcastReceiver(ReceiveChannel):
): ):
state.subs[sub_key] += 1 state.subs[sub_key] += 1
state.cancelled.pop(key, None)
# NOTE: this should ONLY be set if the above task was *NOT* # NOTE: this should ONLY be set if the above task was *NOT*
# cancelled on the `._recv()` call. # cancelled on the `._recv()` call.
event.set() event.set()
@ -312,11 +379,20 @@ class BroadcastReceiver(ReceiveChannel):
# if any one consumer gets an EOC from the underlying # if any one consumer gets an EOC from the underlying
# receiver we need to unblock and send that signal to # receiver we need to unblock and send that signal to
# all other consumers. # all other consumers.
state.cancelled.clear()
self._state.eoc = True self._state.eoc = True
if event.statistics().tasks_waiting: if event.statistics().tasks_waiting:
event.set() event.set()
raise raise
except _BroadcastReceiverClosed:
# `aclose()` cancelled this receiver's source-read scope.
# Wake peers so one of them can take ownership after this
# task clears `recv_ready` in `finally`.
if event.statistics().tasks_waiting:
event.set()
raise trio.ClosedResourceError
except ( except (
trio.Cancelled, trio.Cancelled,
): ):
@ -329,12 +405,33 @@ class BroadcastReceiver(ReceiveChannel):
event.set() event.set()
raise raise
except Exception as receive_exc:
# The underlying receiver is shared by every subscriber,
# 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()
raise
except BaseException:
# 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
finally: finally:
# Reset receiver waiter task event for next blocking condition. # Reset receiver waiter task event for next blocking condition.
# this MUST be reset even if the above ``.recv()`` call # this MUST be reset even if the above ``.recv()`` call
# was cancelled to avoid the next consumer from blocking on # was cancelled to avoid the next consumer from blocking on
# an event that won't be set! # an event that won't be set!
state.recv_ready = None state.recv_ready = None
state.recv_scope = None
async def receive(self) -> ReceiveType: async def receive(self) -> ReceiveType:
key = self.key key = self.key
@ -362,7 +459,23 @@ class BroadcastReceiver(ReceiveChannel):
# seq = state.subs[key] # seq = state.subs[key]
# assert seq == -1 # sanity # assert seq == -1 # sanity
_, ev = state.recv_ready _, ev = state.recv_ready
wait_scope = trio.CancelScope()
self._wait_scope = wait_scope
try:
with wait_scope:
await ev.wait() await ev.wait()
# As with `recv_scope`, only this receiver's
# `aclose()` cancels its private peer-wait scope
# after marking the receiver closed. Outer task
# cancellation remains `trio.Cancelled`.
if wait_scope.cancel_called:
assert self._closed
if self._closed:
raise trio.ClosedResourceError
finally:
self._wait_scope = None
try: try:
return self.receive_nowait( return self.receive_nowait(
_key=key, _key=key,
@ -440,18 +553,35 @@ class BroadcastReceiver(ReceiveChannel):
if self._closed: if self._closed:
return return
# if there are sleeping consumers wake
# them on closure.
rr = self._state.recv_ready
if rr:
_, event = rr
event.set()
# XXX: leaving it like this consumers can still get values # XXX: leaving it like this consumers can still get values
# up to the last received that still reside in the queue. # up to the last received that still reside in the queue.
self._state.subs.pop(self.key) state = self._state
state.subs.pop(self.key)
state.cancelled.pop(self.key, None)
self._closed = True self._closed = True
# A non-owner close must not wake peers waiting behind some
# other receiver's source read. If this receiver owns that
# read, cancel only its private scope; the owner task wakes
# peers after cancellation is delivered and state is ready
# for a clean ownership handoff.
rr = state.recv_ready
if (
rr is not None
# `recv_ready[0]` identifies the receiver which currently
# owns the one shared source read. Only that receiver may
# cancel `BroadcastState.recv_scope`; closing any other
# subscriber must not disturb the owner or its peer tasks.
and
rr[0] == self.key
):
recv_scope = state.recv_scope
assert recv_scope is not None
recv_scope.cancel()
elif (wait_scope := self._wait_scope) is not None:
wait_scope.cancel()
def broadcast_receiver( def broadcast_receiver(
@ -462,6 +592,11 @@ def broadcast_receiver(
) -> BroadcastReceiver: ) -> BroadcastReceiver:
if max_buffer_size < 1:
raise ValueError(
'`max_buffer_size` must be greater than zero'
)
return BroadcastReceiver( return BroadcastReceiver(
recv_chan, recv_chan,
state=BroadcastState( state=BroadcastState(