Compare commits

...

23 Commits

Author SHA1 Message Date
Gud Boi 28e6269a35 Type the remaining niche examples
Finish the examples-typing sweep with the last non-docs-visible
scripts: `-> None` on the two `trio/` behavior-demo mains (plus a
`trio.TaskStatus` on `hold_lock_forever`) and nursery/portal typing
on `integration/mpi4py/inherit_parent_main.py`.

Leaves `concurrent_futures_primes` (a verbatim stdlib baseline) and
`integration/open_context_and_sleep` (its tractor nursery is
commented out) as-is, and the paren-group `trio.open_nursery()`
bindings unannotated (no clean spot for a preceding annotation).
Completes the examples-typing bullet in #472.

(this patch was generated in some part by [`claude-code`][claude-code-gh])
[claude-code-gh]: https://github.com/anthropics/claude-code
2026-08-12 20:09:46 -04:00
Gud Boi 67fa97dbaf Type the docs-visible `examples/` scripts
Type the runtime objects (`ActorNursery`, `Portal`, `Context`,
`trio.Nursery`) + fn signatures across the 16 highest-visibility,
`literalinclude`-d `examples/` scripts, matching the front-page
`we_are_processes.py` style — so the rendered guides show typed
usage throughout, not just on the landing snippet.

Spans the 3 quickstart-backing scripts + `single_func`,
`remote_error_propagation`, `multiple_streams_one_portal`,
`quick_cluster`, `service_discovery`, `service_daemon_discovery`,
`asynchronous_generators`, `nested_actor_tree`,
`concurrent_actors_primes`, `streaming_broadcast_fanout`,
`rpc_bidir_streaming`, `infected_asyncio_echo_server`,
`typed_payloads`.

Annotation-only (no renames/logic changes); each runs green and the
docs build stays warning-free. Part of the examples-typing bullet
in #472.

(this patch was generated in some part by [`claude-code`][claude-code-gh])
[claude-code-gh]: https://github.com/anthropics/claude-code
2026-08-12 20:09:46 -04:00
Gud Boi 545142933e Add basic typing to the `debugging/` examples
Sweep the `examples/debugging/` set for basic typing: add `-> None`
to all 16 bare `async def main()`s and annotate the clean
single-line `open_nursery()` bindings as `tractor.ActorNursery`.

Kept to the unambiguous, runtime-safe cases (these breakpoint/crash
demos can't be run headless); the heterogeneous
multi-line/paren-group nursery bindings + `current_actor()` returns
are left for a later pass. Continues the examples-typing bullet in
#472.

(this patch was generated in some part by [`claude-code`][claude-code-gh])
[claude-code-gh]: https://github.com/anthropics/claude-code
2026-08-12 20:09:46 -04:00
Gud Boi 5067917b9f Add a dedicated-registrar example + discovery guide
Add a runnable `examples/dedicated_registrar.py` + a "A dedicated
registrar" subsection in `guide/discovery.rst` demoing the
registrar decoupled from any app tree's root: boot a bare
`tractor.run_daemon([], registry_addrs=[...])` as its own process
(a root actor that does nothing but hold the registry), point the
app tree at the same `registry_addrs`, and discover a service
*through* that external registrar.

This is the buildable-today form of the #472
"Registrar-as-subsystem (not the root actor)" bullet. Two
constraints are called out inline as follow-ups:
`enable_transports` is single-proto per runtime (no multi-backend
registrar yet), and a registrar can only be a root (no `actor_cls`
hook on `start_actor()` to spawn one as a subactor).

(this patch was generated in some part by [`claude-code`][claude-code-gh])
[claude-code-gh]: https://github.com/anthropics/claude-code
2026-08-12 20:09:46 -04:00
Gud Boi 346878219c Expand caps-based-msging docs w/ #365, #376 + tests
The "Toward capability-based msging" section only pointed at the
`#196`/`#36` epics. Fold in the concrete recent state,

- `#365` as the most recent step: driving the whole `pld_spec` off
  plain type-annotations (e.g. annotating a context's
  `open_stream()` with `msgspec.Struct` subtypes) rather than
  explicit `pld_spec=` kwargs.
- clarify that the decorator-level `@tractor.context(pld_spec=...)`
  is already the higher-level path (vs the lower-level
  `tractor.msg._ops.limit_plds()` escape hatch), pointing at
  `tests/msg/test_pldrx_limiting.py` + `test_ext_types_msgspec.py`
  which exercise both.
- `#376` (from @guilledk, `auto_codecs` branch) as the drafted
  public factory API for the `enc_hook`/`dec_hook` pair (today only
  reachable via `tractor.msg._ops`).

Addresses the caps-based-msging bullet in #472.

(this patch was generated in some part by [`claude-code`][claude-code-gh])
[claude-code-gh]: https://github.com/anthropics/claude-code
2026-08-12 20:09:46 -04:00
Bd 3ad7e7e5dc
Merge pull request #459 from goodboy/dependabot/uv/idna-3.15
Bump idna from 3.10 to 3.18
2026-08-12 19:55:40 -04:00
dependabot[bot] 4b4cc76263 Bump idna from 3.10 to 3.18
Bumps [idna](https://github.com/kjd/idna) from 3.10 to 3.18.
- [Release notes](https://github.com/kjd/idna/releases)
- [Changelog](https://github.com/kjd/idna/blob/master/HISTORY.md)
- [Commits](https://github.com/kjd/idna/compare/v3.10...v3.18)

---
updated-dependencies:
- dependency-name: idna
  dependency-version: '3.18'
  dependency-type: indirect
...

Signed-off-by: dependabot[bot] <support@github.com>
2026-08-12 19:35:17 -04:00
Bd 99f9beccb2
Merge pull request #487 from goodboy/dependabot/uv/setuptools-83.0.0
Bump setuptools from 82.0.1 to 83.0.0
2026-08-12 17:13:22 -04:00
dependabot[bot] 5c4d42c7a7 Bump setuptools from 82.0.1 to 83.0.0
Bumps [setuptools](https://github.com/pypa/setuptools) from 82.0.1 to 83.0.0.
- [Release notes](https://github.com/pypa/setuptools/releases)
- [Changelog](https://github.com/pypa/setuptools/blob/main/NEWS.rst)
- [Commits](https://github.com/pypa/setuptools/compare/v82.0.1...v83.0.0)

---
updated-dependencies:
- dependency-name: setuptools
  dependency-version: 83.0.0
  dependency-type: indirect
...

Signed-off-by: dependabot[bot] <support@github.com>
2026-08-12 16:48:08 -04:00
Bd d887603a1b
Merge pull request #479 from goodboy/wkt/start_or_cancel_tests_474
Add `.trionics.start_or_cancel()` test suite
2026-08-12 16:43:33 -04:00
Gud Boi ae67e2f429 Assert cancellation at the startup boundary
Replace the pre-start wall-clock sleep with an indefinite
checkpoint so cancellation ordering cannot race a timer in slow CI.

Record that `Cancelled` escapes each `start_or_cancel()` await
before the enclosing nursery or cancel scope handles it.

Review: PR #479 (goodboy)
https://github.com/goodboy/tractor/pull/479

(this patch was generated in some part by `opencode` using
`gpt-5.6-sol` (`openai`))
2026-08-12 15:45:30 -04:00
Gud Boi 1c7d0c7f3e Match Trio's startup error exactly
Compare the complete canonical `Nursery.start()` protocol error
before re-surfacing ambient cancellation. Preserve child-owned
`RuntimeError` objects whose messages only resemble Trio's wording.

Cover the colliding prefix and assert the original error remains the
exception group's sole leaf.

Review: PR #479 (goodboy)
https://github.com/goodboy/tractor/pull/479

(this patch was generated in some part by `opencode` using
`gpt-5.6-sol` (`openai`))
2026-08-12 15:45:23 -04:00
Gud Boi 203a0f7e1f Add `.trionics.start_or_cancel()` test suite
Resolve #474 with a new `tests/trionics/test_taskc.py` (9
tests) covering the `trio.Nursery.start()` wrapper landed in
PR #464, incl. the `modden.runtime.progman.open_wks()` use
case dug out as a minimal repro.

Deats,
- the lossy `RuntimeError('child exited without calling
  task_status.started()')` only fires when the child absorbs
  its ambient cancel pre-`.started()` (graceful-teardown
  pattern); a well-behaved child surfaces `Cancelled` direct
  from `.start()` on `trio` 0.29 - verified empirically 1st.
- `test_sibling_err_not_masked_by_startup_rte`: the `modden`
  case; ONLY the root-cause sibling `ValueError` escapes the
  nursery with the wrapper vs. bare-`.start()`'s lossy
  riding-along startup-RTE noise.
- `test_pure_oob_cancel_not_morphed_to_rte`: plain ancestor
  `cs.cancel()` exits clean vs. bare's eg-wrapped RTE.
- `test_genuine_startup_rte_still_raised`: sans cancellation
  the protocol-bug RTE re-raises same as bare.
- `test_childs_own_rte_never_demoted_to_cancel`: exact-msg +
  `isinstance`-guard regression cover; a child's own
  `RuntimeError('never got started!')`/`RuntimeError(1234)`
  never demotes to a `Cancelled`.
- `test_started_value_and_args_passthru`: positional args,
  `name=` and `.started()`-value forwarding.

Also,
- `use_start_or_cancel=False` params pin upstream `trio`'s
  current lossy behaviour as wart-documentation: a break on
  a `trio` upgrade likely means new upstream porcelain and
  the wrapper deserves a re-audit.
- verified 0 flakes over 50 hammer runs + 2 impl mutations
  each caught by exactly the targeted tests.

Prompt-IO: ai/prompt-io/claude/20260702T161624Z_65bf9df5_prompt_io.md
(this patch was generated in some part by [`claude-code`][claude-code-gh])
[claude-code-gh]: https://github.com/anthropics/claude-code
2026-08-12 14:15:31 -04:00
Bd 92c737ad83
Merge pull request #458 from mahmoudhas/fix/guard-hot-path-log-rendering
Guard hot-path log calls to avoid payload rendering when disabled
2026-08-12 14:08:47 -04:00
Gud Boi 935c8cf656 Guard receive-path transport rendering
Skip raw packet, decoded message, peer, and channel formatting when
transport logging is disabled. Keep message processing and wire
reads outside the guards so logging controls never affect IPC flow.

Also narrow the remaining pretty-struct TODO to require a
non-raising formatter with native-repr fallback.

(this patch was generated in some part by `opencode` using
`gpt-5.6-sol` (`openai`))
2026-08-12 13:48:40 -04:00
Gud Boi 84ec895150 Drop resolved channel log-guard TODO
Remove the stale `Channel.from_addr()` design note now that its
`at_least_level()` guard avoids inactive pretty-rendering work.

(this patch was generated in some part by `opencode` using
`gpt-5.6-sol` (`openai`))
2026-08-12 12:51:43 -04:00
Gud Boi 67280c2898 Honor log disable controls in hot-path guards
Use `Logger.isEnabledFor()` in `at_least_level()` so logger-local
and global disable controls short-circuit payload rendering.

Add `Channel.send()` coverage for effective-level, per-logger, and
global suppression while ensuring transport remains unchanged.

Caught-during: review remediation
Found-via: `/run-tests` test_log_guard_skips_payload_formatting

Review: PR #458 (goodboy)
https://github.com/goodboy/tractor/pull/458#issuecomment-5258207470

(this patch was generated in some part by `opencode` using
`gpt-5.6-sol` (`openai`))
2026-08-11 21:56:27 -04:00
root 0e11ff7e9d Guard hot-path log calls to avoid payload rendering when disabled
Wrap `log.transport()` in `Channel.send()` and `log.runtime()` in
`PldRx.decode_pld()` with `log.at_least_level()` checks so that
expensive `pformat(payload)` / `repr(msg)` / `repr(pld)` calls are
skipped entirely when the respective log level is not active.

Previously the f-string arguments were eagerly evaluated before being
passed to the log method, even though `StackLevelAdapter.log()` would
then discard the message internally via its own `isEnabledFor()` check.
On high-frequency IPC paths this caused `pformat` to dominate CPU
usage (~60-70 %) as reported in #455.

This also restores the full diagnostic output (msg type, decoded
payload) that was temporarily commented out in 0373164 as a stopgap.

Resolves #455
2026-08-11 20:41:20 -04:00
mahmoud 148a098ca6 avoid format on the hot send path 2026-08-11 20:41:20 -04:00
Bd 83b3488455
Merge pull request #488 from goodboy/wkt/moc_teardown_completion
Wait for ctx exit in `maybe_open_context()`
2026-08-11 12:16:15 -04:00
Gud Boi daa661aba3 Clarify `maybe_open_context()` teardown notes
Drop the stale sentinel experiment and fix the cancellation-path
comment. Document that cached regular `__aexit__()` failures are
always re-raised at the final consumer boundary.

Review: PR #488 (goodboy)
https://github.com/goodboy/tractor/pull/488

(this patch was generated in some part by `opencode` using
`gpt-5.6-sol` (`openai`))
2026-08-11 11:38:43 -04:00
Gud Boi 55ec3dbf51 Drop unused `_Cache` teardown bindings
Remove the unused `value` assignment after cache eviction and skip
unpacking stale resource state before raising its invariant error.

Review: PR #488 (Copilot)
https://github.com/goodboy/tractor/pull/488#pullrequestreview-4850557500

(this patch was generated in some part by `opencode` using
`gpt-5.6-sol` (`openai`))
2026-08-07 22:24:27 -04:00
Gud Boi a753625fc6 Wait for ctx exit in `maybe_open_context()`
Block the final user on its cached resource's `__aexit__()` and
raise regular cleanup errors at that user's ctx boundary.

Deats,
- serialize user registration and teardown under each cache-key lock
- keep queued entrants on the same lock through resource replacement
- preserve `Cancelled`, `KeyboardInterrupt`, and `SystemExit` flow
- cover successful, failing, cancelled, and re-entry teardown paths

Prompt-IO: ai/prompt-io/opencode/20260804T030309Z_65bf9df5_prompt_io.md

(this patch was generated in some part by `opencode` using `gpt-5.6-sol` (`openai`))
2026-08-03 23:45:42 -04:00
56 changed files with 1420 additions and 193 deletions

View File

@ -0,0 +1,80 @@
---
model: claude-fable-5
service: claude
session: f4da6654-63f9-4641-b567-f200da46d697
timestamp: 2026-07-02T16:16:24Z
git_ref: 65bf9df5
scope: tests
substantive: true
raw_file: 20260702T161624Z_65bf9df5_prompt_io.raw.md
---
## Prompt
> attempt to resolve
> https://github.com/goodboy/tractor/issues/474
> do it with /open-wkt.
Issue #474: "We need test suite for
`.trionics.start_or_cancel()`" — dig out the use case
from `modden` as a unit test (and provide motivation
for upstream `trio` to possibly provide better
porcelain around this case). The wrapper landed via
PR #464 with no tests.
## Response summary
New `tests/trionics/test_taskc.py` (9 tests) covering
`start_or_cancel()`, authored in wkt
`start_or_cancel_tests_474`.
Behaviour was first probed empirically (trio 0.29):
the lossy `RuntimeError('child exited without calling
task_status.started()')` only fires when the child
exits pre-`.started()` WITHOUT propagating the ambient
`Cancelled` — i.e. when the child (or lib code it
calls) absorbs the cancel in a graceful-teardown
pattern; a well-behaved child surfaces `Cancelled`
straight out of `.start()`. The `modden`
`progman.open_wks()` use case was reconstructed from
`modden/runtime/progman.py` accordingly.
Tests (each `use_start_or_cancel` parametrization also
pins upstream trio's current lossy behaviour as
wart-documentation):
- `test_sibling_err_not_masked_by_startup_rte` — the
`modden` case: sibling error OOB-cancels the shared
nursery scope; with the wrapper ONLY the root-cause
`ValueError` escapes; bare `.start()` adds the lossy
RTE alongside.
- `test_pure_oob_cancel_not_morphed_to_rte` — plain
ancestor `cs.cancel()`: wrapper → clean exit; bare
→ eg-wrapped RTE.
- `test_genuine_startup_rte_still_raised` — no
cancellation → protocol-bug RTE re-raised same as
bare.
- `test_childs_own_rte_never_demoted_to_cancel` — a
child's own `RuntimeError('never got started!')` /
`RuntimeError(1234)` under ambient cancel is never
demoted to `Cancelled` (exact-msg-match + str-guard
regression cover).
- `test_started_value_and_args_passthru` — happy path:
positional args, `name=`, `.started()` value.
Verified: 9/9 pass; 0 flakes across 50 hammer runs;
two impl mutations (checkpoint removed; guard relaxed
to substring match) each caught by exactly the
targeted tests; `tests/trionics/` +
`tests/test_trioisms.py` subset green (23 passed,
5 xfailed); ruff clean; 69-col style.
## Files changed
- `tests/trionics/test_taskc.py` — new
`start_or_cancel()` unit-test suite (gh #474).
## Human edits
Pending review — session paused pre-commit per user
deadline; nothing committed as of this entry.

View File

@ -0,0 +1,107 @@
---
model: claude-fable-5
service: claude
timestamp: 2026-07-02T16:16:24Z
git_ref: 65bf9df5
diff_cmd: git diff main..wkt/start_or_cancel_tests_474
---
# Raw output — gh #474 `start_or_cancel()` test suite
## Generated test code
> `git diff main..wkt/start_or_cancel_tests_474 -- tests/trionics/test_taskc.py`
Prose summary of the generated module
(`tests/trionics/test_taskc.py`):
- module docstring framing the `trio.Nursery.start()`
startup-cancellation wart, the wrapper's repair, and
the intent that `use_start_or_cancel=False` params
double as upstream-trio wart-documentation (break on
a trio upgrade → upstream may have shipped porcelain,
re-audit the wrapper); cites gh #474 / PR #464 and
`modden`'s `progman.open_wks()` as the source use
case.
- shared children: `absorbs_cancel_pre_started()` (the
graceful-teardown cancel-absorber which triggers the
lossy RTE path) + `raise_val_err()` (fast-erroring
sibling).
- `test_sibling_err_not_masked_by_startup_rte`
(parametrized `use_start_or_cancel`): asserts eg
contains exactly one `ValueError` and, wrapper-case,
NO residual RTE (`eg.split(ValueError)` remainder is
`None`); bare-case, the residual RTE carries trio's
exact "child exited without calling" wording.
- `test_pure_oob_cancel_not_morphed_to_rte`
(parametrized): wrapper-case runs clean and asserts
`cs.cancelled_caught`; bare-case asserts the
eg-wrapped RTE.
- `test_genuine_startup_rte_still_raised`
(parametrized): no-cancel protocol bug → RTE with
trio's wording from both call forms.
- `test_childs_own_rte_never_demoted_to_cancel`
(parametrized `rte_arg` in `'never got started!'`,
`1234`): child cancels the ambient scope then raises
its own RTE synchronously (no checkpoint between →
deterministically under-cancellation at catch time);
asserts the RTE survives with `args[0]` intact.
- `test_started_value_and_args_passthru`: `.started()`
value, positional args and the `name=` kwarg (via
`trio.lowlevel.current_task().name`) all forward.
## Non-code output (verbatim highlights)
Behaviour probe (trio 0.29, scratchpad scripts) — the
decision basis for the test shapes:
```
== B-sibling-err use_soc=False
start raised: RuntimeError('child exited without
calling task_status.started()')
top-level: ExceptionGroup([ValueError('sibling blew
up!'), RuntimeError('child exited without calling
task_status.started()')])
== B-cs-cancel use_soc=False
top-level: ExceptionGroup([RuntimeError('child
exited without calling task_status.started()')])
== B-sibling-err use_soc=True
start raised: Cancelled()
top-level: ExceptionGroup([ValueError('sibling blew
up!')])
== B-cs-cancel use_soc=True
start raised: Cancelled()
top-level: clean return
== own-rte-under-cancel (both) -> RTE('never got
started!') propagates unchanged
```
Key finding: with a WELL-BEHAVED (non-absorbing) child
an OOB ancestor cancel surfaces `Cancelled` directly
from `.start()` on trio 0.29 — the lossy RTE requires
the child to absorb its cancel pre-`.started()`, which
is what `modden`'s `open_from_wks` teardown did. Trio's
nursery-exit wait defers cancel delivery to children,
so all tested shapes are deterministic (0 flakes / 50
runs).
Mutation verification:
```
mutation 1 (checkpoint_if_cancelled removed):
FAILED test_sibling_err_not_masked_by_startup_rte[True]
FAILED test_pure_oob_cancel_not_morphed_to_rte[True]
mutation 2 (guard relaxed to 'started' substring,
isinstance dropped):
FAILED test_childs_own_rte_never_demoted_to_cancel[never got started!]
FAILED test_childs_own_rte_never_demoted_to_cancel[1234]
```
Final runs:
```
tests/trionics/test_taskc.py: 9 passed in 0.03s
hammer: 0/50 runs failed
tests/trionics/ + tests/test_trioisms.py:
23 passed, 5 xfailed in 3.02s
```

View File

@ -0,0 +1,39 @@
---
model: openai/gpt-5.6-sol
service: opencode
session: moc-teardown-completion-20260804
timestamp: 2026-08-04T03:03:09Z
git_ref: 65bf9df5
scope: code
substantive: true
raw_file: 20260804T030309Z_65bf9df5_prompt_io.raw.md
---
## Prompt
Patch Tractor's `maybe_open_context()` so the final consumer waits for
resource `__aexit__()` completion and receives cleanup errors. Reuse
`outcome.Outcome` for the exit result; use the smaller mutable
`_CtxExit` holder if that makes the implementation simpler. Add and run
the relevant existing unit tests, but do not commit or push the patch.
After reviewing the result, simplify `_CtxExit` back to an optional
exception because the success outcome carries no useful value.
## Response summary
Added an exception-backed completion handshake between
`_Cache.run_ctx()` and the final `maybe_open_context()` consumer.
Serialized consumer registration and final teardown under the per-key
lock, preserving that lock for queued entrants. Added deterministic
regressions for normal exit, cleanup errors, cancellation interactions,
service-nursery cancellation, and teardown re-entry.
## Files changed
- `tractor/trionics/_mngrs.py` - publish and unwrap cached exit outcomes.
- `tests/test_resource_cache.py` - cover completion and cancellation.
## Human edits
The user directed the final simplification from `outcome.Outcome` to an
optional exception field. The patch remains uncommitted.

View File

@ -0,0 +1,36 @@
---
model: openai/gpt-5.6-sol
service: opencode
timestamp: 2026-08-04T03:03:09Z
git_ref: 65bf9df5
diff_cmd: git diff HEAD~1..HEAD
---
Implemented cached-context exit completion in
`tractor.trionics.maybe_open_context()`.
> `git diff HEAD~1..HEAD -- tractor/trionics/_mngrs.py`
The generated implementation adds `_CtxExit`, whose `done` event
publishes an `outcome.Outcome[None]`. `_Cache.run_ctx()` records either
`Value(None)` or `Error(exc)` after the resource exit attempt. The final
MOC consumer signals `no_more_users`, waits for completion under a
shielded cancel scope, removes the per-key lock, and unwraps the outcome
so ordinary cleanup failures are raised at the consumer boundary.
`trio.Cancelled`, `KeyboardInterrupt`, and `SystemExit` continue through
the service task rather than being converted into regular cleanup
errors.
> `git diff HEAD~1..HEAD -- tests/test_resource_cache.py`
The generated regressions cover successful exit blocking, cleanup-error
delivery, final-user cancellation, cancellation combined with a cleanup
error, and service-nursery cancellation. The existing teardown re-entry
test now uses explicit events instead of a ten-second cleanup sleep and
asserts that the replacement resource is a fresh cache miss.
Verification:
`env PYTHONPATH="$PWD" /home/goodboy/repos/tractor/py313/bin/python -m pytest tests/test_resource_cache.py`
Result: `16 passed in 8.37s`.

View File

@ -0,0 +1,22 @@
# AI Prompt I/O Log - OpenCode
This directory tracks prompt inputs and model outputs for AI-assisted
development using `opencode`.
## Policy
Prompt logging follows the [NLNet generative AI policy][nlnet-ai]. All
substantive AI contributions are logged with:
- Model name and version
- Timestamps
- The prompts that produced the output
- Unedited model output (`.raw.md` files)
[nlnet-ai]: https://nlnet.nl/foundation/policies/generativeAI/
## Usage
Entries are created by the prompt-io workflow. Human contributors remain
accountable for all decisions. AI-generated content is never presented as
human-authored work.

View File

@ -62,6 +62,28 @@ the one-and-only registrar; boot then fails loudly with a
``RuntimeError`` if some other process already bound the registry
socket(s).
A dedicated registrar
---------------------
That second rule — *"if a registrar answers, boot as a plain
root"* — is all you need to run the registry as its own
**standalone process**, decoupled from any app tree's root. Boot
a bare ``tractor.run_daemon([], registry_addrs=[...])`` (a root
actor that does nothing but hold the registry), point your app
tree at the same ``registry_addrs``, and every actor discovers
through that *external* registrar instead of a tree-local one:
.. literalinclude:: ../../examples/dedicated_registrar.py
:caption: examples/dedicated_registrar.py
:language: python
This is the "registrar as a subsystem, not the root actor" shape.
Two caveats today (both tracked as #472 follow-ups):
``enable_transports`` is single-proto per runtime, so a registrar
can't yet serve multiple backends at once; and there's no way to
spawn a registrar as a *sub*-actor of a shared tree (only as its
own root), since ``start_actor()`` has no custom-``actor_cls``
hook.
Looking up actors
-----------------

View File

@ -238,9 +238,31 @@ Toward capability-based msging
The ``pld_spec`` + codec-hook layer is the foundation for the
long-game: **capability-based msging** where each dialog's
type contract doubles as a capability grant, negotiated as part
of the protocol itself. That work is tracked in `#196`_ (with the
original typed-proto epic in `#36`_); if strongly-typed
distributed systems get you going, we'd love your input.
of the protocol itself. The epic is tracked in `#196`_ (evolving
the original typed-proto work in `#36`_), and the most recent
concrete step is `#365`_ — driving the whole ``pld_spec`` off
plain type-annotations (e.g. annotating a context's
``open_stream()`` with ``msgspec.Struct`` subtypes) instead of
explicit ``pld_spec=`` kwargs.
You don't have to wait for that, though: the decorator-level
``@tractor.context(pld_spec=...)`` shown above is already the
*higher-level* way to pin a dialog's payload contract, while
``tractor.msg._ops.limit_plds()`` is the lower-level, per-block
escape hatch. Both are exercised end-to-end in
``tests/msg/test_pldrx_limiting.py`` and
``tests/msg/test_ext_types_msgspec.py``.
On the codec-hook side, the ``enc_hook``/``dec_hook`` pair is
today only reachable via ``tractor.msg._ops``; a public *factory*
API for them is drafted in `#376`_ (from
`@guilledk <https://github.com/guilledk>`_, on the
`auto_codecs <https://github.com/goodboy/tractor/tree/auto_codecs>`_
branch) — the likely long-term home for custom-type
(de)serialization.
If strongly-typed distributed systems get you going, we'd love
your input on any of the above.
Where to next?
--------------
@ -258,3 +280,5 @@ Where to next?
.. _(un)protocol: https://zguide.zeromq.org/docs/chapter7/#Unprotocols
.. _#196: https://github.com/goodboy/tractor/issues/196
.. _#36: https://github.com/goodboy/tractor/issues/36
.. _#365: https://github.com/goodboy/tractor/issues/365
.. _#376: https://github.com/goodboy/tractor/pull/376

View File

@ -8,29 +8,31 @@ the_line = 'Hi my name is {}'
tractor.log.get_console_log("INFO")
async def hi():
async def hi() -> str:
return the_line.format(tractor.current_actor().name)
async def say_hello(other_actor):
async def say_hello(other_actor: str) -> str:
portal: tractor.Portal
async with tractor.wait_for_actor(other_actor) as portal:
return await portal.run(hi)
async def main():
async def main() -> None:
"""Main tractor entry point, the "master" process (for now
acts as the "director").
"""
n: tractor.ActorNursery
async with tractor.open_nursery() as n:
print("Alright... Action!")
donny = await n.run_in_actor(
donny: tractor.Portal = await n.run_in_actor(
say_hello,
name='donny',
# arguments are always named
other_actor='gretchen',
)
gretchen = await n.run_in_actor(
gretchen: tractor.Portal = await n.run_in_actor(
say_hello,
name='gretchen',
other_actor='donny',

View File

@ -2,17 +2,18 @@ import trio
import tractor
async def cellar_door():
async def cellar_door() -> str:
assert not tractor.is_root_process()
return "Dang that's beautiful"
async def main():
async def main() -> None:
"""The main ``tractor`` routine.
"""
n: tractor.ActorNursery
async with tractor.open_nursery() as n:
portal = await n.run_in_actor(
portal: tractor.Portal = await n.run_in_actor(
cellar_door,
name='some_linguist',
)

View File

@ -2,19 +2,20 @@ import trio
import tractor
async def movie_theatre_question():
async def movie_theatre_question() -> str:
"""A question asked in a dark theatre, in a tangent
(errr, I mean different) process.
"""
return 'have you ever seen a portal?'
async def main():
async def main() -> None:
"""The main ``tractor`` routine.
"""
n: tractor.ActorNursery
async with tractor.open_nursery() as n:
portal = await n.start_actor(
portal: tractor.Portal = await n.start_actor(
'frank',
# enable the actor to run funcs from this current module
enable_modules=[__name__],

View File

@ -13,11 +13,12 @@ async def stream_forever() -> AsyncIterator[int]:
await trio.sleep(0.01)
async def main():
async def main() -> None:
n: tractor.ActorNursery
async with tractor.open_nursery() as n:
portal = await n.start_actor(
portal: tractor.Portal = await n.start_actor(
'donny',
enable_modules=[__name__],
)
@ -25,7 +26,7 @@ async def main():
# this async for loop streams values from the above
# async generator running in a separate process
async with portal.open_stream_from(stream_forever) as stream:
count = 0
count: int = 0
async for letter in stream:
print(letter)
count += 1

View File

@ -35,7 +35,7 @@ async def open_ctx(
assert first is None
async def main():
async def main() -> None:
async with tractor.open_nursery(
debug_mode=True,

View File

@ -20,7 +20,7 @@ async def name_error():
getattr(doggypants) # noqa
async def main():
async def main() -> None:
'''
Test breakpoint in a streaming actor.

View File

@ -21,6 +21,7 @@ async def breakpoint_forever():
async def spawn_until(depth=0):
""""A nested nursery that triggers another ``NameError``.
"""
n: tractor.ActorNursery
async with tractor.open_nursery() as n:
if depth < 1:
@ -46,7 +47,7 @@ async def spawn_until(depth=0):
# TODO: notes on the new boxed-relayed errors through proxy actors
async def main():
async def main() -> None:
"""The main ``tractor`` routine.
The process tree should look as approximately as follows when the debugger

View File

@ -15,6 +15,7 @@ async def name_error():
async def spawn_error():
""""A nested nursery that triggers another ``NameError``.
"""
n: tractor.ActorNursery
async with tractor.open_nursery() as n:
portal = await n.run_in_actor(
name_error,
@ -23,7 +24,7 @@ async def spawn_error():
return await portal.result()
async def main():
async def main() -> None:
"""The main ``tractor`` routine.
The process tree should look as approximately as follows:

View File

@ -17,6 +17,7 @@ async def name_error():
async def spawn_error():
""""A nested nursery that triggers another ``NameError``.
"""
n: tractor.ActorNursery
async with tractor.open_nursery() as n:
portal = await n.run_in_actor(
name_error,
@ -25,7 +26,7 @@ async def spawn_error():
return await portal.result()
async def main():
async def main() -> None:
"""The main ``tractor`` routine.
The process tree should look as approximately as follows:

View File

@ -5,7 +5,8 @@ async def die():
raise RuntimeError
async def main():
async def main() -> None:
tn: tractor.ActorNursery
async with tractor.open_nursery() as tn:
debug_actor = await tn.start_actor(

View File

@ -18,7 +18,7 @@ async def name_error(
raise
async def main():
async def main() -> None:
'''
Test 3 `PdbREPL` entries:
- one in the child due to manual `.post_mortem()`,

View File

@ -2,7 +2,7 @@ import trio
import tractor
async def main():
async def main() -> None:
async with tractor.open_root_actor(
debug_mode=True,

View File

@ -2,7 +2,7 @@ import trio
import tractor
async def main():
async def main() -> None:
async with tractor.open_root_actor(
debug_mode=True,
):

View File

@ -10,6 +10,7 @@ async def name_error():
async def spawn_until(depth=0):
""""A nested nursery that triggers another ``NameError``.
"""
n: tractor.ActorNursery
async with tractor.open_nursery() as n:
if depth < 1:
# await n.run_in_actor('breakpoint_forever', breakpoint_forever)
@ -23,7 +24,7 @@ async def spawn_until(depth=0):
)
async def main():
async def main() -> None:
'''
The process tree should look as approximately as follows when the
debugger first engages:

View File

@ -2,7 +2,7 @@ import trio
import tractor
async def main():
async def main() -> None:
async with tractor.open_root_actor(
debug_mode=True,
loglevel='cancel',

View File

@ -7,7 +7,7 @@ async def key_error():
return {}['doggy']
async def main():
async def main() -> None:
'''
Root is fail-after-cancelled while blocking and child RPC fails
simultaneously.

View File

@ -71,7 +71,7 @@ async def cancelled_before_pause(
await pm_on_cancelled()
async def main():
async def main() -> None:
async with tractor.open_nursery(
debug_mode=True,
) as n:

View File

@ -34,7 +34,7 @@ async def just_bp(
async def main():
async def main() -> None:
# !TODO, parametrize the --tpt-proto={key} with osenv vars just
# like we do for loglevel/spawn-backend!

View File

@ -12,7 +12,7 @@ async def breakpoint_forever():
await tractor.pause()
async def main():
async def main() -> None:
async with tractor.open_nursery(
debug_mode=True,

View File

@ -6,7 +6,7 @@ async def name_error():
getattr(doggypants) # noqa (on purpose)
async def main():
async def main() -> None:
async with tractor.open_nursery(
debug_mode=True,
) as an:

View File

@ -0,0 +1,122 @@
'''
Run a *dedicated* registrar as its own standalone process decoupled
from your app's root actor — and discover a service *through* it.
Normally the registrar **is** the root actor of your tree. Here we
instead boot a separate `tractor.run_daemon([], registry_addrs=[...])`
process whose *sole* job is to be the registry, then point our app
tree at it via `registry_addrs`. Because a registrar is already
reachable at that addr, our app's root actor does NOT become one — it
registers with (and discovers through) the external daemon. That's the
"registrar as a subsystem, not the root actor" pattern.
NB: `enable_transports` is single-proto per-runtime today (see
`tractor._root`), so this demos one transport; a genuinely
multi-backend registrar (and spawning one as a *sub*actor of a shared
tree) are future runtime work see the #472 follow-ups.
'''
from contextlib import suppress
import signal
import socket
import subprocess
import sys
import time
import trio
import tractor
# the fixed addr the dedicated registrar binds and everyone points at.
REG_ADDR: tuple[str, int] = ('127.0.0.1', 1717)
def _wait_registrar_ready(
addr: tuple[str, int],
proc: subprocess.Popen,
deadline: float = 10.0,
) -> None:
'''
Active-poll the registrar's bind addr until it accepts a
connection (proving it's booted + listening), bailing early if
the daemon proc dies during startup.
'''
end: float = time.monotonic() + deadline
while time.monotonic() < end:
if proc.poll() is not None:
raise RuntimeError(
f'registrar died on startup (rc={proc.returncode})'
)
with suppress(OSError):
with socket.create_connection(addr, timeout=0.1):
return
time.sleep(0.05)
raise TimeoutError(f'registrar never came up @ {addr}')
async def greet() -> str:
'''A trivial service task any peer can RPC by name.'''
return f'hello from {tractor.current_actor().name}!'
async def app() -> None:
'''
Point our app tree at the EXTERNAL registrar (not its own root)
via `registry_addrs`, register a named service, then discover +
RPC it purely by name.
'''
an: tractor.ActorNursery
async with tractor.open_nursery(
registry_addrs=[REG_ADDR],
enable_transports=['tcp'],
) as an:
# this subactor registers with the DEDICATED registrar @
# REG_ADDR (our root is a plain peer, not the registry).
await an.start_actor(
'greeter',
enable_modules=[__name__],
)
# discover it *through the external registrar*, by name only.
portal: tractor.Portal
async with tractor.wait_for_actor('greeter') as portal:
print(f'found `greeter` via dedicated registrar @ {REG_ADDR}')
print(await portal.run(greet))
await an.cancel()
def main() -> None:
# boot the dedicated registrar as its own process/tree: an empty
# `enable_modules` `run_daemon()` is just a root actor that does
# nothing but hold + serve the registry.
code: str = (
'import tractor; '
f'tractor.run_daemon([], registry_addrs={[REG_ADDR]!r}, '
"enable_transports=['tcp'], loglevel='error')"
)
registrar: subprocess.Popen = subprocess.Popen(
[sys.executable, '-c', code],
# the registry is a quiet background service; hush its logs +
# expected SIGINT-teardown traceback so the demo output stays
# focused on the discovery flow.
stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL,
)
try:
_wait_registrar_ready(REG_ADDR, registrar)
print(
f'dedicated registrar up @ {REG_ADDR} '
f'(pid {registrar.pid})'
)
trio.run(app)
finally:
# graceful SIGINT teardown of the standalone registrar.
registrar.send_signal(signal.SIGINT)
with suppress(subprocess.TimeoutExpired):
registrar.wait(timeout=10)
print('dedicated registrar shut down')
if __name__ == '__main__':
main()

View File

@ -9,14 +9,14 @@ from tractor import (
# this is the first 2 actors, streamer_1 and streamer_2
async def stream_data(seed):
async def stream_data(seed: int):
for i in range(seed):
yield i
await trio.sleep(0.0001) # trigger scheduler
# this is the third actor; the aggregator
async def aggregate(seed):
async def aggregate(seed: int):
'''
Ensure that the two streams we receive match but only stream
a single set of values to the parent.
@ -28,7 +28,7 @@ async def aggregate(seed):
for i in range(1, 3):
# fork/spawn call
portal = await an.start_actor(
portal: Portal = await an.start_actor(
name=f'streamer_{i}',
enable_modules=[__name__],
)
@ -37,7 +37,7 @@ async def aggregate(seed):
send_chan, recv_chan = trio.open_memory_channel(500)
async def push_to_chan(portal, send_chan):
async def push_to_chan(portal: Portal, send_chan):
# TODO: https://github.com/goodboy/tractor/issues/207
async with send_chan:
@ -49,6 +49,7 @@ async def aggregate(seed):
print(f"FINISHED ITERATING {portal.channel.uid}")
# spawn 2 trio tasks to collect streams and push to a local queue
n: trio.Nursery
async with trio.open_nursery() as n:
for portal in portals:

View File

@ -28,7 +28,7 @@ async def aio_echo_server(
@tractor.context
async def trio_to_aio_echo_server(
ctx: tractor.Context,
):
) -> None:
# this will block until the ``asyncio`` task sends a "first"
# message.
async with tractor.to_asyncio.open_channel_from(
@ -48,10 +48,11 @@ async def trio_to_aio_echo_server(
await stream.send(out)
async def main():
async def main() -> None:
n: tractor.ActorNursery
async with tractor.open_nursery() as n:
p = await n.start_actor(
p: tractor.Portal = await n.start_actor(
'aio_server',
enable_modules=[__name__],
infect_asyncio=True,

View File

@ -33,15 +33,16 @@ async def main() -> None:
rank = MPI.COMM_WORLD.Get_rank()
print(f"[parent] rank={rank} pid={os.getpid()}", flush=True)
an: tractor.ActorNursery
async with tractor.open_nursery(start_method='trio') as an:
portal = await an.start_actor(
portal: tractor.Portal = await an.start_actor(
'mpi-child',
enable_modules=[child_fn.__module__],
# Without this the child replays __main__, which
# re-imports mpi4py and crashes on MPI_Init.
inherit_parent_main=False,
)
result = await portal.run(child_fn)
result: str = await portal.run(child_fn)
print(f"[parent] got: {result}", flush=True)
await portal.cancel_actor()

View File

@ -5,7 +5,7 @@ import tractor
log = tractor.log.get_logger('multiportal')
async def stream_data(seed=10):
async def stream_data(seed: int = 10):
log.info("Starting stream task")
for i in range(seed):
@ -13,7 +13,10 @@ async def stream_data(seed=10):
await trio.sleep(0) # trigger scheduler
async def stream_from_portal(p, consumed):
async def stream_from_portal(
p: tractor.Portal,
consumed: list,
) -> None:
async with p.open_stream_from(stream_data) as stream:
async for item in stream:
@ -23,14 +26,19 @@ async def stream_from_portal(p, consumed):
consumed.append(item)
async def main():
async def main() -> None:
an: tractor.ActorNursery
async with tractor.open_nursery(loglevel='info') as an:
p = await an.start_actor('stream_boi', enable_modules=[__name__])
p: tractor.Portal = await an.start_actor(
'stream_boi',
enable_modules=[__name__],
)
consumed = []
consumed: list = []
n: trio.Nursery
async with trio.open_nursery() as n:
for i in range(2):
n.start_soon(stream_from_portal, p, consumed)

View File

@ -41,6 +41,7 @@ async def fan_out_squares(
aggregated squares to our parent.
'''
an: tractor.ActorNursery
async with tractor.open_nursery() as an:
portals: list[tractor.Portal] = []
for i in (1, 2):
@ -67,6 +68,7 @@ async def fan_out_squares(
)
# fan out one sub-RPC per input val, concurrently.
tn: trio.Nursery
async with trio.open_nursery() as tn:
for i, x in enumerate(vals):
tn.start_soon(
@ -83,8 +85,9 @@ async def fan_out_squares(
async def main() -> None:
an: tractor.ActorNursery
async with tractor.open_nursery() as an:
portal = await an.start_actor(
portal: tractor.Portal = await an.start_actor(
'supervisor',
enable_modules=[__name__],
)

View File

@ -31,7 +31,7 @@ PRIMES = [
]
async def is_prime(n):
async def is_prime(n: int) -> bool:
if n < 2:
return False
if n == 2:
@ -47,7 +47,7 @@ async def is_prime(n):
@acm
async def worker_pool(workers=4):
async def worker_pool(workers: int = 4):
"""Though it's a trivial special case for ``tractor``, the well
known "worker pool" seems to be the defacto "but, I want this
process pattern!" for most parallelism pilgrims.
@ -55,9 +55,10 @@ async def worker_pool(workers=4):
Yes, the workers stay alive (and ready for work) until you close
the context.
"""
tn: tractor.ActorNursery
async with tractor.open_nursery() as tn:
portals = []
portals: list[tractor.Portal] = []
snd_chan, recv_chan = trio.open_memory_channel(len(PRIMES))
for i in range(workers):
@ -77,9 +78,14 @@ async def worker_pool(workers=4):
) -> list[bool]:
# define an async (local) task to collect results from workers
async def send_result(func, value, portal):
async def send_result(
func: Callable,
value: int,
portal: tractor.Portal,
):
await snd_chan.send((value, await portal.run(func, n=value)))
n: trio.Nursery
async with trio.open_nursery() as n:
for value, portal in zip(sequence, itertools.cycle(portals)):
@ -101,7 +107,7 @@ async def worker_pool(workers=4):
await tn.cancel()
async def main():
async def main() -> None:
async with worker_pool() as actor_map:

View File

@ -12,7 +12,7 @@ import tractor
import trio
async def burn_cpu():
async def burn_cpu() -> int:
pid = os.getpid()
@ -23,17 +23,18 @@ async def burn_cpu():
return os.getpid()
async def main():
async def main() -> None:
n: tractor.ActorNursery
async with tractor.open_nursery() as n:
portal = await n.run_in_actor(burn_cpu)
portal: tractor.Portal = await n.run_in_actor(burn_cpu)
# burn rubber in the parent too
await burn_cpu()
# wait on result from target function
pid = await portal.wait_for_result()
pid: int = await portal.wait_for_result()
# end of nursery block
print(f"Collected subproc {pid}")

View File

@ -9,12 +9,13 @@ async def sleepy_jane() -> None:
await trio.sleep_forever()
async def main():
async def main() -> None:
'''
Spawn a flat actor cluster, with one process per detected core.
'''
portal_map: dict[str, tractor.Portal]
tn: trio.Nursery
# look at this hip new syntax!
async with (

View File

@ -2,13 +2,14 @@ import trio
import tractor
async def assert_err():
async def assert_err() -> None:
assert 0
async def main():
async def main() -> None:
n: tractor.ActorNursery
async with tractor.open_nursery() as n:
real_actors = []
real_actors: list[tractor.Portal] = []
for i in range(3):
real_actors.append(await n.start_actor(
f'actor_{i}',

View File

@ -31,9 +31,10 @@ async def simple_rpc(
async def main() -> None:
n: tractor.ActorNursery
async with tractor.open_nursery() as n:
portal = await n.start_actor(
portal: tractor.Portal = await n.start_actor(
'rpc_server',
enable_modules=[__name__],
)

View File

@ -49,13 +49,15 @@ async def client_task() -> None:
async def main() -> None:
an: tractor.ActorNursery
async with tractor.open_nursery() as an:
portal = await an.start_actor(
portal: tractor.Portal = await an.start_actor(
'quote_svc',
enable_modules=[__name__],
)
# run the client in a separate task which discovers
# the daemon purely by its registered name.
tn: trio.Nursery
async with trio.open_nursery() as tn:
tn.start_soon(client_task)
# explicit graceful teardown of the daemon.

View File

@ -4,14 +4,17 @@ import tractor
tractor.log.get_console_log("INFO")
async def main(service_name):
async def main(service_name: str) -> None:
an: tractor.ActorNursery
async with tractor.open_nursery() as an:
await an.start_actor(service_name)
portal: tractor.Portal
async with tractor.get_registry() as portal:
print(f"Registrar is listening on {portal.channel}")
sockaddr: tractor.Portal
async with tractor.wait_for_actor(service_name) as sockaddr:
print(f"my_service is found at {sockaddr}")

View File

@ -54,8 +54,9 @@ async def consume(
async def main() -> None:
an: tractor.ActorNursery
async with tractor.open_nursery() as an:
portal = await an.start_actor(
portal: tractor.Portal = await an.start_actor(
'ticker',
enable_modules=[__name__],
)
@ -67,6 +68,7 @@ async def main() -> None:
ctx.open_stream() as stream,
):
assert first == 5
tn: trio.Nursery
async with trio.open_nursery() as tn:
# use `.start()` so each consumer is known
# to be subscribed before the ticks flow.

View File

@ -32,8 +32,8 @@ async def acquire_singleton_lock(
async def hold_lock_forever(
task_status=trio.TASK_STATUS_IGNORED
):
task_status: trio.TaskStatus = trio.TASK_STATUS_IGNORED,
) -> None:
async with (
tractor.trionics.maybe_raise_from_masking_exc(),
acquire_singleton_lock() as lock,
@ -46,7 +46,7 @@ async def main(
ignore_special_cases: bool,
loglevel: str = 'info',
debug_mode: bool = True,
):
) -> None:
async with (
trio.open_nursery() as tn,

View File

@ -134,7 +134,7 @@ async def main(
raise_unmasked: bool = False,
loglevel: str = 'info',
):
) -> None:
tractor.log.get_console_log(level=loglevel)
# the `.aclose()` being checkpoints on these

View File

@ -59,8 +59,9 @@ async def point_doubler(
async def main() -> None:
an: tractor.ActorNursery
async with tractor.open_nursery() as an:
portal = await an.start_actor(
portal: tractor.Portal = await an.start_actor(
'point_doubler',
enable_modules=[__name__],
)

View File

@ -27,10 +27,11 @@ async def report_addr() -> str:
async def main() -> None:
an: tractor.ActorNursery
async with tractor.open_nursery(
enable_transports=['uds'],
) as an:
portal = await an.start_actor(
portal: tractor.Portal = await an.start_actor(
'uds_child',
enable_modules=[__name__],
)

View File

@ -2,16 +2,19 @@
`tractor.log`-wrapping unit tests.
'''
import logging
from pathlib import Path
import shutil
from types import ModuleType
import pytest
import tractor
import trio
from tractor import (
_code_load,
log,
)
from tractor.ipc import _chan
def test_root_pkg_not_duplicated_in_logger_name():
@ -222,6 +225,88 @@ def test_add_log_level_pluggable():
delattr(log.StackLevelAdapter, name.lower())
@pytest.mark.parametrize(
'suppression',
[
'level',
'logger',
'global',
],
)
def test_log_guard_skips_payload_formatting(
monkeypatch: pytest.MonkeyPatch,
suppression: str,
):
'''
Suppressed transport logs must not render payloads.
The original hot-path guard compared only the effective logger
level. A logger disabled through its `Logger.disabled` flag or
the global `logging.disable()` threshold could therefore still
call `pformat()` before `Logger.isEnabledFor()` discarded the
record.
Exercise effective-level, per-logger, and global suppression
independently. A poisoned `_chan.pformat()` proves rendering is
skipped, while the fake transport proves `Channel.send()` still
transmits the original payload and traceback-hiding flag.
'''
sent: list[tuple[object, bool]] = []
class FakeTransport:
async def send(
self,
payload: object,
hide_tb: bool = False,
) -> None:
sent.append((payload, hide_tb))
def fail_pformat(payload: object) -> str:
raise AssertionError(
f'suppressed log rendered payload: {payload!r}'
)
chan_log = log.get_logger(
name=f'guard_test.{suppression}',
)
std_log = chan_log.logger
orig_level: int = std_log.level
orig_disable: int = logging.root.manager.disable
transport_level: int = log.CUSTOM_LEVELS['TRANSPORT']
monkeypatch.setattr(_chan, 'log', chan_log)
monkeypatch.setattr(_chan, 'pformat', fail_pformat)
try:
logging.disable(logging.NOTSET)
std_log.setLevel(transport_level)
if suppression == 'level':
std_log.setLevel(logging.INFO)
elif suppression == 'logger':
monkeypatch.setattr(std_log, 'disabled', True)
else:
logging.disable(logging.CRITICAL)
assert not chan_log.isEnabledFor(transport_level)
transport = FakeTransport()
chan = _chan.Channel(transport=transport)
payload = object()
async def send_payload() -> None:
await chan.send(
payload,
hide_tb=True,
)
trio.run(send_payload)
assert sent == [(payload, True)]
finally:
std_log.setLevel(orig_level)
logging.disable(orig_disable)
# TODO, moar tests against existing feats:
# ------ - ------
# - [ ] color settings?

View File

@ -9,6 +9,7 @@ from typing import Awaitable
import pytest
import trio
from trio.testing import wait_all_tasks_blocked
import tractor
from tractor.trionics import (
maybe_open_context,
@ -94,6 +95,232 @@ def test_resource_only_entered_once(key_on):
trio.run(main)
def test_last_moc_user_waits_for_resource_exit():
'''
Verify the final user cannot return before resource teardown.
Previously the final `maybe_open_context()` user only signalled
`_Cache.run_ctx()` through its `no_more_users` event. The user
then returned while the service task was still running the
resource's `__aexit__()`, so callers could observe stale external
state immediately after their `async with` block.
The resource sets `exit_started` before blocking on
`allow_exit`. The user task must remain inside MOC until the test
releases that deterministic checkpoint and `__aexit__()` sets
`exit_finished`.
'''
async def main():
exit_started = trio.Event()
allow_exit = trio.Event()
exit_finished = trio.Event()
user_returned = trio.Event()
@acm
async def open_resource():
try:
yield
finally:
exit_started.set()
await allow_exit.wait()
exit_finished.set()
async def use_resource():
async with maybe_open_context(open_resource):
pass
assert exit_finished.is_set()
user_returned.set()
async with (
tractor.open_root_actor(),
trio.open_nursery() as tn,
):
tn.start_soon(use_resource)
await exit_started.wait()
assert not user_returned.is_set()
allow_exit.set()
await user_returned.wait()
trio.run(main)
def test_moc_delivers_resource_exit_error():
'''
Verify a resource exit error reaches the final MOC user.
Previously `_Cache.run_ctx()` executed the cached resource's
`__aexit__()` after the final user had returned. An exit failure
therefore surfaced later through the actor service nursery rather
than at the user's `async with maybe_open_context()` boundary.
This resource raises a unique `ResourceExitError` during exit.
Catching that exact instance around MOC proves the service task
delivered the failure to the final user without replacing it.
'''
class ResourceExitError(Exception):
pass
exit_error = ResourceExitError('resource exit failed')
async def main():
@acm
async def open_resource():
yield
raise exit_error
async with tractor.open_root_actor():
with pytest.raises(ResourceExitError) as exc_info:
async with maybe_open_context(open_resource):
pass
assert exc_info.value is exit_error
trio.run(main)
def test_moc_final_user_cancellation_waits_for_exit():
'''
Verify final-user cancellation still waits for successful exit.
Previously cancellation escaped the final MOC user immediately
after it signalled `_Cache.run_ctx()`, leaving resource exit to
finish later in the actor service task. This violated the context
manager boundary even when cleanup itself succeeded.
The consumer cancels its own scope while holding the sole cached
resource. The resource sets `exit_finished` from its `finally`
block, and the consumer checks that event immediately after its
cancel scope catches `trio.Cancelled`. This proves MOC's
completion wait is shielded without suppressing the original
cancellation.
'''
async def main():
exit_finished = trio.Event()
@acm
async def open_resource():
try:
yield
finally:
exit_finished.set()
async with tractor.open_root_actor():
with trio.CancelScope() as cs:
async with maybe_open_context(open_resource):
cs.cancel()
await trio.sleep_forever()
assert cs.cancelled_caught
assert exit_finished.is_set()
trio.run(main)
def test_moc_exit_error_masks_final_user_cancellation():
'''
Verify cleanup errors survive final-user cancellation.
A cancelled final user previously signalled `no_more_users` and
propagated `trio.Cancelled` before `_Cache.run_ctx()` completed
resource exit. If `__aexit__()` then failed, its error was
detached from the API call which caused teardown.
The consumer cancels its own scope at a deterministic checkpoint
inside MOC. Resource exit raises `ResourceExitError`; observing
that exact error outside the cancel scope proves MOC shields the
completion wait and applies normal context-manager masking, where
a cleanup failure replaces the active cancellation.
'''
class ResourceExitError(Exception):
pass
exit_error = ResourceExitError('resource exit failed')
async def main():
@acm
async def open_resource():
yield
raise exit_error
async with tractor.open_root_actor():
with pytest.raises(ResourceExitError) as exc_info:
with trio.CancelScope() as cs:
async with maybe_open_context(open_resource):
cs.cancel()
await trio.sleep_forever()
assert exc_info.value is exit_error
trio.run(main)
def test_moc_service_nursery_cancellation_completes_exit():
'''
Verify service-nursery cancellation cannot strand a final user.
`_Cache.run_ctx()` and an MOC consumer may share a
caller-provided service nursery. Cancelling that nursery
interrupts the service task's `no_more_users` wait and the
consumer body together. A shielded final-user wait would deadlock
if `run_ctx()` failed to publish completion while propagating its
own `trio.Cancelled`.
The outer task waits for resource entry, then cancels the exact
nursery containing both tasks. The resource shields one cleanup
checkpoint and sets `exit_finished`; observing both it and
`service_finished` proves cancellation propagated normally while
MOC's completion handshake terminated deterministically.
'''
async def main():
resource_entered = trio.Event()
exit_finished = trio.Event()
service_finished = trio.Event()
service_tn: trio.Nursery|None = None
@acm
async def open_resource():
try:
resource_entered.set()
yield
finally:
with trio.CancelScope(shield=True):
await trio.lowlevel.checkpoint()
exit_finished.set()
async def use_resource(tn: trio.Nursery):
async with maybe_open_context(
open_resource,
tn=tn,
):
await trio.sleep_forever()
async def run_service():
nonlocal service_tn
async with trio.open_nursery() as tn:
service_tn = tn
tn.start_soon(use_resource, tn)
service_finished.set()
async with trio.open_nursery() as outer_tn:
outer_tn.start_soon(run_service)
await resource_entered.wait()
assert service_tn is not None
service_tn.cancel_scope.cancel()
await service_finished.wait()
assert exit_finished.is_set()
trio.run(main)
@tractor.context
async def streamer(
ctx: tractor.Context,
@ -548,43 +775,52 @@ def test_moc_reentry_during_teardown(
loglevel: str,
):
'''
Reproduce the piker `open_cached_client('kraken')` race:
Reproduce re-entry while an identical cached context exits.
- same `acm_func`, NO kwargs (identical `ctx_key`)
- multiple tasks share the cached resource
- all users exit -> teardown starts
- a NEW task enters during `_Cache.run_ctx.__aexit__`
- `values[ctx_key]` is gone (popped in inner finally)
but `resources[ctx_key]` still exists (outer finally
hasn't run yet bc the acm cleanup has checkpoints)
- old code: `assert not resources.get(ctx_key)` FIRES
- multiple tasks use the same `acm_func` with no kwargs,
producing an identical `ctx_key`;
- all users leave and the final user starts resource teardown;
- `_Cache.run_ctx()` removes the cached value and resource entry
before entering the resource's blocking `__aexit__()` body;
- a new task attempts to enter that same `ctx_key` during exit;
- the per-key lock keeps that entrant queued until exit
completes;
- the entrant then receives a fresh cache miss and resource.
This models the real-world scenario where `brokerd.kraken`
tasks concurrently call `open_cached_client('kraken')`
(same `acm_func`, empty kwargs, shared `ctx_key`) and
the teardown/re-entry race triggers intermittently.
Without teardown sharing the registration lock, re-entry could
race resource replacement while the prior generation was still
exiting. The final user could also return before that exit
completed.
The first resource generation signals `in_aexit` and waits on
`allow_aexit`. The re-entry task signals `reentry_started` and
blocks inside MOC; only after `wait_all_tasks_blocked()` confirms
that ordering does the coordinator release cleanup. The entrant
must then receive a fresh cache miss. `first_done` additionally
proves the first MOC user observed completed teardown before
returning.
'''
async def main():
in_aexit = trio.Event()
allow_aexit = trio.Event()
reentry_started = trio.Event()
generation: int = 0
@acm
async def cached_client():
'''
Simulates `kraken.api.get_client()`:
- no params (all callers share one `ctx_key`)
- slow-ish cleanup to widen the race window
between `values.pop()` and `resources.pop()`
inside `_Cache.run_ctx`.
Simulate a no-argument `kraken.api.get_client()`.
'''
nonlocal generation
generation += 1
resource_generation: int = generation
yield 'the-client'
# Signal that we're in __aexit__ — at this
# point `values` has already been popped by
# `run_ctx`'s inner finally, but `resources`
# is still alive (outer finally hasn't run).
if resource_generation == 1:
in_aexit.set()
await trio.sleep(10)
await allow_aexit.wait()
first_done = trio.Event()
@ -598,16 +834,25 @@ def test_moc_reentry_during_teardown(
async def reenter_during_teardown():
'''
Wait for the acm's `__aexit__` to start (meaning
`values` is popped but `resources` still exists),
then re-enter triggering the assert.
the cached value is no longer available), then re-enter.
'''
await in_aexit.wait()
# Tell the coordinator this task is about to enter MOC.
# `Event.set()` is not a checkpoint. Though `async with`
# awaits MOC's `__aenter__()`, its async generator runs
# synchronously until the held per-key `lock.acquire()`
# actually suspends this task.
reentry_started.set()
async with maybe_open_context(
cached_client,
) as (cache_hit, value):
assert not cache_hit
assert value == 'the-client'
await first_done.wait()
with trio.fail_after(5):
async with (
tractor.open_root_actor(
@ -619,5 +864,15 @@ def test_moc_reentry_during_teardown(
):
tn.start_soon(use_and_exit)
tn.start_soon(reenter_during_teardown)
await reentry_started.wait()
# Wait until the re-entry task is queued on MOC's
# per-key lock while `_Cache.run_ctx()` remains
# blocked in the first generation's `__aexit__()`.
# Only then release cleanup, making the intended
# enter-during-sibling-exit ordering deterministic.
await wait_all_tasks_blocked()
assert not first_done.is_set()
allow_aexit.set()
trio.run(main)

View File

@ -0,0 +1,333 @@
'''
`tractor.trionics._taskc.start_or_cancel()` unit tests.
`trio.Nursery.start()` collapses an out-of-band (ancestor)
cancellation into a lossy,
`RuntimeError('child exited without calling
task_status.started()')`
whenever the started child exits pre-`.started()` WITHOUT
propagating the ambient `trio.Cancelled`; a common outcome
when the child (or any lib code it calls) runs a graceful
teardown which absorbs the cancel and returns early. Our
`start_or_cancel()` wrapper re-surfaces the real in-flight
cancellation in that case so the true root error/cancel
propagates to the `.start()` caller instead.
These tests verify both that repair AND document upstream
`trio`'s current lossy behaviour via the
`use_start_or_cancel=False` parametrizations; if a `trio`
upgrade breaks one of THOSE cases it likely means upstream
shipped better startup-cancellation porcelain and our
wrapper deserves a re-audit!
The core use case was dug out of `modden`'s
`progman.open_wks()` program-spawn machinery as per gh
issue #474; the wrapper landed originally via gh PR #464.
'''
import pytest
import trio
from trio import TaskStatus
from tractor.trionics import start_or_cancel
async def absorbs_cancel_pre_started(
task_status: TaskStatus[None] = trio.TASK_STATUS_IGNORED,
):
'''
Swallow the ambient (ancestor-scope) cancel and return
early, a naughty-but-realistic graceful-teardown pattern
and the exact shape which causes `trio.Nursery.start()`
to raise its lossy startup `RuntimeError` in place of
the real `trio.Cancelled`.
'''
try:
await trio.sleep_forever()
except trio.Cancelled:
return
async def raise_val_err():
'''
Sibling task which blows up (fast) thus OOB-cancelling
the shared parent-nursery's cancel-scope.
'''
await trio.lowlevel.checkpoint()
raise ValueError('sibling blew up!')
@pytest.mark.parametrize(
'use_start_or_cancel',
[
True,
False,
],
)
def test_sibling_err_not_masked_by_startup_rte(
use_start_or_cancel: bool,
):
'''
The `modden.runtime.progman` use case: a sibling task
errors while the `.start()`-ed child is still
pre-`.started()`, OOB-cancelling the shared nursery
scope; the child absorbs its cancel (graceful teardown)
and exits early.
- with `start_or_cancel()` the in-flight cancellation
is re-surfaced as the real `trio.Cancelled` (then
absorbed by the cancelled nursery scope) so ONLY the
root-cause sibling error escapes the nursery.
- with a bare `.start()`, upstream `trio` (currently)
also delivers its lossy startup `RuntimeError`
alongside, obscuring that the child was in fact
cancelled due to the sibling's error.
`cancelled_at_start` records that the wrapper's own await
raises `Cancelled`, rather than merely relying on the
nursery's eventual exception-group shape.
'''
cancelled_at_start: list[bool] = []
async def main():
async with trio.open_nursery() as tn:
tn.start_soon(raise_val_err)
if use_start_or_cancel:
try:
await start_or_cancel(
tn,
absorbs_cancel_pre_started,
)
except trio.Cancelled:
cancelled_at_start.append(True)
raise
else:
await tn.start(absorbs_cancel_pre_started)
with pytest.raises(ExceptionGroup) as excinfo:
trio.run(main)
eg: ExceptionGroup = excinfo.value
val_eg, rest_eg = eg.split(ValueError)
assert len(val_eg.exceptions) == 1
if use_start_or_cancel:
assert cancelled_at_start == [True]
# the re-surfaced `Cancelled` is absorbed by the
# (sibling-error cancelled) nursery scope leaving
# NO startup-noise, just the root cause.
assert rest_eg is None
else:
# the `trio` wart: a lossy startup RTE rides along
# with (and distracts from) the root cause.
rte = rest_eg.exceptions[0]
assert isinstance(rte, RuntimeError)
assert 'child exited without calling' in rte.args[0]
@pytest.mark.parametrize(
'use_start_or_cancel',
[
True,
False,
],
)
def test_pure_oob_cancel_not_morphed_to_rte(
use_start_or_cancel: bool,
):
'''
A plain (error-free) ancestor `CancelScope.cancel()`
fired while the (cancel-absorbing) child is still
pre-`.started()`:
- `start_or_cancel()` re-surfaces the `Cancelled` so
the cancelled scope exits CLEAN, no error at all.
- a bare `.start()` (currently) morphs the plain
cancel into an (eg-wrapped) startup `RuntimeError`.
`cancelled_at_start` proves cancellation interrupts the
wrapper call itself before the cancelled scope exits.
'''
cancelled_at_start: list[bool] = []
async def main():
with trio.CancelScope() as cs:
async with trio.open_nursery() as tn:
async def canceller():
await trio.lowlevel.checkpoint()
cs.cancel()
tn.start_soon(canceller)
if use_start_or_cancel:
try:
await start_or_cancel(
tn,
absorbs_cancel_pre_started,
)
except trio.Cancelled:
cancelled_at_start.append(True)
raise
else:
await tn.start(
absorbs_cancel_pre_started,
)
assert cs.cancelled_caught
if use_start_or_cancel:
trio.run(main)
assert cancelled_at_start == [True]
else:
with pytest.raises(ExceptionGroup) as excinfo:
trio.run(main)
rte = excinfo.value.exceptions[0]
assert isinstance(rte, RuntimeError)
assert 'child exited without calling' in rte.args[0]
@pytest.mark.parametrize(
'use_start_or_cancel',
[
True,
False,
],
)
def test_genuine_startup_rte_still_raised(
use_start_or_cancel: bool,
):
'''
Absent ANY in-flight cancellation, a child exiting
cleanly without calling `task_status.started()` is a
genuine startup-protocol bug; `start_or_cancel()` must
re-raise the resulting `RuntimeError` exactly like a
bare `.start()` does.
'''
async def exits_wo_started(
task_status: TaskStatus[None] = (
trio.TASK_STATUS_IGNORED
),
):
await trio.lowlevel.checkpoint()
async def main():
async with trio.open_nursery() as tn:
with pytest.raises(RuntimeError) as excinfo:
if use_start_or_cancel:
await start_or_cancel(
tn,
exits_wo_started,
)
else:
await tn.start(exits_wo_started)
rte = excinfo.value
assert (
'child exited without calling'
in
rte.args[0]
)
trio.run(main)
@pytest.mark.parametrize(
'rte_arg',
[
# Broad substring matches would wrongly demote either
# child-owned error to `Cancelled` under cancellation.
'never got started!',
'child exited without calling user hook',
# non-`str` first-arg edge; must not `TypeError`
# inside the wrapper's msg-match guard.
1234,
],
)
def test_childs_own_rte_never_demoted_to_cancel(
rte_arg: str|int,
):
'''
A child's OWN `RuntimeError`, one which merely smells
like `trio`'s startup wording (or carries a non-`str`
first arg), raised under ambient cancellation must NOT
be demoted to a `trio.Cancelled` by the exact-msg-match
guard inside `start_or_cancel()`; the real error must
always propagate to the caller as the sole exception-group
leaf, preserving object identity.
'''
child_rte = RuntimeError(rte_arg)
async def cancels_cs_then_raises(
task_status: TaskStatus[None] = (
trio.TASK_STATUS_IGNORED
),
):
# cancel the ambient (ancestor) scope then raise
# sync-ly, no checkpoint between, so the child
# deterministically dies with ITS error while the
# caller is under effective cancellation.
cs.cancel()
raise child_rte
cs = trio.CancelScope()
async def main():
with cs:
async with trio.open_nursery() as tn:
await start_or_cancel(
tn,
cancels_cs_then_raises,
)
with pytest.raises(ExceptionGroup) as excinfo:
trio.run(main)
assert excinfo.value.exceptions == (child_rte,)
def test_started_value_and_args_passthru():
'''
Happy path: positional args, the `name=` kwarg and the
`.started(value)`-delivered value all pass through
`start_or_cancel()` identically to a bare `.start()`.
'''
async def echo_started(
*args,
task_status: TaskStatus[tuple] = (
trio.TASK_STATUS_IGNORED
),
):
task_name: str = trio.lowlevel.current_task().name
task_status.started((
args,
task_name,
))
async def main():
async with trio.open_nursery() as tn:
(
args,
task_name,
) = await start_or_cancel(
tn,
echo_started,
'chillin',
10,
name='doggy',
)
assert args == ('chillin', 10)
assert task_name == 'doggy'
trio.run(main)

View File

@ -198,9 +198,6 @@ class Channel:
# assert transport.raddr == addr
chan = Channel(transport=transport)
# ?TODO, compact this into adapter level-methods?
# -[ ] would avoid extra repr-calcs if level not active?
# |_ how would the `calc_if_level` look though? func?
if log.at_least_level('runtime'):
from tractor.devx import (
pformat as _pformat,
@ -325,6 +322,8 @@ class Channel:
'''
__tracebackhide__: bool = hide_tb
try:
if log.at_least_level('transport'):
# don't materialize the payload repr if not necessary
log.transport(
'=> send IPC msg:\n\n'
f'{pformat(payload)}\n'

View File

@ -309,7 +309,10 @@ class MsgpackTransport(MsgTransport):
log.transport(f'received header {size}') # type: ignore
msg_bytes: bytes = await self.recv_stream.receive_exactly(size)
log.transport(f"received {msg_bytes}") # type: ignore
if log.at_least_level('transport'):
log.transport( # type: ignore
f'received {msg_bytes}'
)
try:
# NOTE: lookup the `trio.Task.context`'s var for
# the current `MsgCodec`.

View File

@ -111,9 +111,7 @@ def at_least_level(
if isinstance(level, str):
level: int = CUSTOM_LEVELS[level.upper()]
if log.getEffectiveLevel() <= level:
return True
return False
return log.isEnabledFor(level)
# TODO, compare with using a "filter" instead?

View File

@ -306,12 +306,12 @@ class PldRx(Struct):
):
try:
pld: PayloadT = self._pld_dec.decode(pld)
if log.at_least_level('runtime'):
# don't materialize the payload repr if not necessary
log.runtime(
f'Decoded payload for\n'
# f'\n'
f'\n'
f'{msg}\n'
# ^TODO?, ideally just render with `,
# pld={decode}` in the `msg.pformat()`??
f'where, '
f'{type(msg).__name__}.pld={pld!r}\n'
)

View File

@ -1003,17 +1003,15 @@ async def process_messages(
task_status.started(loop_cs)
async for msg in chan:
if log.at_least_level('transport'):
log.transport( # type: ignore
f'IPC msg from peer\n'
f'<= {chan.aid.reprol()}\n\n'
# TODO: use of the pprinting of structs is
# FRAGILE and should prolly not be
#
# avoid fmting depending on loglevel for perf?
# -[ ] specifically `pretty_struct.pformat()` sub-call..?
# - how to only log-level-aware actually call this?
# -[ ] use `.msg.pretty_struct` here now instead!
# TODO: pretty-printing structs is FRAGILE;
# -[ ] add a non-raising log formatter with
# native-repr fallback before using
# `.msg.pretty_struct` here.
# f'{pretty_struct.pformat(msg)}\n'
f'{msg}\n'
)
@ -1262,6 +1260,7 @@ async def process_messages(
log.exception(message)
raise RuntimeError(message)
if log.at_least_level('transport'):
log.transport(
'Waiting on next IPC msg from\n'
f'peer: {chan.aid.reprol()}\n'

View File

@ -198,6 +198,16 @@ async def gather_contexts(
# Further potential examples of interest:
# https://gist.github.com/njsmith/cf6fc0a97f53865f2c671659c88c1798#file-cache-py-L8
class _CtxExit:
'''
Completion state for a cached context's shared exit.
'''
def __init__(self) -> None:
self.done = trio.Event()
self.error: Exception|None = None
class _Cache:
'''
Globally (actor-processs scoped) cached, task access to
@ -213,7 +223,11 @@ class _Cache:
values: dict[Any, Any] = {}
resources: dict[
Hashable,
tuple[trio.Nursery, trio.Event]
tuple[
trio.Nursery,
trio.Event,
_CtxExit,
],
] = {}
# nurseries: dict[int, trio.Nursery] = {}
no_more_users: trio.Event|None = None
@ -223,19 +237,39 @@ class _Cache:
cls,
mng,
ctx_key: tuple,
ctx_exit: _CtxExit,
task_status: trio.TaskStatus[T] = trio.TASK_STATUS_IGNORED,
) -> None:
entered: bool = False
try:
async with mng as value:
_, no_more_users = cls.resources[ctx_key]
entered = True
(
_,
no_more_users,
_,
) = cls.resources[ctx_key]
cls.values[ctx_key] = value
task_status.started(value)
try:
await no_more_users.wait()
finally:
value = cls.values.pop(ctx_key)
cls.values.pop(ctx_key)
cls.resources.pop(ctx_key)
except Exception as exc:
if not entered:
raise
# Deliver regular `__aexit__()` failures to the final
# consumer instead of raising into the service nursery.
ctx_exit.error = exc
finally:
if entered:
ctx_exit.done.set()
class _UnresolvedCtx:
'''
@ -281,9 +315,10 @@ async def maybe_open_context(
)
# yielded output
# sentinel = object()
yielded: Any = _UnresolvedCtx
user_registered: bool = False
ctx_exit: _CtxExit|None = None
exit_error: Exception|None = None
# Lock resource acquisition around task racing / ``trio``'s
# scheduler protocol.
@ -300,7 +335,6 @@ async def maybe_open_context(
] = trio.StrictFIFOLock()
header: str = 'Allocated NEW lock for @acm_func,\n'
else:
await trio.lowlevel.checkpoint()
header: str = 'Reusing OLD lock for @acm_func,\n'
log.debug(
@ -368,7 +402,6 @@ async def maybe_open_context(
resources = _Cache.resources
entry: tuple|None = resources.get(ctx_key)
if entry:
service_tn, ev = entry
raise RuntimeError(
f'Caching resources ALREADY exist?!\n'
f'ctx_key={ctx_key!r}\n'
@ -376,12 +409,18 @@ async def maybe_open_context(
f'task: {task}\n'
)
resources[ctx_key] = (service_tn, trio.Event())
ctx_exit = _CtxExit()
resources[ctx_key] = (
service_tn,
trio.Event(),
ctx_exit,
)
try:
yielded: Any = await service_tn.start(
_Cache.run_ctx,
mngr,
ctx_key,
ctx_exit,
)
except BaseException:
# If `run_ctx` (wrapping the acm's `__aenter__`)
@ -427,6 +466,11 @@ async def maybe_open_context(
raise taskc
else:
# XXX, cached-entry-path
(
_,
_,
ctx_exit,
) = _Cache.resources[ctx_key]
_Cache.users[ctx_key] += 1
user_registered = True
log.debug(
@ -445,24 +489,17 @@ async def maybe_open_context(
)
finally:
if lock.locked():
stats: trio.LockStatistics = lock.statistics()
owner: trio.Task|None = stats.owner
log.error(
f'Lock never released by last owner={owner!r} !?\n'
f'{stats}\n'
f'\n'
f'task={task!r}\n'
f'ctx_key={ctx_key!r}\n'
f'acm_func={acm_func}\n'
)
if user_registered:
# Serialize user registration and teardown under the same
# per-key lock so no entrant can acquire a resource after
# its final user has committed to exiting it.
with trio.CancelScope(shield=True):
await lock.acquire()
try:
_Cache.users[ctx_key] -= 1
if yielded is not _UnresolvedCtx:
# if no more consumers, teardown the client
# If no consumers remain, keep entrants queued
# until the cached context has completely exited.
if _Cache.users[ctx_key] <= 0:
log.debug(
f'De-allocating @acm-func entry\n'
@ -470,19 +507,39 @@ async def maybe_open_context(
f'acm_func={acm_func!r}\n'
)
# XXX: if we're cancelled we the entry may have never
# been entered since the nursery task was killed.
# _, no_more_users = _Cache.resources[ctx_key]
# XXX: if we're cancelled, the entry may
# have never been entered since the nursery
# task was killed.
entry = _Cache.resources.get(ctx_key)
if entry:
_, no_more_users = entry
(
_,
no_more_users,
ctx_exit,
) = entry
no_more_users.set()
maybe_lock = _Cache.locks.pop(
ctx_key,
None,
)
if maybe_lock is None:
assert ctx_exit is not None
await ctx_exit.done.wait()
exit_error = ctx_exit.error
# A queued entrant already holds a reference
# to this lock. Keep it registered until that
# task has acquired and released it.
stats = lock.statistics()
if not stats.tasks_waiting:
maybe_lock = _Cache.locks.get(ctx_key)
if maybe_lock is lock:
_Cache.locks.pop(ctx_key)
else:
log.error(
f'Resource lock for {ctx_key} ALREADY POPPED?'
f'Resource lock for {ctx_key} '
f'was replaced before teardown?'
)
finally:
lock.release()
if exit_error is not None:
# Always re-raise a regular `__aexit__()` error at the
# final consumer's context boundary.
raise exit_error

View File

@ -349,7 +349,10 @@ async def start_or_cancel(
# demote it to a `Cancelled`, losing the real error. The
# `isinstance` guard also avoids a `TypeError` when
# `rte.args[0]` isn't a `str`.
'child exited without calling' in rte.args[0]
rte.args[0] == (
'child exited without calling '
'task_status.started()'
)
):
# re-raises the in-flight `trio.Cancelled` IFF we're
# under effective cancellation; else a cheap no-op and

12
uv.lock
View File

@ -308,11 +308,11 @@ wheels = [
[[package]]
name = "idna"
version = "3.10"
version = "3.18"
source = { registry = "https://pypi.org/simple" }
sdist = { url = "https://files.pythonhosted.org/packages/f1/70/7703c29685631f5a7590aa73f1f1d3fa9a380e654b86af429e0934a32f7d/idna-3.10.tar.gz", hash = "sha256:12f65c9b470abda6dc35cf8e63cc574b1c52b11df2c86030af0ac09b01b13ea9", size = 190490, upload-time = "2024-09-15T18:07:39.745Z" }
sdist = { url = "https://files.pythonhosted.org/packages/cd/63/9496c57188a2ee585e0f1db071d75089a11e98aa86eb99d9d7618fc1edce/idna-3.18.tar.gz", hash = "sha256:ffb385a7e039654cef1ab9ef32c6fafe283c0c0467bba1d9029738ce4a14a848", size = 196711, upload-time = "2026-06-02T14:34:07.794Z" }
wheels = [
{ url = "https://files.pythonhosted.org/packages/76/c6/c88e154df9c4e1a2a66ccf0005a88dfb2650c1dffb6f5ce603dfbd452ce3/idna-3.10-py3-none-any.whl", hash = "sha256:946d195a0d259cbba61165e88e65941f16e9b36ea6ddb97f00452bae8b1287d3", size = 70442, upload-time = "2024-09-15T18:07:37.964Z" },
{ url = "https://files.pythonhosted.org/packages/1e/5e/d4e9f1a599fb8e573b7b87160658329fbf28d19eac2718f51fc3def3aa5a/idna-3.18-py3-none-any.whl", hash = "sha256:7f952cbe720b688055e3f87de14f5c3e5fdaa8bc3928985c4077ca689de849a2", size = 65455, upload-time = "2026-06-02T14:34:06.319Z" },
]
[[package]]
@ -900,11 +900,11 @@ wheels = [
[[package]]
name = "setuptools"
version = "82.0.1"
version = "83.0.0"
source = { registry = "https://pypi.org/simple" }
sdist = { url = "https://files.pythonhosted.org/packages/4f/db/cfac1baf10650ab4d1c111714410d2fbb77ac5a616db26775db562c8fab2/setuptools-82.0.1.tar.gz", hash = "sha256:7d872682c5d01cfde07da7bccc7b65469d3dca203318515ada1de5eda35efbf9", size = 1152316, upload-time = "2026-03-09T12:47:17.221Z" }
sdist = { url = "https://files.pythonhosted.org/packages/34/26/f5d29e25ffdb535afef2d35cdb55b325298f96debd670da4c325e08d70f4/setuptools-83.0.0.tar.gz", hash = "sha256:025bccbbf0fa05b6192bc64ae1e7b16e001fd6d6d4d5de03c97b1c1ade523bef", size = 1154254, upload-time = "2026-07-04T15:31:22.699Z" }
wheels = [
{ url = "https://files.pythonhosted.org/packages/9d/76/f789f7a86709c6b087c5a2f52f911838cad707cc613162401badc665acfe/setuptools-82.0.1-py3-none-any.whl", hash = "sha256:a59e362652f08dcd477c78bb6e7bd9d80a7995bc73ce773050228a348ce2e5bb", size = 1006223, upload-time = "2026-03-09T12:47:15.026Z" },
{ url = "https://files.pythonhosted.org/packages/5d/40/e1e72872c6354b306daef1703549e8e83b4d43cfea356311bf722a043752/setuptools-83.0.0-py3-none-any.whl", hash = "sha256:29b23c360f22f414dc7336bb39178cc7bcbf6021ed2733cde173f09dba19abb3", size = 1008090, upload-time = "2026-07-04T15:31:20.885Z" },
]
[[package]]