Compare commits
6 Commits
6429e6a515
...
3798aa9c9a
| Author | SHA1 | Date |
|---|---|---|
|
|
3798aa9c9a | |
|
|
b881512d65 | |
|
|
933374de2c | |
|
|
08dfd834a6 | |
|
|
8ab6281a68 | |
|
|
8db13375fc |
|
|
@ -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.
|
||||||
|
|
@ -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.
|
||||||
|
|
@ -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.
|
||||||
|
|
@ -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.
|
||||||
|
|
@ -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.
|
||||||
|
|
@ -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.
|
||||||
|
|
@ -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.
|
||||||
|
|
@ -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.
|
||||||
|
|
@ -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.
|
||||||
|
|
@ -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.
|
||||||
|
|
@ -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.
|
||||||
|
|
@ -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.
|
||||||
|
|
@ -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
|
||||||
----------------------
|
----------------------
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -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
|
||||||
|
|
|
||||||
|
|
@ -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.
|
||||||
|
|
|
||||||
|
|
@ -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)
|
||||||
|
|
|
||||||
|
|
@ -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
|
||||||
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:
|
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
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -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
|
||||||
|
|
|
||||||
|
|
@ -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,
|
||||||
)
|
)
|
||||||
|
|
|
||||||
|
|
@ -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?
|
||||||
value = await self._recv()
|
receive_exc: BaseException|None = None
|
||||||
|
with recv_scope:
|
||||||
|
try:
|
||||||
|
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
|
||||||
await ev.wait()
|
wait_scope = trio.CancelScope()
|
||||||
|
self._wait_scope = wait_scope
|
||||||
|
try:
|
||||||
|
with wait_scope:
|
||||||
|
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(
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue