Compare commits

..

No commits in common. "b645f7fa8c82c4bcdec4f791e9e45a6c7b68b1cc" and "0ac9fa1c0cc43c29812b4cf5b1959850ea9b8c90" have entirely different histories.

59 changed files with 867 additions and 944 deletions

View File

@ -367,94 +367,9 @@ user's failing-test-first convention). The poll-based reap
fix in `_supervise.py` is UNCOMMITTED and likely SUPERSEDED
by the re-scoping — do NOT land it as-is.
## RESOLVED (2026-07-06): migrate everything, remove the API
The PAUSED re-assessment concluded decisively: rather than
re-scope `_reap_ria_portals` (or bolt any hack onto it), the
`run_in_actor()` API itself was REMOVED — its non-blocking
"result at teardown" semantic predates streaming and confused
more than it served. Every in-repo caller was migrated
per-file/-group (each its own commit, each gated):
- tests: `test_infected_asyncio` `test_runtime` `test_rpc`
`test_spawning` `test_pubsub` `test_registrar`
`test_cancellation` (3 groups) `test_advanced_streaming`.
- examples: 4 non-debugging + all 8 `debugging/` REPL scripts
(debugger suite byte-identical green, 28p/6s).
- docs: 8 rst pages + the `experimental/_pubsub` docstring.
Migration patterns (the `run_in_actor` shape -> successor):
- blocking result -> `to_actor.run(fn, an=an, ...)`
- fire-&-forget/forever -> bg `to_actor.run()` task in a local
`trio` task-nursery (or `start_actor`
+ bg `Portal.run()` when a portal
handle is needed)
- concurrent fan-out -> N bg `to_actor.run()` tasks / or
`gather_contexts([p.open_context(..)])`
- reap-all-error-collect -> the "collect don't cancel" pattern:
each one-shot catches + stashes its
`RemoteActorError`, group raised
after the task-nursery joins (see
`examples/debugging/multi_subactors.py`)
- mutual-rendezvous -> peers must OUTLIVE both dialogs:
`start_actor()` daemons + concurrent
`Portal.run()`s + explicit
`an.cancel()` (eager one-shot reap
races the slower peer's dial of the
winner's dead sockaddr; found via
`test_trynamic_trio` flake).
Semantic deltas (tests loosened accordingly):
- teardown-reap-all BEG-of-N is GONE: local task-nurseries are
cancel-on-first, raced siblings' `Cancelled`s are absorbed,
and the runtime's `collapse_eg()` unwraps every single-member
group at each actor boundary — a fully-raced nested tree
relays a bare (annotated) `RemoteActorError` chain.
- `test_multierror_fast_nursery` deleted (pure reap-stress);
`test_nested_multierrors` re-purposed as deep-tree
cancel-cascade stress w/ a race-tolerant shape walk.
Final excision (after zero callers remained): `run_in_actor()`,
`._cancel_after_result_on_exit`, `_reap_ria_portals()`,
`Portal._submit_for_result/._expect_result_ctx/
.wait_for_result()/.result()`, `exhaust_portal()`,
`cancel_on_completion()`, `NoResult` — net -402 lines. The
reap-hang class (unbounded `wait_for_result` in machinery
scope) dissolves structurally: the only result-wait left lives
in the caller's task inside its own cancel-scope; the
`d1fb4a1a` anti-hang guard test passes by construction. The
poll-vs-`proc_waiter` debate is moot as predicted.
## Follow-up sketch: `to_actor.open_one_shot()` (run-async parity)
If deferred-result parity is ever wanted, the design that needs
NO runtime coupling, NO returned `Portal` and NO cancel-relay
`trio.Event` machinery:
async with to_actor.open_one_shot(
fn, an=an, **kws,
) as one_shot:
... # concurrent caller work
val = await one_shot.wait() # optional; errors always
# propagate at scope exit
an `@acm` that opens a private task-nursery, `start_soon`s ONE
task running the existing blocking `run()` and stashes the
value in a slot + sets a done-`trio.Event` (a memo, not a
cancel relay). Cancellation = plain scope-cancel of the acm's
nursery (the parked `Portal.run()` unwinds via `Cancelled`, the
shielded `cancel_actor()` reap still runs); a child error
raises into the acm scope so an un-`wait()`ed one-shot can
never silently drop its error. i.e. the old reaper's job is
done by scoping, not machinery. ~40 lines, all in
`to_actor/_api.py`, zero `_supervise` involvement.
## Verification gate
- per-migration-commit module gates on `trio` (+ `mp_spawn`
spot-gates incl. `test_infected_asyncio` per the B2 lesson);
`tests/devx/test_debugger.py` for the REPL flows.
- full suite on `trio` + `mp_spawn` at branch tip + CI matrix
via draft PR #484.
- `tests/test_cancellation.py test_spawning.py test_local.py
test_rpc.py` on `trio` + `mp_spawn` + `mp_forkserver`
backends, then full suite; `tests/devx/test_debugger.py`
for risk 3.

View File

@ -1,81 +0,0 @@
---
model: claude-fable-5
service: claude
session: 6db64ac6-6986-4505-9343-df4ee31e67db
timestamp: 2026-07-06T17:28:18Z
git_ref: ad42871e
scope: code
substantive: true
raw_file: 20260706T172818Z_ad42871e_prompt_io.raw.md
---
## Prompt
Resuming the `drop_ria_nursery` (wkt) branch after the human's
review pause at `ba4af0ad`, with the directive:
> we need to resolve the remaining outstandings ... to finish
> this endeavour we ideally,
> - finish the ria-tn hoist into the machinery in `to_actor` if
> possible.
> - avoid supporting the `.run_in_actor()` non-blocking API
> going forward other then if we can re-wrap `to_actor.run()`
> per prior discussion with opus4.8.
> - get the test harness completely green locally and in CI -
> AND ESPECIALLY without any hangs due to the "the old
> spawn-backend reaper raced `soft_kill()`'s scope-cancel"
> - update all docs to reflect these changes once we're happy.
plus mid-flight additions:
> make sure both all the test and example code uses the naming
> `an: ActorNursery` for variable names
> where we at in terms of a potential hackaround design for
> supporting a to_actor.run_async() ... without requiring the
> same malarky ... with a `trio.Event` relay ... i'm totes fine
> with rejigging the iface to NOT return a Portal as well.
## Response summary
Completed the entire #477 `run_in_actor` removal on
`drop_ria_nursery` (16 commits this session), gated per-commit:
- ported the remaining `test_cancellation` groups
(`some_cancels_all`, `nested_multierrors`+`spawn_and_error`,
SIGINT + sync-sleep) — one-shots become concurrent
`to_actor.run()` tasks in local task-nurseries; BEG-shape
assertions loosened for cancel-on-first + the runtime's
`collapse_eg()` single-member unwrap (a fully-raced nested
tree relays a bare annotated `RemoteActorError` chain).
- fixed a pre-existing `UnboundLocalError` (`timeout` `match`
had no default arm for non-trio/MTF backends).
- ported `test_dynamic_pub_sub`, 4 non-debugging examples, all
8 `debugging/` examples (debugger suite byte-identical,
28p/6s; `multi_subactors` introduces the "collect don't
cancel" reap-all replacement pattern), 8 docs pages + the
`experimental/_pubsub` docstring.
- EXCISED the API + cluster: `run_in_actor`,
`_reap_ria_portals`, `_cancel_after_result_on_exit`,
`Portal._submit_for_result/_expect_result_ctx/
wait_for_result/result`, `exhaust_portal`,
`cancel_on_completion`, `NoResult` — net -402 lines. The
reap-hang class dissolves structurally (result-waits now only
in caller task-scope).
- found + fixed a real migration race: mutual-rendezvous peers
(`test_trynamic_trio`, `a_trynamic_first_scene.py`) flaked
because an eagerly-reaped one-shot dies while its peer still
dials the registry-resolved (dead) sockaddr — such peers now
pin lifetimes via `start_actor()` + concurrent `Portal.run()`
+ explicit `an.cancel()`.
- `an: ActorNursery` naming sweep across tests/examples (±82
lines, scoped renames, prose untouched).
- parked a `to_actor.open_one_shot()` design sketch (acm +
private task-nursery over blocking `run()`; done-Event as
memo not cancel-relay; no Portal) in the plan doc.
## Files changed
See commits `d01a2123..ad42871e` on `drop_ria_nursery`
(tests, examples, docs, `tractor/{runtime,spawn,to_actor,msg}`
+ `_exceptions/_context/experimental`).

View File

@ -1,39 +0,0 @@
---
model: claude-fable-5
service: claude
timestamp: 2026-07-06T17:28:18Z
git_ref: ad42871e
diff_cmd: git diff ba4af0ad..ad42871e
---
# Raw AI output (diff-ref mode)
This session's output spans the 16 migration/excision commits
`d01a2123..ad42871e` on `drop_ria_nursery`; per diff-ref mode
the verbatim content is reachable via the pointer below.
## Generated files
> `git diff ba4af0ad..ad42871e`
Commit-wise (each `Gate:`-footed msg documents its own module
gate):
- `d01a2123` port `test_some_cancels_all`
- `697c6152` fix unbound `timeout` (non-trio/MTF `match` arm)
- `fa8799d5` port `test_nested_multierrors`
- `f11754ce` port SIGINT + sync-sleep cancel tests
- `cb6202e3` port `test_dynamic_pub_sub`
- `d8af5f12` port non-debugging examples
- `a3057cb2` port debugging examples (+ `test_debugger`
nested-nurseries final-shape expectations)
- `d6bed7c4` port docs (8 rst pages)
- `07e1669e` fix stale `@pub` docstring example
- `2a59cefb` REMOVE `run_in_actor()` + the ria reap cluster
(net -402 lines)
- `a297a32a` fix mutual-rendezvous premature-reap race
- `ad42871e` `an: ActorNursery` naming sweep
Plan/design record updated in
`ai/conc-anal/ria_nursery_removal_plan.md` (RESOLVED section +
the `to_actor.open_one_shot()` follow-up sketch).

View File

@ -1,56 +0,0 @@
---
model: openai/gpt-5.6-sol
service: opencode
session: 76c5d31c-5a2f-4503-9b16-410ee7f4fab3
timestamp: 2026-08-19T18:46:40Z
git_ref: 481ba003
scope: code
substantive: true
raw_file: 20260819T184640Z_481ba003_prompt_io.raw.md
---
## Prompt
Rebase PR #484 onto final PR #481, migrate every affected one-shot call
to the new positional target API and continue through downstream tests,
examples and documentation review.
## Response summary
Converted stale target keyword calls to target partials so previously
named inputs remain explicit while placement/runtime controls stay
direct. Updated error expectations for local signature validation and
linked remote error propagation, then corrected docs which still
described the removed one-shot implementation. Linked spawning and
context lifecycle prose to the corresponding API methods and detailed
context guide.
## Files changed
- `docs/api/core.rst` - describe linked one-shot context execution.
- `docs/guide/rpc.rst` - update placement and target call semantics.
- `docs/guide/spawning.rst` - document positional target inputs.
- `examples/debugging/multi_nested_subactors_error_up_through_nurseries.py` - migrate nested actor target inputs.
- `examples/debugging/root_cancelled_but_child_is_in_tty_lock.py` - preserve named recursive target inputs with partials.
- `tests/test_advanced_streaming.py` - migrate streaming target inputs.
- `tests/test_cancellation.py` - migrate calls and tighten errors.
- `tests/test_infected_asyncio.py` - bind asyncio target options.
- `tests/test_rpc.py` - migrate RPC target argument binding.
- `tests/test_runtime.py` - preserve named runtime target inputs.
- `tests/test_spawning.py` - preserve named spawning target inputs.
## Human edits
The human selected the stack order and final PR #481 base, asked the
agent to continue after each diagnostic step and required a complete
commit plan after independently force-pushing the rebased history.
After reviewing the migration, the human required every formerly named
target input to remain visibly named through `functools.partial()`
rather than becoming positional. These were human-directed agent edits;
the human also required plain `start_actor()` and `open_context()`
references in the spawning and RPC guides to link to their API methods
and the detailed context guide, then clarified that `to_actor.run()`
already uses the full context API while `Portal.run()` should share
linked lifecycle machinery without necessarily delegating through
`Portal.open_context()` or adding a `Started` message. The human made
no direct source-line edits.

View File

@ -1,30 +0,0 @@
---
model: openai/gpt-5.6-sol
service: opencode
timestamp: 2026-08-19T18:46:40Z
git_ref: 481ba003
diff_cmd: git diff HEAD~1..HEAD
---
Migrate PR #484's downstream one-shot calls to PR #481's final
`tractor.to_actor.run()` contract after the stack rebase.
> `git diff HEAD~1..HEAD -- docs examples tests`
Pass target arguments positionally and bind target keyword-only inputs
with `functools.partial()`. Keep placement and runtime controls as
direct `to_actor.run()` keywords. Update the invalid-target-argument
test to expect local signature binding before actor startup and require
direct `RemoteActorError` propagation from linked one-shots.
Update API and guide prose to describe positional target inputs,
linked `Portal.open_context()` execution and per-child reaping instead
of the removed `Portal.run()` and target-`**kwargs` conventions.
Verification:
- core and migrated runtime batches: `97 passed`
- discovery and related lifecycle batch: `33 passed, 1 skipped`
- changed executable examples: `9 passed`
- mapped debugger cases: `12 passed, 6 skipped`
- Ruff, compilation and `git diff --check`: clean

View File

@ -37,6 +37,7 @@ Spawning actors
.. autoclass:: ActorNursery
:members: start_actor,
run_in_actor,
cancel,
cancel_called,
cancelled_caught
@ -46,22 +47,10 @@ Spawning actors
:meth:`ActorNursery.start_actor` (daemon actor + portal) is the
blessed spawning primitive; pair it with
``Portal.open_context()`` for SC-linked remote tasks.
One-shot task actors
--------------------
.. autofunction:: tractor.to_actor.run
.. note::
:func:`tractor.to_actor.run` (parlance of
``trio.to_thread.run_sync()`` and friends) is the
*convenience* one-shot — spawn, run a single task, block on
its result, reap — built entirely on
:meth:`ActorNursery.start_actor`, a linked
:meth:`Portal.open_context` call and per-child cancellation/reaping,
so don't design around it as the core model. It supersedes the
removed (legacy, non-blocking) ``ActorNursery.run_in_actor()``.
:meth:`ActorNursery.run_in_actor` is a *convenience* one-shot —
spawn, run a single task, auto-cancel after the result — slated
to be rebuilt as a high-level wrapper, so don't design around
it as the core model.
.. deprecated:: 0.1.0a6
@ -82,12 +71,14 @@ flowing back `exactly like trio`_.
:members: run,
run_from_ns,
open_stream_from,
wait_for_result,
cancel_actor,
chan
.. deprecated:: 0.1.0a6
The str-form ``Portal.run('mod.path', 'fn_name')`` warns;
``Portal.result()`` warns; use :meth:`Portal.wait_for_result`.
The str-form ``Portal.run('mod.path', 'fn_name')`` also warns;
pass a function *object* whose module is listed in the target's
``enable_modules``. ``Portal.channel`` is the legacy spelling
of :attr:`Portal.chan`.

View File

@ -76,8 +76,8 @@ Just flip the flag on :meth:`tractor.ActorNursery.start_actor`:
infect_asyncio=True,
)
The one-shot convenience ``tractor.to_actor.run()`` accepts the
same flag. The ``to_asyncio`` APIs may **only** be called from
The one-shot convenience ``ActorNursery.run_in_actor()`` accepts
the same flag. The ``to_asyncio`` APIs may **only** be called from
tasks inside an infected actor; calling them anywhere else raises
a loud ``RuntimeError``. You can introspect at runtime with
``tractor.current_actor().is_infected_aio()``.
@ -229,7 +229,7 @@ dialog, skip the channel ceremony and use
It schedules the fn as an ``asyncio.Task``, waits for completion
and hands the return value back to ``trio``; think of it as the
cross-loop sibling of ``tractor.to_actor.run()``. Errors and
cross-loop sibling of ``ActorNursery.run_in_actor()``. Errors and
cancellation are translated exactly as for channels.
Cross-loop errors and cancellation

View File

@ -64,13 +64,11 @@ What's going on here?
- three healthy actors are spawned as daemons via
:meth:`tractor.ActorNursery.start_actor`; left alone they'd
happily idle forever,
- a fourth actor runs ``assert_err()`` via a blocking
``tractor.to_actor.run()`` one-shot and promptly trips its
``assert 0``,
- a fourth actor runs ``assert_err()`` via ``.run_in_actor()`` and
promptly trips its ``assert 0``,
- the resulting ``AssertionError`` ships back over IPC as a
serialized error msg and re-raises *boxed* right at the call
inside the nursery block as a
:class:`tractor.RemoteActorError`,
serialized error msg and re-raises *boxed* inside the nursery
block as a :class:`tractor.RemoteActorError`,
- the nursery reacts like any ``trio`` nursery would: it cancels
the three healthy siblings (graceful runtime-cancel requests,
acks awaited), reaps all four processes, then re-raises,

View File

@ -15,8 +15,8 @@ a single `structured concurrency`_ (SC) scope over IPC.
:alt: sequence diagram of the context handshake msg flow
Pretty much everything else is (or is slated to be) built on this
one primitive: ``tractor.to_actor.run()`` is a convenience for
"spawn, run the lone task, await the result, tear down"; plain
one primitive: ``ActorNursery.run_in_actor()`` is a convenience
for "spawn, open a context, await the result, tear down"; plain
``Portal.run()`` RPC is planned to be re-implemented on top of it;
the multi-process debugger's tree-wide REPL lock rides one. Grok
this page and the rest of the library reads as convenience

View File

@ -119,16 +119,15 @@ Run a func in a process
Even a pool can be overkill; "run this one async func in a
subprocess and give me the result" is a one-liner via
:func:`tractor.to_actor.run`,
:meth:`tractor.ActorNursery.run_in_actor`,
.. literalinclude:: ../../examples/parallelism/single_func.py
:caption: examples/parallelism/single_func.py
:language: python
``to_actor.run()`` is a *convenience wrapper* — spawn an actor,
run exactly one task in it, block on and return its result, reap
— not the core spawning model (that's
:meth:`tractor.ActorNursery.start_actor` plus
``run_in_actor()`` is a *convenience wrapper* — spawn an actor, run
exactly one task in it, reap on result — not the core spawning
model (that's :meth:`tractor.ActorNursery.start_actor` plus
:meth:`tractor.Portal.open_context`; see :doc:`/guide/context`).
But for this fire-and-collect shape it's exactly the right amount
of typing.

View File

@ -80,36 +80,28 @@ One special namespace exists: ``'self'`` resolves to the remote
how internal machinery (cancel requests, registry ops) travels;
don't build your app on it.
One-shot subactors: ``to_actor.run()``
--------------------------------------
When a subactor's *entire job* is a single function call, skip
the portal plumbing with :func:`tractor.to_actor.run`: spawn,
run the lone task, return its result and reap the process — all
in one blocking call:
One-shot results: ``wait_for_result()``
---------------------------------------
A portal returned from
:meth:`~tractor.ActorNursery.run_in_actor` has exactly one
"main" task running remotely; that task's ``return`` value is
delivered as the portal's *final result*:
.. code:: python
from functools import partial
final = await tractor.to_actor.run(
partial(fib, n=10),
an=an,
)
portal = await an.run_in_actor(fib, n=10)
final = await portal.wait_for_result()
Semantics worth knowing:
- it blocks until the remote task returns, re-raising any
remote error in the usual boxed form right in the calling
task.
- "placement" is composable: ``an=`` spawns from an existing
actor-nursery, ``portal=`` reuses an already-running actor
(no spawn/reap, just a linked
:meth:`~tractor.Portal.open_context` call; see the
:doc:`context guide </guide/context>`), and passing neither
opens a private call-scoped nursery (booting the runtime if needed).
- concurrency composes the plain ``trio`` way: schedule
multiple ``run()`` calls into a local task nursery (see
``examples/parallelism/to_actor_one_shots.py``).
remote error in the usual boxed form.
- once resolved it's idempotent: later calls return the same
cached value.
- a *daemon* portal (from ``start_actor()``) has no main task,
so there's no final result to wait for: you'll get a warning
plus a ``NoResult`` sentinel. Results of individual daemon
calls come straight back from each ``await portal.run()``.
Pure RPC daemons: ``run_daemon()``
----------------------------------
@ -155,8 +147,7 @@ call tears down the entire sub-tree — SC, transitively.
When to graduate to ``Context``
-------------------------------
The :meth:`~tractor.Portal.run` method is great for one-shot,
request-response calls.
``portal.run()`` is great for one-shot, request-response calls.
Reach for :meth:`~tractor.Portal.open_context` with an
``@tractor.context`` endpoint as soon as you want:
@ -169,15 +160,10 @@ Reach for :meth:`~tractor.Portal.open_context` with an
:meth:`~tractor.Portal.cancel_actor` nukes the **entire**
remote runtime and its process.
:func:`tractor.to_actor.run` already enters the full
:meth:`~tractor.Portal.open_context` lifecycle. The older
:meth:`~tractor.Portal.run` path instead uses the ``Context`` returned
by the lower-level ``Actor.start_remote_task()`` directly, avoiding a
``Started`` handshake but owning less lifecycle machinery. A follow-up
should factor their shared linked-task lifecycle without requiring
``Portal.run()`` to delegate through the public context API or add
another wire message. Take the full tour in
:doc:`the context guide </guide/context>`.
In fact the source plans for ``Portal.run()`` itself to be
rebuilt on top of ``open_context()`` — contexts *are* the core
inter-actor protocol. Take the full tour in
:doc:`/guide/context`.
.. seealso::

View File

@ -91,34 +91,31 @@ somebody-ing:
What's going on here?
- :meth:`~tractor.ActorNursery.start_actor` forks off
- ``start_actor('frank', enable_modules=[__name__])`` forks off
a new process, boots a ``tractor`` runtime inside it, and
allows it to serve functions from the current module (see the
allowlist section below).
- each :meth:`~tractor.Portal.run` call schedules a *new* task in
- each ``await portal.run(...)`` schedules a *new* task in
frank's task tree and waits on its result — the full RPC story
lives in :doc:`/guide/rpc`.
- frank has no main task to complete, so without the final
:meth:`~tractor.Portal.cancel_actor` call the nursery block would
wait on him **forever**. Daemon lifetimes are *yours* to end;
that explicitness is the point.
``await portal.cancel_actor()`` the nursery block would wait
on him **forever**. Daemon lifetimes are *yours* to end; that
explicitness is the point.
``to_actor.run()``: quick one-shot parallelism
``run_in_actor()``: quick one-shot parallelism
----------------------------------------------
:func:`tractor.to_actor.run` is the convenience wrapper: spawn
an actor, run exactly one async function in it, block on the
result, then reap the process — the distributed sibling of
``trio.to_thread.run_sync()``.
:meth:`~tractor.ActorNursery.run_in_actor` is the convenience
wrapper: spawn an actor, run exactly one async function in it,
then reap the process as soon as the result arrives.
.. code:: python
async with (
tractor.open_nursery() as an,
trio.open_nursery() as tn,
):
async with tractor.open_nursery() as an:
portal = await an.run_in_actor(burn_cpu)
# burn rubber in the parent too...
tn.start_soon(burn_cpu)
total = await tractor.to_actor.run(burn_cpu, an=an)
await burn_cpu()
total = await portal.wait_for_result()
A few details worth knowing:
@ -126,52 +123,43 @@ A few details worth knowing:
``name='something_cuter'``.
- the function's module is auto-added to the child's
``enable_modules`` allowlist.
- target arguments are positional; use ``functools.partial()``
to bind target keyword arguments. Keywords passed directly to
``run()`` configure actor placement and spawning.
- the call blocks until the result (or error) lands and the
child is *auto-cancelled* (reaped) right after — so remote
errors raise directly in your calling task (causality_ is
paramount!).
- "placement" composes: ``an=`` spawns from a caller-managed
actor-nursery, ``portal=`` reuses an already-running actor
(no spawn/reap), and passing neither opens a private
call-scoped nursery (booting the runtime if needed).
- extra ``**kwargs`` are forwarded to the function itself.
- the child is *auto-cancelled* once its "main" result lands;
at nursery exit these run-once children are always reaped
first (causality_ is paramount!).
.. note::
:func:`tractor.to_actor.run` is a convenience, **not** the core
model — it's built *entirely* on
:meth:`~tractor.ActorNursery.start_actor` plus a linked
:meth:`~tractor.Portal.open_context` call and per-child
cancellation/reaping. Teach your fingers to use it for quick
fire-and-collect parallelism — think a per-function trio-parallel_
style one-shot — and reach for
:meth:`~tractor.ActorNursery.start_actor` plus
:meth:`~tractor.Portal.open_context` for anything long-lived,
stateful or streaming; see :doc:`/guide/context`.
``run_in_actor()`` is a convenience, **not** the core model.
The source literally marks it for an eventual rebuild as
a thin "hilevel" wrapper on top of
:meth:`~tractor.Portal.open_context` (the modern inter-actor
task API). Teach your fingers to use it for quick
fire-and-collect parallelism — think a per-function
trio-parallel_ style one-shot — and reach for
``start_actor()`` + ``open_context()`` for anything
long-lived, stateful or streaming
(:doc:`/guide/context`).
Actor lifetimes and teardown order
----------------------------------
So we have two lifetime flavors:
- **one-shot** (``to_actor.run()``): lives exactly as long as
- **run-once** (``run_in_actor()``): lives exactly as long as
its single task; reaped the moment its result (or error)
arrives back in the (blocking) call.
- **daemon** (:meth:`~tractor.ActorNursery.start_actor`): lives
until *someone* cancels it — an explicit
:meth:`~tractor.Portal.cancel_actor`, a bulk
:meth:`~tractor.ActorNursery.cancel`, or the one-cancels-all
strategy kicking in on error.
arrives.
- **daemon** (``start_actor()``): lives until *someone* cancels
it — an explicit ``await portal.cancel_actor()``, a bulk
``await an.cancel()``, or the one-cancels-all strategy kicking
in on error.
On a clean exit of the nursery block the teardown order is:
1. one-shot actors never make it to nursery exit: each is
reaped inside its own ``to_actor.run()`` call, any error
raising immediately in the calling task so your code
(acting as supervisor) gets first crack at handling it.
2. the nursery then waits on daemon actors — **indefinitely**.
If you spawned a daemon, you own its lifetime.
1. the nursery waits on every run-once actor's final result;
any errors from these are raised immediately so your code
(acting as supervisor) gets first crack at handling them.
2. then it waits on daemon actors — **indefinitely**. If you
spawned a daemon, you own its lifetime.
When a child *is* cancelled, teardown is graceful-first per SC
discipline: the runtime sends an IPC cancel request and gives

View File

@ -43,20 +43,24 @@ Run it::
What's going on here?
- ``trio.run(main)`` starts the **root actor**; the ``tractor``
runtime boots *implicitly* inside ``tractor.to_actor.run()``
runtime boots *implicitly* inside ``tractor.open_nursery()``
whenever it isn't already up. No special entrypoint, no
framework takeover - it's just a ``trio`` app,
- inside ``main()`` a *subactor* is spawned via
``tractor.to_actor.run()`` and told to run exactly one
``ActorNursery.run_in_actor()`` and told to run exactly one
function: ``cellar_door()``,
- you get back a ``Portal``: your handle for invoking tasks in
the new process's (separate!) memory domain. We lean on it
much harder in the next section,
- the subactor, *some_linguist*, boots a fresh ``trio.run()`` in
a **new process** and executes ``cellar_door()`` as its *main
task* (note the child proving it is *not* the root with
``tractor.is_root_process()``), then ships the return value
back over IPC,
- the call *blocks* until that final result arrives, then
returns it - causality is preserved: your task only proceeds
once the child is *done*, dead, and reaped.
- the parent grabs that *final result* with
``await portal.wait_for_result()``, much like you'd expect
from a "future" - except causality is preserved: the nursery
block only exits once the child is *done*, dead, and reaped.
.. margin:: Just need a worker pool?
@ -67,22 +71,19 @@ What's going on here?
.. note::
``to_actor.run()`` (parlance of ``trio.to_thread`` and
friends) is the *convenience* wrapper: one-shot
``run_in_actor()`` is the *convenience* wrapper: one-shot
spawn-run-reap semantics for when a subactor's entire job is
a single function call. The core primitives are
``ActorNursery.start_actor()`` (next up) — which hands you
a ``Portal``, your handle for invoking tasks in the new
process's (separate!) memory domain — paired with
``ActorNursery.start_actor()`` (next up) paired with
``Portal.open_context()`` for full, SC-linked cross-actor
dialogs - see :doc:`/guide/context`.
Daemon actors and RPC
---------------------
A ``to_actor.run()`` one-shot subactor terminates when its lone
task returns. But often you want long-lived *daemon* actors
instead: spawned once, then serving (allowlisted) RPC requests
until told otherwise. That's ``start_actor()``:
A ``run_in_actor()``-spawned actor terminates when its main task
returns. But often you want long-lived *daemon* actors instead:
spawned once, then serving (allowlisted) RPC requests until told
otherwise. That's ``start_actor()``:
.. literalinclude:: ../../examples/actor_spawning_and_causality_with_daemon.py
:caption: examples/actor_spawning_and_causality_with_daemon.py
@ -90,9 +91,9 @@ until told otherwise. That's ``start_actor()``:
Two lifetime rules to internalize:
- a ``to_actor.run()`` one-shot actor lives exactly as long as
its lone task; the call blocks until that function (and thus
the process) completes,
- a ``run_in_actor()`` actor lives exactly as long as its main
task; the nursery waits for that function (and thus the
process) to complete before unblocking,
- a ``start_actor()`` actor *lives forever* - an RPC daemon the
nursery will happily wait on **indefinitely** - until some
task explicitly cancels it via ``Portal.cancel_actor()`` (as

View File

@ -21,34 +21,22 @@ async def main():
"""Main tractor entry point, the "master" process (for now
acts as the "director").
"""
async with tractor.open_nursery() as an:
async with tractor.open_nursery() as n:
print("Alright... Action!")
# both actors wait on (then dial!) the *other*, so each
# must outlive both hellos: spawn as daemons, run the
# hellos concurrently, reap only once both complete.
portals: dict[str, tractor.Portal] = {
name: await an.start_actor(
name,
enable_modules=[__name__],
)
for name in ('donny', 'gretchen')
}
async def run_and_print(name: str, other_actor: str):
print(
await portals[name].run(
donny = await n.run_in_actor(
say_hello,
other_actor=other_actor,
name='donny',
# arguments are always named
other_actor='gretchen',
)
gretchen = await n.run_in_actor(
say_hello,
name='gretchen',
other_actor='donny',
)
async with trio.open_nursery() as tn:
tn.start_soon(run_and_print, 'donny', 'gretchen')
tn.start_soon(run_and_print, 'gretchen', 'donny')
await an.cancel()
print(await gretchen.wait_for_result())
print(await donny.wait_for_result())
print("CUTTTT CUUTT CUT!!! Donny!! You're supposed to say...")

View File

@ -10,14 +10,17 @@ async def cellar_door():
async def main():
"""The main ``tractor`` routine.
"""
# spawn a subactor, run ``cellar_door()`` as its lone task,
# block until its result arrives and the subactor is reaped.
print(
await tractor.to_actor.run(
async with tractor.open_nursery() as n:
portal = await n.run_in_actor(
cellar_door,
name='some_linguist',
)
)
# The ``async with`` will unblock here since the 'some_linguist'
# actor has completed its main task ``cellar_door``.
print(await portal.wait_for_result())
if __name__ == '__main__':

View File

@ -12,9 +12,9 @@ async def movie_theatre_question():
async def main():
"""The main ``tractor`` routine.
"""
async with tractor.open_nursery() as an:
async with tractor.open_nursery() as n:
portal = await an.start_actor(
portal = await n.start_actor(
'frank',
# enable the actor to run funcs from this current module
enable_modules=[__name__],

View File

@ -15,9 +15,9 @@ async def stream_forever() -> AsyncIterator[int]:
async def main():
async with tractor.open_nursery() as an:
async with tractor.open_nursery() as n:
portal = await an.start_actor(
portal = await n.start_actor(
'donny',
enable_modules=[__name__],
)

View File

@ -1,5 +1,3 @@
from functools import partial
import trio
import tractor
@ -23,39 +21,26 @@ async def breakpoint_forever():
async def spawn_until(depth=0):
""""A nested nursery that triggers another ``NameError``.
"""
async with (
tractor.open_nursery() as an,
trio.open_nursery() as tn,
):
async with tractor.open_nursery() as n:
if depth < 1:
tn.start_soon(
partial(
tractor.to_actor.run,
breakpoint_forever,
an=an,
)
)
await n.run_in_actor(breakpoint_forever)
p = await n.run_in_actor(
name_error,
name='name_error'
)
await trio.sleep(0.5)
# rx and propagate error from child
await tractor.to_actor.run(
name_error,
an=an,
name='name_error',
)
await p.result()
else:
# recusrive call to spawn another process branching layer of
# the tree; blocks (up) each level until the leaf's
# `name_error` relays through.
# the tree
depth -= 1
await tractor.to_actor.run(
partial(
await n.run_in_actor(
spawn_until,
depth=depth,
),
an=an,
name=f'spawn_until_{depth}',
)
@ -80,37 +65,34 @@ async def main():
python -m tractor._child --uid ('spawn_until_0', 'de918e6d ...)
"""
async with (
tractor.open_nursery(
async with tractor.open_nursery(
debug_mode=True,
loglevel='pdb',
) as an,
trio.open_nursery() as tn,
):
# spawn both spawner trees as concurrent one-shots; the
# first tree's (relayed) error cancels the other.
tn.start_soon(
partial(
tractor.to_actor.run,
partial(
) as n:
# spawn both actors
portal = await n.run_in_actor(
spawn_until,
depth=3,
),
an=an,
name='spawner0',
)
)
tn.start_soon(
partial(
tractor.to_actor.run,
partial(
portal1 = await n.run_in_actor(
spawn_until,
depth=4,
),
an=an,
name='spawner1',
)
)
# TODO: test this case as well where the parent don't see
# the sub-actor errors by default and instead expect a user
# ctrl-c to kill the root.
with trio.move_on_after(3):
await trio.sleep_forever()
# gah still an issue here.
await portal.result()
# should never get here
await portal1.result()
if __name__ == '__main__':

View File

@ -15,12 +15,12 @@ async def name_error():
async def spawn_error():
""""A nested nursery that triggers another ``NameError``.
"""
async with tractor.open_nursery() as an:
return await tractor.to_actor.run(
async with tractor.open_nursery() as n:
portal = await n.run_in_actor(
name_error,
an=an,
name='name_error_1',
)
return await portal.result()
async def main():
@ -38,36 +38,29 @@ async def main():
- root actor should then fail on assert
- program termination
"""
async with (
tractor.open_nursery(
async with tractor.open_nursery(
debug_mode=True,
loglevel='devx',
) as an,
trio.open_nursery() as tn,
):
# spawn both actors..
portal = await an.start_actor(
'name_error',
enable_modules=[__name__],
)
portal1 = await an.start_actor(
'spawn_error',
enable_modules=[__name__],
)
) as n:
# ..and bg-schedule their erroring tasks.
tn.start_soon(portal.run, name_error)
tn.start_soon(portal1.run, spawn_error)
# yield to the bg tasks so both RPC requests are
# submitted (and start crashing) before the root's own
# error below (the legacy `run_in_actor()` submitted
# in-line with each spawn).
await trio.sleep(0.5)
# spawn both actors
portal = await n.run_in_actor(
name_error,
name='name_error',
)
portal1 = await n.run_in_actor(
spawn_error,
name='spawn_error',
)
# trigger a root actor error
assert 0
# attempt to collect results (which raises error in parent)
# still has some issues where the parent seems to get stuck
await portal.result()
await portal1.result()
if __name__ == '__main__':
trio.run(main)

View File

@ -17,12 +17,12 @@ async def name_error():
async def spawn_error():
""""A nested nursery that triggers another ``NameError``.
"""
async with tractor.open_nursery() as an:
return await tractor.to_actor.run(
async with tractor.open_nursery() as n:
portal = await n.run_in_actor(
name_error,
an=an,
name='name_error_1',
)
return await portal.result()
async def main():
@ -36,39 +36,17 @@ async def main():
`-python -m tractor._child --uid ('spawn_error', '52ee14a5 ...)
`-python -m tractor._child --uid ('name_error', '3391222c ...)
"""
errors: list[BaseException] = []
async with tractor.open_nursery(
debug_mode=True,
# loglevel='runtime',
) as an:
) as n:
async def run_and_collect(fn):
'''
One-shot whose (boxed) error is stashed instead of
raised so a sibling's crash never cancels the others
before they've had their own debugger sessions (the
"collect all errors" the legacy `run_in_actor()` API
did implicitly at nursery teardown).
'''
try:
await tractor.to_actor.run(fn, an=an)
except tractor.RemoteActorError as rae:
errors.append(rae)
# Spawn all one-shot task actors, collecting (vs.
# raising) their errors.
async with trio.open_nursery() as tn:
tn.start_soon(run_and_collect, breakpoint_forever)
tn.start_soon(run_and_collect, name_error)
tn.start_soon(run_and_collect, spawn_error)
if errors:
raise BaseExceptionGroup(
'multi_subactors errored!',
errors,
)
# Spawn both actors, don't bother with collecting results
# (would result in a different debugger outcome due to parent's
# cancellation).
await n.run_in_actor(breakpoint_forever)
await n.run_in_actor(name_error)
await n.run_in_actor(spawn_error)
if __name__ == '__main__':

View File

@ -21,8 +21,8 @@ async def main() -> None:
async with tractor.open_nursery(
debug_mode=True,
) as an:
portal = await an.start_actor(
) as n:
portal = await n.start_actor(
'ctx_child',
# XXX: we don't enable the current module in order

View File

@ -6,14 +6,14 @@ async def die():
async def main():
async with tractor.open_nursery() as an:
async with tractor.open_nursery() as tn:
debug_actor = await an.start_actor(
debug_actor = await tn.start_actor(
'debugged_boi',
enable_modules=[__name__],
debug_mode=True,
)
crash_boi = await an.start_actor(
crash_boi = await tn.start_actor(
'crash_boi',
enable_modules=[__name__],
# debug_mode=True,

View File

@ -1,5 +1,3 @@
from functools import partial
import trio
import tractor
@ -12,17 +10,15 @@ async def name_error():
async def spawn_until(depth=0):
""""A nested nursery that triggers another ``NameError``.
"""
async with tractor.open_nursery() as an:
async with tractor.open_nursery() as n:
if depth < 1:
await tractor.to_actor.run(name_error, an=an)
# await n.run_in_actor('breakpoint_forever', breakpoint_forever)
await n.run_in_actor(name_error)
else:
depth -= 1
await tractor.to_actor.run(
partial(
await n.run_in_actor(
spawn_until,
depth=depth,
),
an=an,
name=f'spawn_until_{depth}',
)
@ -41,37 +37,28 @@ async def main():
python -m tractor._child --uid ('name_error', '6c2733b8 ...)
'''
async with (
tractor.open_nursery(
async with tractor.open_nursery(
debug_mode=True,
enable_transports=['uds'], # TODO, apss this via osenv?
loglevel='devx', # XXX, required for test!
) as an,
trio.open_nursery() as tn,
):
# spawn the deeper tree in the bg..
tn.start_soon(
partial(
tractor.to_actor.run,
partial(
spawn_until,
depth=1,
),
an=an,
name='spawner1',
)
)
) as n:
# ..while blocking on the shallow (faster to fail) tree
# whose propagated error triggers nursery cancellation.
await tractor.to_actor.run(
partial(
# spawn both actors
portal = await n.run_in_actor(
spawn_until,
depth=0,
),
an=an,
name='spawner0',
)
portal1 = await n.run_in_actor(
spawn_until,
depth=1,
name='spawner1',
)
# nursery cancellation should be triggered due to propagated
# error from child.
await portal.result()
await portal1.result()
if __name__ == '__main__':

View File

@ -13,24 +13,17 @@ async def main():
simultaneously.
'''
async with (
tractor.open_nursery(
async with tractor.open_nursery(
debug_mode=True,
# loglevel='debug' # ?XXX required?
) as an,
trio.open_nursery() as tn,
):
# spawn the actor..
portal = await an.start_actor(
'key_error',
enable_modules=[__name__],
)
) as n:
# spawn both actors
portal = await n.run_in_actor(key_error)
print(
f'Child is up @ {portal.chan.aid.reprol()}'
)
# ..then schedule its erroring task in the bg while the
# root blocks below.
tn.start_soon(portal.run, key_error)
# XXX: originally a bug caused by this is where root would enter
# the debugger and clobber the tty used by the repl even though

View File

@ -74,11 +74,11 @@ async def cancelled_before_pause(
async def main():
async with tractor.open_nursery(
debug_mode=True,
) as an:
await tractor.to_actor.run(
) as n:
portal: tractor.Portal = await n.run_in_actor(
cancelled_before_pause,
an=an,
)
await portal.wait_for_result()
# ensure the same works in the root actor!
await pm_on_cancelled()

View File

@ -58,8 +58,8 @@ async def main():
debug_mode=True,
enable_transports=[tpt],
loglevel='devx',
) as an:
p = await an.start_actor(
) as n:
p = await n.start_actor(
'bp_boi',
enable_modules=[__name__],
)

View File

@ -17,14 +17,12 @@ async def main():
async with tractor.open_nursery(
debug_mode=True,
loglevel='cancel',
) as an:
) as n:
# parks awaiting a result which only arrives once the
# user quits (`BdbQuit`s) the child's REPL loop.
await tractor.to_actor.run(
portal = await n.run_in_actor(
breakpoint_forever,
an=an,
)
await portal.wait_for_result()
if __name__ == '__main__':

View File

@ -12,12 +12,16 @@ async def main():
) as an:
# TODO: ideally the REPL arrives at this frame in the parent,
# ABOVE the @api_frame of `to_actor.run()` ..
# ABOVE the @api_frame of `Portal.run_in_actor()` (which
# should eventually not even be a portal method ... XD)
# await tractor.pause()
p: tractor.Portal = await an.run_in_actor(name_error)
# the one-shot blocks on the subactor's result so the
# boxed `NameError` raises right here.
await tractor.to_actor.run(name_error, an=an)
# with this style, should raise on this line
await p.wait_for_result()
# with this alt style should raise at `open_nusery()`
# return await p.wait_for_result()
if __name__ == '__main__':

View File

@ -90,7 +90,7 @@ async def main() -> None:
# TODO: 3 sub-actor usage cases:
# -[x] via a `.open_context()`
# -[ ] via a `to_actor.run()` call
# -[ ] via a `.run_in_actor()` call
# -[ ] via a `.run()`
# -[ ] via a `.to_thread.run_sync()` in subactor
async with p.open_context(

View File

@ -50,8 +50,8 @@ async def trio_to_aio_echo_server(
async def main():
async with tractor.open_nursery() as an:
p = await an.start_actor(
async with tractor.open_nursery() as n:
p = await n.start_actor(
'aio_server',
enable_modules=[__name__],
infect_asyncio=True,

View File

@ -29,9 +29,9 @@ async def main() -> None:
))
await proc.wait()
# await trio.sleep_forever()
# async with tractor.open_nursery() as an:
# async with tractor.open_nursery() as n:
# portal = await an.start_actor(
# portal = await n.start_actor(
# 'rpc_server',
# enable_modules=[__name__],
# )

View File

@ -55,7 +55,7 @@ async def worker_pool(workers=4):
Yes, the workers stay alive (and ready for work) until you close
the context.
"""
async with tractor.open_nursery() as an:
async with tractor.open_nursery() as tn:
portals = []
snd_chan, recv_chan = trio.open_memory_channel(len(PRIMES))
@ -65,7 +65,7 @@ async def worker_pool(workers=4):
# this starts a new sub-actor (process + trio runtime) and
# stores it's "portal" for later use to "submit jobs" (ugh).
portals.append(
await an.start_actor(
await tn.start_actor(
f'worker_{i}',
enable_modules=[__name__],
)
@ -80,10 +80,10 @@ async def worker_pool(workers=4):
async def send_result(func, value, portal):
await snd_chan.send((value, await portal.run(func, n=value)))
async with trio.open_nursery() as tn:
async with trio.open_nursery() as n:
for value, portal in zip(sequence, itertools.cycle(portals)):
tn.start_soon(
n.start_soon(
send_result,
worker_func,
value,
@ -98,7 +98,7 @@ async def worker_pool(workers=4):
yield _map
# tear down all "workers" on pool close
await an.cancel()
await tn.cancel()
async def main():

View File

@ -25,15 +25,17 @@ async def burn_cpu():
async def main():
async with trio.open_nursery() as tn:
async with tractor.open_nursery() as n:
portal = await n.run_in_actor(burn_cpu)
# burn rubber in the parent too
tn.start_soon(burn_cpu)
await burn_cpu()
# run the same func as the lone task in a subactor,
# block on (and collect) its result
pid = await tractor.to_actor.run(burn_cpu)
# wait on result from target function
pid = await portal.wait_for_result()
# end of nursery block
print(f"Collected subproc {pid}")

View File

@ -7,20 +7,19 @@ async def assert_err():
async def main():
async with tractor.open_nursery() as an:
async with tractor.open_nursery() as n:
real_actors = []
for i in range(3):
real_actors.append(await an.start_actor(
real_actors.append(await n.start_actor(
f'actor_{i}',
enable_modules=[__name__],
))
# run one one-shot task actor that will fail immediately;
# its error raises right here in the caller's task..
await tractor.to_actor.run(assert_err, an=an)
# start one actor that will fail immediately
await n.run_in_actor(assert_err)
# ..as a ``RemoteActorError`` containing an ``AssertionError``
# and all the other actors have been cancelled
# should error here with a ``RemoteActorError`` containing
# an ``AssertionError`` and all the other actors have been cancelled
if __name__ == '__main__':

View File

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

View File

@ -849,36 +849,43 @@ def test_multi_nested_subactors_error_through_nurseries(
break
# boxed source errors
#
# NB post-#477 (`to_actor.run()` one-shots in local
# task-nurseries) the final relay is the LAST-released
# (leaf) REPL's error chain: it wins each level's
# relay-vs-cancel race so every level's single-member
# group gets unwrapped by the runtime's `collapse_eg()`
# (annotated at each actor boundary) while the sibling
# tree ('spawner1') is cancelled + absorbed. The legacy
# `run_in_actor()` teardown-reap instead grouped BOTH the
# `name_error` and bp-quit chains into the final dump
# (the previously-unexplained "extra" patterns).
expect_patts: list[str] = [
"NameError: name 'doggypants' is not defined",
"tractor._exceptions.RemoteActorError:",
"('name_error'",
# each level's unwrapped-single-member-group
# annotation + the first-level subtree's boundary
# footer.
"( ^^^ this exc was collapsed from a group ^^^ )",
"------ ('spawner0'",
# first level subtrees
# "tractor._exceptions.RemoteActorError: ('spawner0'",
"src_uid=('spawner0'",
# "tractor._exceptions.RemoteActorError: ('spawner1'",
# propagation of errors up through nested subtrees
# "tractor._exceptions.RemoteActorError: ('spawn_until_0'",
# "tractor._exceptions.RemoteActorError: ('spawn_until_1'",
# "tractor._exceptions.RemoteActorError: ('spawn_until_2'",
# ^-NOTE-^ old RAE repr, new one is below with a field
# showing the src actor's uid.
"src_uid=('spawn_until_2'",
]
# XXX, I HAVE NO IDEA why these patts only show on the
# `trio`-spawner but it seems to have something to do with
# what gets dumped in prior-prompt latches somehow??
# TODO for claude, explain and or work through how this is
# happening but ONLY WHEN RUN FROM THE TEST, bc when i try to
# run the test script manually the correct output ALWAYS seems
# to be in the last `str(child.before.decode())` output !?!?
if (
not is_forking_spawner
and
last_send_char == 'q'
):
expect_patts += [
# expect the pdb-quit exc relayed from the leaf's
# bp-loop child.
# expect the pdb-quit exc.
"bdb.BdbQuit",
"src_uid=('breakpoint_forever'",
# BUT WHY these dude!?
"src_uid=('spawn_until_0'",
"relay_uid=('spawn_until_1'",
]
assert_before(

View File

@ -46,9 +46,9 @@ async def test_reg_then_unreg(
async with tractor.open_nursery(
registry_addrs=[reg_addr],
) as an:
) as n:
portal = await an.start_actor('actor', enable_modules=[__name__])
portal = await n.start_actor('actor', enable_modules=[__name__])
uid = portal.channel.aid.uid
async with tractor.get_registry(reg_addr) as aportal:
@ -62,7 +62,7 @@ async def test_reg_then_unreg(
# XXX: can we figure out what the listen addr will be?
assert sockaddrs
await an.cancel() # tear down nursery
await n.cancel() # tear down nursery
await trio.sleep(0.1)
assert uid not in aportal.actor._registry
@ -89,9 +89,9 @@ async def test_reg_then_unreg_maddr(
async with tractor.open_nursery(
registry_addrs=[maddr_str],
) as an:
) as n:
portal = await an.start_actor(
portal = await n.start_actor(
'actor_maddr',
enable_modules=[__name__],
)
@ -105,7 +105,7 @@ async def test_reg_then_unreg_maddr(
sockaddrs = actor._registry[uid]
assert sockaddrs
await an.cancel()
await n.cancel()
await trio.sleep(0.1)
assert uid not in aportal.actor._registry
@ -152,37 +152,27 @@ async def test_trynamic_trio(
for the directed subs.
'''
async with tractor.open_nursery() as an:
async with (
tractor.open_nursery() as n,
trio.open_nursery() as tn,
):
print("Alright... Action!")
# donny + gretchen each wait on (then dial!) the *other*, so
# both actors must OUTLIVE both hellos: spawn as daemons and
# only reap after both tasks complete. NB a pair of eagerly
# reaped `to_actor.run()` one-shots races: the first to
# finish dies while the other may still be dialing its
# registry-resolved (now dead) sockaddr -> conn-refused.
portals: dict[str, tractor.Portal] = {
name: await an.start_actor(
name,
enable_modules=[__name__],
)
for name in ('donny', 'gretchen')
}
# donny + gretchen each wait on the *other* to register, so
# they must run CONCURRENTLY — schedule both one-shots into a
# local task-nursery (was two non-blocking `run_in_actor()`s).
async def _direct(this_name: str, other_actor: str):
res = await portals[this_name].run(
res = await tractor.to_actor.run(
ria_fn,
an=n,
other_actor=other_actor,
reg_addr=reg_addr,
name=this_name,
)
print(res)
async with trio.open_nursery() as tn:
tn.start_soon(_direct, 'donny', 'gretchen')
tn.start_soon(_direct, 'gretchen', 'donny')
# both hellos have completed; reap the thespians.
await an.cancel()
print("CUTTTT CUUTT CUT!!?! Donny!! You're supposed to say...")

View File

@ -3,7 +3,6 @@ Advanced streaming patterns using bidirectional streams and contexts.
'''
from collections import Counter
from functools import partial
import itertools
import platform
from typing import Type
@ -174,8 +173,8 @@ def test_dynamic_pub_sub(
# test. Picked backend-aware: under `trio` backend spawn is
# cheap (~1s for `cpus` actors) but fork-based backends pay
# a per-spawn cost (forkserver round-trip + IPC peer-handshake)
# that can stack up over the `cpus - 1` one-shot
# (`to_actor.run()`) spawns — especially on UDS under cross-pytest contention
# that can stack up over `cpus - 1` sequential `n.run_in_actor()`
# calls — especially on UDS under cross-pytest contention
# (#451 / #452). 4s was flaking right at the edge under fork
# backends — bumped to 8s with diag-snapshot-on-timeout via
# `fail_after_w_trace` so a borderline run still fails loud
@ -215,59 +214,33 @@ def test_dynamic_pub_sub(
f'enter `fail_after_w_trace({fail_after_s})` scope'
)
try:
async with (
tractor.open_nursery(
async with tractor.open_nursery(
registry_addrs=[reg_addr],
debug_mode=debug_mode,
) as an,
# bg-schedules the forever-streaming
# one-shots below; the user-cancel raise
# cancels them all, each reaping its
# subactor via `to_actor.run()`'s
# (shielded) `Portal.cancel_actor()`.
trio.open_nursery() as tn,
):
) as n:
test_log.cancel(
'test_dynamic_pub_sub: '
'actor nursery opened'
)
# name of this actor will be same as target func
tn.start_soon(
partial(
tractor.to_actor.run,
publisher,
an=an,
)
)
await n.run_in_actor(publisher)
for i, sub in zip(
range(cpus - 2),
itertools.cycle(_registry.keys())
):
tn.start_soon(
partial(
tractor.to_actor.run,
partial(
await n.run_in_actor(
consumer,
subs=[sub],
),
an=an,
name=f'consumer_{sub}',
)
subs=[sub],
)
# make one dynamic subscriber
tn.start_soon(
partial(
tractor.to_actor.run,
partial(
await n.run_in_actor(
consumer,
subs=list(_registry.keys()),
),
an=an,
name='consumer_dynamic',
)
subs=list(_registry.keys()),
)
# block until "cancelled by user"
@ -374,10 +347,10 @@ def test_reqresp_ontopof_streaming():
timeout = 4
with trio.move_on_after(timeout):
async with tractor.open_nursery() as an:
async with tractor.open_nursery() as n:
# name of this actor will be same as target func
portal = await an.start_actor(
portal = await n.start_actor(
'dual_tasks',
enable_modules=[__name__]
)
@ -440,9 +413,9 @@ def test_sigint_both_stream_types():
async def main():
with trio.fail_after(timeout):
async with tractor.open_nursery() as an:
async with tractor.open_nursery() as n:
# name of this actor will be same as target func
portal = await an.start_actor(
portal = await n.start_actor(
'2_way',
enable_modules=[__name__]
)
@ -555,8 +528,8 @@ def test_local_task_fanout_from_stream(
async with tractor.open_nursery(
debug_mode=debug_mode,
) as an:
p: tractor.Portal = await an.start_actor(
) as tn:
p: tractor.Portal = await tn.start_actor(
'inf_streamer',
enable_modules=[__name__],
)

View File

@ -104,7 +104,7 @@ async def do_nuthin():
[
# expected to be thrown in assert_err
({}, AssertionError),
# argument mismatch rejected locally before spawn
# argument mismatch raised in _invoke()
({'unexpected': 10}, TypeError)
],
ids=['no_args', 'unexpected_args'],
@ -126,38 +126,52 @@ def test_remote_error(
async def main():
async with tractor.open_nursery(
registry_addrs=[reg_addr],
) as an:
) as nursery:
# `to_actor.run()` blocks on the one-shot's result and
# raises the remote error directly here in the caller's
# task. Invalid target args fail local signature binding
# before any one-shot actor is spawned.
# task (a bad-arg `TypeError` likewise relays as a
# `RemoteActorError`).
try:
await tractor.to_actor.run(
partial(assert_err, **args),
an=an,
assert_err,
an=nursery,
name='errorer',
**args
)
except tractor.RemoteActorError as err:
assert err.boxed_type == errtype
print("Look Maa that actor failed hard, hehh")
raise
# Invalid args never cross the process boundary.
# ensure boxed errors
if args:
with pytest.raises(errtype):
trio.run(main)
else:
# The linked one-shot raises the child's boxed error
# directly in this caller task.
with pytest.raises(
tractor.RemoteActorError,
) as excinfo:
with pytest.raises(tractor.RemoteActorError) as excinfo:
trio.run(main)
assert excinfo.value.boxed_type == errtype
else:
# the root task will also error on the `Portal.result()`
# call so we expect an error from there AND the child.
# |_ tho seems like on new `trio` this doesn't always
# happen?
with pytest.raises((
BaseExceptionGroup,
tractor.RemoteActorError,
)) as excinfo:
trio.run(main)
# ensure boxed errors are `errtype`
err: BaseException = excinfo.value
if isinstance(err, BaseExceptionGroup):
suberrs: list[BaseException] = err.exceptions
else:
suberrs: list[BaseException] = [err]
for exc in suberrs:
assert exc.boxed_type == errtype
def test_multierror(
reg_addr: tuple[str, int],
@ -234,16 +248,16 @@ def test_cancel_single_subactor(
'''
async with tractor.open_nursery(
registry_addrs=[reg_addr],
) as an:
) as nursery:
portal = await an.start_actor(
portal = await nursery.start_actor(
'nothin', enable_modules=[__name__],
)
assert (await portal.run(do_nothing)) is None
if mechanism == 'nursery_cancel':
# would hang otherwise
await an.cancel()
await nursery.cancel()
else:
raise mechanism
@ -275,8 +289,8 @@ async def test_cancel_infinite_streamer(
trio.fail_after(4),
trio.move_on_after(1) as cancel_scope
):
async with tractor.open_nursery() as an:
portal = await an.start_actor(
async with tractor.open_nursery() as n:
portal = await n.start_actor(
'donny',
enable_modules=[__name__],
)
@ -289,7 +303,7 @@ async def test_cancel_infinite_streamer(
# we support trio's cancellation system
assert cancel_scope.cancelled_caught
assert an.cancel_called
assert n.cancel_called
@pytest.mark.parametrize(
@ -377,9 +391,10 @@ async def test_some_cancels_all(
tn.start_soon(
partial(
tractor.to_actor.run,
partial(func, **kwargs),
func,
an=an,
name=f'actor_{i}',
**kwargs,
)
)
@ -454,14 +469,12 @@ async def spawn_and_error(
if depth > 0:
args = (
partial(
spawn_and_error,
breadth=breadth,
depth=depth - 1,
),
)
kwargs = {
'name': f'spawner_{i}_depth_{depth}',
'breadth': breadth,
'depth': depth - 1,
}
else:
args = (
@ -685,20 +698,18 @@ async def test_nested_multierrors(
async with fail_after_w_trace(timeout):
try:
async with (
tractor.open_nursery() as an,
tractor.open_nursery() as nursery,
trio.open_nursery() as tn,
):
for i in range(subactor_breadth):
tn.start_soon(
partial(
tractor.to_actor.run,
partial(
spawn_and_error,
an=nursery,
name=f'spawner_{i}',
breadth=subactor_breadth,
depth=depth,
),
an=an,
name=f'spawner_{i}',
)
)
except (
@ -781,8 +792,8 @@ def test_cancel_via_SIGINT(
with trio.fail_after(2):
async with tractor.open_nursery(
registry_addrs=[reg_addr],
) as an:
await an.start_actor('sucka')
) as tn:
await tn.start_actor('sucka')
if 'mp' in start_method:
time.sleep(0.1)
os.kill(pid, signal.SIGINT)
@ -1043,8 +1054,8 @@ def test_fast_graceful_cancel_when_spawn_task_in_soft_proc_wait_for_daemon(
start = time.time()
try:
async with trio.open_nursery() as nurse:
async with tractor.open_nursery() as an:
p = await an.start_actor(
async with tractor.open_nursery() as tn:
p = await tn.start_actor(
'fast_boi',
enable_modules=[__name__],
)

View File

@ -156,8 +156,8 @@ def test_actor_managed_trio_nursery_task_error_cancels_aio(
async def main():
# cancel the nursery shortly after boot
async with tractor.open_nursery() as an:
p = await an.start_actor(
async with tractor.open_nursery() as n:
p = await n.start_actor(
'nursery_mngr',
infect_asyncio=asyncio_mode, # TODO, is this enabling debug mode?
enable_modules=[__name__],

View File

@ -5,7 +5,7 @@ The hipster way to force SC onto the stdlib's "async": 'infection mode'.
import asyncio
import builtins
from contextlib import ExitStack
from functools import partial
# from functools import partial
import itertools
import importlib
import os
@ -209,12 +209,10 @@ def test_aio_simple_error(
debug_mode=debug_mode,
) as an:
await to_actor.run(
partial(
asyncio_actor,
an=an,
target='sleep_and_err',
expect_err='AssertionError',
),
an=an,
infect_asyncio=True,
)
@ -457,14 +455,10 @@ def test_aio_cancelled_from_aio_causes_trio_cancelled(
# relays the remote error here in the caller's task.
with trio.fail_after(1 + delay):
await to_actor.run(
partial(
asyncio_actor,
target='aio_cancel',
expect_err=(
'tractor.to_asyncio.AsyncioCancelled'
),
),
an=an,
target='aio_cancel',
expect_err='tractor.to_asyncio.AsyncioCancelled',
infect_asyncio=True,
)
@ -669,12 +663,10 @@ def test_basic_interloop_channel_stream(
) as an:
# should raise RAE diectly
await to_actor.run(
partial(
stream_from_aio,
fan_out=fan_out,
),
an=an,
infect_asyncio=True,
fan_out=fan_out,
)
trio.run(main)
@ -690,11 +682,9 @@ def test_trio_error_cancels_intertask_chan(
) as an:
# should trigger remote actor error
await to_actor.run(
partial(
stream_from_aio,
trio_raise_err=True,
),
an=an,
trio_raise_err=True,
infect_asyncio=True,
)
@ -729,11 +719,9 @@ def test_trio_closes_early_causes_aio_checkpoint_raise(
# should raise RAE diectly
print('waiting on final infected subactor result..')
res: None = await to_actor.run(
partial(
stream_from_aio,
trio_exit_early=True,
),
an=an,
trio_exit_early=True,
infect_asyncio=True,
)
assert res is None
@ -782,13 +770,11 @@ def test_aio_exits_early_relays_AsyncioTaskExited(
# should raise RAE diectly
print('waiting on final infected subactor result..')
res: None = await to_actor.run(
partial(
stream_from_aio,
trio_exit_early=False,
aio_exit_early=True,
),
an=an,
infect_asyncio=True,
trio_exit_early=False,
aio_exit_early=True,
)
assert res is None
print(f'infected subactor returned result: {res!r}\n')
@ -823,11 +809,9 @@ def test_aio_errors_and_channel_propagates_and_closes(
) as an:
# should trigger RAE directly, not an eg.
await to_actor.run(
partial(
stream_from_aio,
aio_raise_err=True,
),
an=an,
aio_raise_err=True,
infect_asyncio=True,
)

View File

@ -163,12 +163,12 @@ def test_do_not_swallow_error_before_started_by_remote_contextcancelled(
async def main():
async with tractor.open_nursery(
debug_mode=debug_mode,
) as an:
portal = await an.start_actor(
) as n:
portal = await n.start_actor(
'errorer',
enable_modules=[__name__],
)
await an.start_actor(
await n.start_actor(
'sleeper',
enable_modules=[__name__],
)

View File

@ -139,9 +139,9 @@ async def test_required_args(callwith_expecterror):
with pytest.raises(err):
await func(**kwargs)
else:
async with tractor.open_nursery() as an:
async with tractor.open_nursery() as n:
portal = await an.start_actor(
portal = await n.start_actor(
name='pubber',
enable_modules=[__name__],
)
@ -180,7 +180,7 @@ def test_multi_actor_subs_arbiter_pub(
tractor.open_nursery(
registry_addrs=[reg_addr],
enable_modules=[__name__],
) as an,
) as n,
trio.open_nursery() as tn,
):
@ -188,7 +188,7 @@ def test_multi_actor_subs_arbiter_pub(
if pub_actor == 'streamer':
# start the publisher as a daemon
master_portal = await an.start_actor(
master_portal = await n.start_actor(
'streamer',
enable_modules=[__name__],
)
@ -215,11 +215,11 @@ def test_multi_actor_subs_arbiter_pub(
):
pass # expected once we `cancel_actor()` below
even_portal = await an.start_actor(
even_portal = await n.start_actor(
'evens',
enable_modules=[__name__],
)
odd_portal = await an.start_actor(
odd_portal = await n.start_actor(
'odds',
enable_modules=[__name__],
)
@ -294,9 +294,9 @@ def test_single_subactor_pub_multitask_subs(
async with tractor.open_nursery(
registry_addrs=[reg_addr],
enable_modules=[__name__],
) as an:
) as n:
portal = await an.start_actor(
portal = await n.start_actor(
'streamer',
enable_modules=[__name__],
)

View File

@ -143,12 +143,9 @@ def test_ringbuf(
child_read_shm,
**common_kwargs,
total_bytes=total_bytes,
) as (rctx, _sent),
) as (sctx, _sent),
):
# ctx-acm exits await each child task's
# `Return` (the prior `recv_p.result()` here
# was a daemon-portal no-op).
pass
await recv_p.result()
await send_p.cancel_actor()
await recv_p.cancel_actor()

View File

@ -4,7 +4,6 @@ related API and error checks.
'''
import itertools
from functools import partial
from unittest.mock import (
AsyncMock,
Mock,
@ -235,26 +234,24 @@ def test_rpc_errors(
# do that if actually debugging subactor but keep it
# disabled for the test.
# debug_mode=True,
) as an:
) as n:
actor = tractor.current_actor()
assert actor.is_registrar
await tractor.to_actor.run(
partial(
sleep_back_actor,
an=n,
actor_name=subactor_requests_to,
func_name=funcname,
func_defined=bool(func_defined),
exposed_mods=exposed_mods,
reg_addr=reg_addr,
),
an=an,
name='subactor',
# Function from the local exposed module space the
# subactor invokes when it RPCs back to this actor.
# function from the local exposed module space
# the subactor will invoke when it RPCs back to this actor
func_name=funcname,
exposed_mods=exposed_mods,
func_defined=True if func_defined else False,
enable_modules=subactor_exposed_mods,
reg_addr=reg_addr,
)
def run():

View File

@ -2,7 +2,6 @@
Verifying internal runtime state and undocumented extras.
"""
from functools import partial
import os
import pytest
@ -85,12 +84,10 @@ async def test_lifetime_stack_wipes_tmpfile(
loglevel=loglevel,
) as an:
await tractor.to_actor.run(
partial(
crash_and_clean_tmpdir,
an=an,
tmp_file_path=path,
error=error_in_child,
),
an=an,
)
except (
tractor.RemoteActorError,

View File

@ -51,18 +51,18 @@ async def spawn(
# recursively spawn this same `spawn()` fn as the lone
# task of a one-shot child subactor and get its result.
result = await tractor.to_actor.run(
partial(
spawn,
should_be_root=False,
data=data_to_pass_down,
reg_addr=reg_addr,
),
an=an,
# spawning args
name='sub-actor',
enable_modules=[__name__],
# passed to a subactor-recursive RPC invoke
# of this same `spawn()` fn.
should_be_root=False,
data=data_to_pass_down,
reg_addr=reg_addr,
)
assert result == 10
return result
@ -152,11 +152,9 @@ async def test_most_beautiful_word(
debug_mode=debug_mode,
) as an:
res: Any = await tractor.to_actor.run(
partial(
cellar_door,
return_value=return_value,
),
an=an,
return_value=return_value,
name='some_linguist',
)
assert res == return_value
@ -204,14 +202,12 @@ def test_loglevel_propagated_to_subactor(
start_method=start_method,
registry_addrs=[reg_addr],
) as an:
) as tn:
await tractor.to_actor.run(
partial(
check_loglevel,
level=level,
),
an=an,
an=tn,
loglevel=level,
level=level,
)
trio.run(main)
@ -277,23 +273,19 @@ def test_to_actor_run_can_skip_parent_main_inheritance(
# Default: child receives parent __main__ bootstrap data
await tractor.to_actor.run(
partial(
check_parent_main_inheritance,
expect_inherited=True,
),
an=an,
name='replaying-parent-main',
expect_inherited=True,
)
# Opt-out: child gets no parent __main__ data
await tractor.to_actor.run(
partial(
check_parent_main_inheritance,
expect_inherited=False,
),
an=an,
name='isolated-parent-main',
inherit_parent_main=False,
expect_inherited=False,
)
trio.run(main)

View File

@ -1,8 +1,8 @@
'''
`tractor.to_actor`: one-shot single-remote-task API suite.
Verifies the "spiritual successor" to (and replacement of)
the removed legacy `ActorNursery.run_in_actor()`; see
Verifies the "spiritual successor" to (and eventual
replacement of) `ActorNursery.run_in_actor()`; see
https://github.com/goodboy/tractor/issues/477
'''
@ -178,7 +178,7 @@ async def test_remote_error_relayed_to_caller_task(
A remote task error is raised directly in the
caller's task as a boxed `RemoteActorError` instead
of surfacing at actor-nursery teardown as with the
removed legacy `.run_in_actor()` API.
legacy `.run_in_actor()` API.
'''
with pytest.raises(RemoteActorError) as excinfo:

View File

@ -780,7 +780,7 @@ class Context:
# `Portal.open_context()` has been opened since it's
# assumed that other portal APIs like,
# - `Portal.run()`,
# - `to_actor.run()`
# - `ActorNursery.run_in_actor()`
# do their own error checking at their own call points and
# result processing.

View File

@ -1161,6 +1161,10 @@ class TransportClosed(Exception):
)
class NoResult(RuntimeError):
"No final result is expected for this actor"
class ModuleNotExposed(ModuleNotFoundError):
"The requested module is not exposed for RPC"

View File

@ -221,18 +221,15 @@ def pub(
import tractor
async with tractor.open_nursery() as n:
portal = await n.start_actor(
portal = n.run_in_actor(
'publisher', # actor name
enable_modules=[__name__],
)
async with portal.open_stream_from(
partial( # func to execute in it
pub_service,
topics=('clicks', 'users'),
task_name='source1',
)
) as stream:
async for value in stream:
)
async for value in await portal.result():
print(f"Subscriber received {value}")

View File

@ -306,11 +306,12 @@ class Start(
It is called by all the following public APIs:
- `to_actor.run()`
- `ActorNursery.run_in_actor()`
- `Portal.run()`
`|_.run_from_ns()`
`|_.open_stream_from()`
`|_._submit_for_result()`
- `Context.open_context()`

View File

@ -50,11 +50,13 @@ from ..ipc import Channel
from ..log import get_logger
from ..msg import (
# Error,
PayloadMsg,
NamespacePath,
Return,
)
from .._exceptions import (
ActorTooSlowError,
NoResult,
TransportClosed,
)
from .._context import (
@ -100,6 +102,14 @@ class Portal:
) -> None:
self._chan: Channel = channel
# during the portal's lifetime
self._final_result_pld: Any|None = None
self._final_result_msg: PayloadMsg|None = None
# When set to a ``Context`` (when _submit_for_result is called)
# it is expected that ``result()`` will be awaited at some
# point.
self._expect_result_ctx: Context|None = None
self._streams: set[MsgStream] = set()
# TODO, this should be PRIVATE (and never used publicly)! since it's just
@ -127,6 +137,102 @@ class Portal:
)
return self.chan
# TODO: factor this out into a `.highlevel` API-wrapper that uses
# a single `.open_context()` call underneath.
async def _submit_for_result(
self,
ns: str,
func: str,
**kwargs
) -> None:
if self._expect_result_ctx is not None:
raise RuntimeError(
'A pending main result has already been submitted'
)
self._expect_result_ctx: Context = await self.actor.start_remote_task(
self.channel,
nsf=NamespacePath(f'{ns}:{func}'),
kwargs=kwargs,
portal=self,
)
# TODO: we should deprecate this API right? since if we remove
# `.run_in_actor()` (and instead move it to a `.highlevel`
# wrapper api (around a single `.open_context()` call) we don't
# really have any notion of a "main" remote task any more?
#
# @api_frame
async def wait_for_result(
self,
hide_tb: bool = True,
) -> Any:
'''
Return the final result delivered by a `Return`-msg from the
remote peer actor's "main" task's `return` statement.
'''
__tracebackhide__: bool = hide_tb
# Check for non-rpc errors slapped on the
# channel for which we always raise
exc = self.channel._exc
if exc:
raise exc
# not expecting a "main" result
if self._expect_result_ctx is None:
peer_id: str = f'{self.channel.aid.reprol()!r}'
log.warning(
f'Portal to peer {peer_id} will not deliver a final result?\n'
f'\n'
f'Context.result() can only be called by the parent of '
f'a sub-actor when it was spawned with '
f'`ActorNursery.run_in_actor()`'
f'\n'
f'Further this `ActorNursery`-method-API will deprecated in the'
f'near fututre!\n'
)
return NoResult
# expecting a "main" result
assert self._expect_result_ctx
if self._final_result_msg is None:
try:
(
self._final_result_msg,
self._final_result_pld,
) = await self._expect_result_ctx._pld_rx.recv_msg(
ipc=self._expect_result_ctx,
expect_msg=Return,
)
except BaseException as err:
# TODO: wrap this into `@api_frame` optionally with
# some kinda filtering mechanism like log levels?
__tracebackhide__: bool = False
raise err
return self._final_result_pld
# TODO: factor this out into a `.highlevel` API-wrapper that uses
# a single `.open_context()` call underneath.
async def result(
self,
*args,
**kwargs,
) -> Any|Exception:
typname: str = type(self).__name__
log.warning(
f'`{typname}.result()` is DEPRECATED!\n'
f'\n'
f'Use `{typname}.wait_for_result()` instead!\n'
)
return await self.wait_for_result(
*args,
**kwargs,
)
async def _cancel_streams(self):
# terminate all locally running async generator
# IPC calls

View File

@ -20,6 +20,7 @@
"""
from contextlib import asynccontextmanager as acm
from functools import partial
import inspect
from typing import (
TYPE_CHECKING,
)
@ -258,6 +259,14 @@ class ActorNursery:
# and syncing purposes to any actor opened nurseries.
self._implicit_runtime_started: bool = False
# TODO, factor this into a .hilevel api!
#
# portals spawned with ``run_in_actor()`` are
# cancelled when their "main" result arrives. Reaped by
# `_reap_ria_portals()` at nursery-block exit now that
# the 2ndary `._ria_nursery` is gone (see issue #477).
self._cancel_after_result_on_exit: set = set()
# trio.Nursery-like cancel (request) statuses
self._cancelled_caught: bool = False
self._cancel_called: bool = False
@ -463,6 +472,84 @@ class ActorNursery:
)
)
# TODO: DEPRECATE THIS:
# -[x] impl instead as a hilevel wrapper on top of
# the lower level daemon-spawn + portal APIs
# |_ see `.to_actor.run()` (issue #477) which does
# `.start_actor()` + `Portal.run()` + a one-shot
# reap via `Portal.cancel_actor()`.
# -[ ] emit a `DeprecationWarning` here (requires
# migrating all in-repo usage first!)
# -[ ] use @api_frame on the wrapper
async def run_in_actor(
self,
fn: typing.Callable,
*,
name: str | None = None,
bind_addrs: UnwrappedAddress|None = None,
rpc_module_paths: list[str] | None = None,
enable_modules: list[str] | None = None,
loglevel: str | None = None, # set log level per subactor
infect_asyncio: bool = False,
inherit_parent_main: bool = True,
proc_kwargs: dict[str, typing.Any] | None = None,
**kwargs, # explicit args to ``fn``
) -> Portal:
'''
Spawn a new actor, run a lone task, then terminate the actor and
return its result.
Actors spawned using this method are kept alive at nursery teardown
until the task spawned by executing ``fn`` completes at which point
the actor is terminated.
NOTE: prefer the (eventual) replacement API
`tractor.to_actor.run()` which delivers the same
one-shot semantics decoupled from this nursery's
internal spawn machinery; see issue #477.
'''
__runtimeframe__: int = 1 # noqa
mod_path: str = fn.__module__
if name is None:
# use the explicit function name if not provided
name = fn.__name__
proc_kwargs = dict(proc_kwargs or {})
portal: Portal = await self.start_actor(
name,
enable_modules=[mod_path] + (
enable_modules or rpc_module_paths or []
),
bind_addrs=bind_addrs,
loglevel=loglevel,
infect_asyncio=infect_asyncio,
inherit_parent_main=inherit_parent_main,
proc_kwargs=proc_kwargs
)
# XXX: don't allow stream funcs
if not (
inspect.iscoroutinefunction(fn) and
not getattr(fn, '_tractor_stream_function', False)
):
raise TypeError(f'{fn} must be an async function!')
# this marks the actor to be cancelled after its portal result
# is retreived, see logic in `open_nursery()` below.
self._cancel_after_result_on_exit.add(portal)
await portal._submit_for_result(
mod_path,
fn.__name__,
**kwargs
)
return portal
# @api_frame
async def cancel(
self,
@ -583,6 +670,51 @@ class ActorNursery:
self._request_reap_all()
async def _reap_ria_portals(
an: ActorNursery,
errors: dict[tuple[str, str], BaseException],
ria_children: list[tuple[Portal, Actor]]|None = None,
) -> None:
'''
Wait on and stash the final result/error from every
`.run_in_actor()`-spawned child then cancel its actor
runtime, one `_spawn.cancel_on_completion()` task per
child.
Replaces the per-child reaper task formerly spawned by
the spawn backends (keyed off
`._cancel_after_result_on_exit` membership) which
required routing such children into the (now removable)
`._ria_nursery`. Only call AFTER `._join_procs` is set
so user code inside the nursery block retains exclusive
result-await access; see the "manually await results"
note in `spawn._mp.mp_proc()`.
'''
if ria_children is None:
ria_children: list[tuple[Portal, Actor]] = [
(portal, subactor)
for subactor, _, portal in an._children.values()
if portal in an._cancel_after_result_on_exit
]
if not ria_children:
return
async with (
collapse_eg(),
trio.open_nursery() as tn,
):
portal: Portal
subactor: Actor
for portal, subactor in ria_children:
tn.start_soon(
_spawn.cancel_on_completion,
portal,
subactor,
errors,
)
@acm
async def _open_and_supervise_one_cancels_all_nursery(
actor: Actor,
@ -597,10 +729,11 @@ async def _open_and_supervise_one_cancels_all_nursery(
errors: dict[tuple[str, str], BaseException] = {}
# The single "daemon actor" nursery into which ALL subactors
# are spawned; one-shot (`to_actor.run()`) subactors are
# result-waited and reaped in their caller's own task-scope
# (see the #477 `.run_in_actor()`/`._ria_nursery` removal);
# errors from this nursery bubble up to the caller.
# are spawned — both `.start_actor()` daemons AND
# `.run_in_actor()` one-shots. The latter's result-reaping now
# runs via `_reap_ria_portals()` at block-exit rather than a
# 2ndary `._ria_nursery` (see the #477 removal); errors from
# this nursery bubble up to the caller.
async with (
collapse_eg(),
trio.open_nursery() as da_nursery,
@ -624,6 +757,11 @@ async def _open_and_supervise_one_cancels_all_nursery(
)
an._request_reap_all()
# collect results (and errors) from all
# `.run_in_actor()` children then cancel
# each, one reaper task per child.
await _reap_ria_portals(an, errors)
# Single one-cancels-all handler for the (now single)
# daemon nursery. Pre-#477 a 2ndary `._ria_nursery`
# required a separate *outer* handler to catch errors
@ -697,14 +835,46 @@ async def _open_and_supervise_one_cancels_all_nursery(
# '------ - ------'
)
# snapshot `.run_in_actor()` children
# BEFORE cancelling: each backend
# spawn-task pops its `._children`
# entry as the proc gets reaped.
ria_children: list = [
(portal, subactor)
for subactor, _, portal
in an._children.values()
if portal in
an._cancel_after_result_on_exit
]
# cancel all subactors
await an.cancel()
# then collect any already-relayed
# results/errors from ria children.
# Tightly bounded: anything
# collectable is already queued in
# the local ctx (relayed BEFORE the
# cancel above); a child hard-killed
# without relaying just parks its
# reaper which then self-cleans (a
# `trio.Cancelled` result is never
# stashed), mirroring the old
# backend-side reaper-vs-`soft_kill`
# cancel race.
with trio.move_on_after(0.5):
await _reap_ria_portals(
an,
errors,
ria_children=ria_children,
)
finally:
# an error was stashed by the handler above (or by
# a spawn task via the shared `errors` dict) so
# cancel any remaining subactors, summarize and
# re-raise.
# No errors were raised while awaiting ".run_in_actor()"
# actors but those actors may have returned remote errors as
# results (meaning they errored remotely and have relayed
# those errors back to this parent actor). The errors are
# collected in ``errors`` so cancel all actors, summarize
# all errors and re-raise.
if errors:
if an._children:
with trio.CancelScope(shield=True):
@ -748,13 +918,15 @@ async def open_nursery(
Create and yield a new ``ActorNursery`` to be used for spawning
structured concurrent subactors.
When an actor is spawned a new trio task invokes one of the
process spawning backends to create and start a new subprocess.
These tasks are started in the supervisor's process nursery.
Spawning from a task is required because ``trio_run_in_process``
creates an internal nursery which the opening task **must** close;
this also makes each task's cancellation scope correspond to its
spawned subactor.
When an actor is spawned a new trio task is started which
invokes one of the process spawning backends to create and start
a new subprocess. These tasks are started by one of two nurseries
detailed below. The reason for spawning processes from within
a new task is because ``trio_run_in_process`` itself creates a new
internal nursery and the same task that opens a nursery **must**
close it. It turns out this approach is probably more correct
anyway since it is more clear from the following nested nurseries
which cancellation scopes correspond to each spawned subactor set.
'''
__tracebackhide__: bool = hide_tb

View File

@ -191,7 +191,9 @@ async def mp_proc(
# This is a "soft" (cancellable) join/reap which
# will remote cancel the actor on a ``trio.Cancelled``
# condition.
# condition. Any `.run_in_actor()` result-reaping
# happens up in the `ActorNursery` machinery (see
# `_supervise._reap_ria_portals()`), NOT here.
await soft_kill(
proc,
proc_waiter,

View File

@ -126,6 +126,98 @@ def try_set_start_method(
return _ctx
async def exhaust_portal(
portal: Portal,
actor: Actor
) -> Any:
'''
Pull final result from portal (assuming it has one).
If the main task is an async generator do our best to consume
what's left of it.
'''
__tracebackhide__ = True
try:
log.debug(
f'Waiting on final result from {actor.aid.uid}'
)
# XXX: streams should never be reaped here since they should
# always be established and shutdown using a context manager api
final: Any = await portal.wait_for_result()
except (
Exception,
BaseExceptionGroup,
) as err:
# we reraise in the parent task via a ``BaseExceptionGroup``
return err
except trio.Cancelled as err:
# lol, of course we need this too ;P
# TODO: merge with above?
log.warning(
'Cancelled portal result waiter task:\n'
f'uid: {portal.channel.aid}\n'
f'error: {err}\n'
)
return err
else:
log.debug(
f'Returning final result from portal:\n'
f'uid: {portal.channel.aid}\n'
f'result: {final}\n'
)
return final
async def cancel_on_completion(
portal: Portal,
actor: Actor,
errors: dict[tuple[str, str], Exception],
) -> None:
'''
Cancel actor gracefully once its "main" portal's
result arrives.
Should only be called for actors spawned via the
`Portal.run_in_actor()` API.
=> and really this API will be deprecated and should be
re-implemented as a `.hilevel.one_shot_task_nursery()`..)
'''
# if this call errors we store the exception for later
# in ``errors`` which will be reraised inside
# an exception group and we still send out a cancel request
result: Any|Exception = await exhaust_portal(
portal,
actor,
)
if isinstance(result, Exception):
errors[actor.aid.uid]: Exception = result
log.cancel(
'Cancelling subactor runtime due to error:\n\n'
f'Portal.cancel_actor() => {portal.channel.aid}\n\n'
f'error: {result}\n'
)
else:
log.runtime(
'Cancelling subactor gracefully:\n\n'
f'Portal.cancel_actor() => {portal.channel.aid}\n\n'
f'result: {result}\n'
)
# cancel the process now that we have a final result
await portal.cancel_actor()
async def hard_kill(
proc: trio.Process,
@ -372,8 +464,8 @@ async def new_proc(
# NOTE: bottom-of-module to avoid a circular import since the
# backend submodules pull `soft_kill`/`hard_kill`/`proc_waiter`
# from this module.
# backend submodules pull `cancel_on_completion`/`soft_kill`/
# `hard_kill`/`proc_waiter` from this module.
from ._trio import trio_proc
from ._mp import mp_proc

View File

@ -200,7 +200,9 @@ async def trio_proc(
# This is a "soft" (cancellable) join/reap which
# will remote cancel the actor on a ``trio.Cancelled``
# condition.
# condition. Any `.run_in_actor()` result-reaping
# happens up in the `ActorNursery` machinery (see
# `_supervise._reap_ria_portals()`), NOT here.
await soft_kill(
proc,
trio.Process.wait, # XXX, uses `pidfd_open()` below.

View File

@ -25,8 +25,8 @@ its result and (when the call owns the subactor) reap it.
Target arguments follow Trio's positional convention; use
`functools.partial()` to bind target keyword arguments.
The "spiritual successor" to (and replacement of) the removed
legacy `ActorNursery.run_in_actor()` API; see
The "spiritual successor" to (and eventual replacement of)
the `ActorNursery.run_in_actor()` API; see
https://github.com/goodboy/tractor/issues/477
'''

View File

@ -31,7 +31,7 @@ the lower level daemon-actor spawn + portal APIs,
such that error collection and propagation happens in the
*caller's task* (and thus whatever `trio` nursery/scope
encloses it) instead of inside the actor-nursery's
spawn-machinery nurseries as with the (now removed) legacy
spawn-machinery nurseries as with the (to be deprecated)
`ActorNursery.run_in_actor()` API.
'''
@ -280,8 +280,8 @@ async def run(
module in its `enable_modules` list. Calls that spawn their own
actor add the trampoline module automatically.
Unlike the removed legacy `ActorNursery.run_in_actor()` (which
returned a `Portal` whose result was only collected at
Unlike `ActorNursery.run_in_actor()` (which returns
a `Portal` whose result is only collected at
actor-nursery teardown) this is a plain "call and
wait" primitive: any remote error is raised HERE, in
the caller's task. Concurrency is composed the usual