Compare commits
No commits in common. "b645f7fa8c82c4bcdec4f791e9e45a6c7b68b1cc" and "0ac9fa1c0cc43c29812b4cf5b1959850ea9b8c90" have entirely different histories.
b645f7fa8c
...
0ac9fa1c0c
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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`).
|
||||
|
|
@ -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).
|
||||
|
|
@ -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.
|
||||
|
|
@ -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
|
||||
|
|
@ -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`.
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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::
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -21,35 +21,23 @@ 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(
|
||||
say_hello,
|
||||
other_actor=other_actor,
|
||||
)
|
||||
)
|
||||
|
||||
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("CUTTTT CUUTT CUT!!! Donny!! You're supposed to say...")
|
||||
donny = await n.run_in_actor(
|
||||
say_hello,
|
||||
name='donny',
|
||||
# arguments are always named
|
||||
other_actor='gretchen',
|
||||
)
|
||||
gretchen = await n.run_in_actor(
|
||||
say_hello,
|
||||
name='gretchen',
|
||||
other_actor='donny',
|
||||
)
|
||||
print(await gretchen.wait_for_result())
|
||||
print(await donny.wait_for_result())
|
||||
print("CUTTTT CUUTT CUT!!! Donny!! You're supposed to say...")
|
||||
|
||||
|
||||
if __name__ == '__main__':
|
||||
|
|
|
|||
|
|
@ -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__':
|
||||
|
|
|
|||
|
|
@ -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__],
|
||||
|
|
|
|||
|
|
@ -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__],
|
||||
)
|
||||
|
|
|
|||
|
|
@ -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(
|
||||
spawn_until,
|
||||
depth=depth,
|
||||
),
|
||||
an=an,
|
||||
await n.run_in_actor(
|
||||
spawn_until,
|
||||
depth=depth,
|
||||
name=f'spawn_until_{depth}',
|
||||
)
|
||||
|
||||
|
|
@ -80,38 +65,35 @@ async def main():
|
|||
└─ python -m tractor._child --uid ('spawn_until_0', 'de918e6d ...)
|
||||
|
||||
"""
|
||||
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(
|
||||
spawn_until,
|
||||
depth=3,
|
||||
),
|
||||
an=an,
|
||||
name='spawner0',
|
||||
)
|
||||
async with tractor.open_nursery(
|
||||
debug_mode=True,
|
||||
loglevel='pdb',
|
||||
) as n:
|
||||
|
||||
# spawn both actors
|
||||
portal = await n.run_in_actor(
|
||||
spawn_until,
|
||||
depth=3,
|
||||
name='spawner0',
|
||||
)
|
||||
tn.start_soon(
|
||||
partial(
|
||||
tractor.to_actor.run,
|
||||
partial(
|
||||
spawn_until,
|
||||
depth=4,
|
||||
),
|
||||
an=an,
|
||||
name='spawner1',
|
||||
)
|
||||
portal1 = await n.run_in_actor(
|
||||
spawn_until,
|
||||
depth=4,
|
||||
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__':
|
||||
trio.run(main)
|
||||
|
|
|
|||
|
|
@ -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(
|
||||
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__],
|
||||
)
|
||||
async with tractor.open_nursery(
|
||||
debug_mode=True,
|
||||
loglevel='devx',
|
||||
) 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)
|
||||
|
|
|
|||
|
|
@ -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__':
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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(
|
||||
spawn_until,
|
||||
depth=depth,
|
||||
),
|
||||
an=an,
|
||||
await n.run_in_actor(
|
||||
spawn_until,
|
||||
depth=depth,
|
||||
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(
|
||||
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',
|
||||
)
|
||||
)
|
||||
async with tractor.open_nursery(
|
||||
debug_mode=True,
|
||||
enable_transports=['uds'], # TODO, apss this via osenv?
|
||||
loglevel='devx', # XXX, required for test!
|
||||
) as n:
|
||||
|
||||
# ..while blocking on the shallow (faster to fail) tree
|
||||
# whose propagated error triggers nursery cancellation.
|
||||
await tractor.to_actor.run(
|
||||
partial(
|
||||
spawn_until,
|
||||
depth=0,
|
||||
),
|
||||
an=an,
|
||||
# spawn both actors
|
||||
portal = await n.run_in_actor(
|
||||
spawn_until,
|
||||
depth=0,
|
||||
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__':
|
||||
|
|
|
|||
|
|
@ -13,24 +13,17 @@ async def main():
|
|||
simultaneously.
|
||||
|
||||
'''
|
||||
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__],
|
||||
)
|
||||
async with tractor.open_nursery(
|
||||
debug_mode=True,
|
||||
# loglevel='debug' # ?XXX required?
|
||||
) 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
|
||||
|
|
|
|||
|
|
@ -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()
|
||||
|
|
|
|||
|
|
@ -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__],
|
||||
)
|
||||
|
|
|
|||
|
|
@ -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__':
|
||||
|
|
|
|||
|
|
@ -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__':
|
||||
|
|
|
|||
|
|
@ -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(
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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__],
|
||||
# )
|
||||
|
|
|
|||
|
|
@ -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():
|
||||
|
|
|
|||
|
|
@ -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:
|
||||
|
||||
# burn rubber in the parent too
|
||||
tn.start_soon(burn_cpu)
|
||||
portal = await n.run_in_actor(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)
|
||||
# burn rubber in the parent too
|
||||
await burn_cpu()
|
||||
|
||||
# wait on result from target function
|
||||
pid = await portal.wait_for_result()
|
||||
|
||||
# end of nursery block
|
||||
print(f"Collected subproc {pid}")
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -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__':
|
||||
|
|
|
|||
|
|
@ -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__],
|
||||
)
|
||||
|
|
|
|||
|
|
@ -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(
|
||||
|
|
|
|||
|
|
@ -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()
|
||||
tn.start_soon(_direct, 'donny', 'gretchen')
|
||||
tn.start_soon(_direct, 'gretchen', 'donny')
|
||||
print("CUTTTT CUUTT CUT!!?! Donny!! You're supposed to say...")
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -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(
|
||||
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,
|
||||
):
|
||||
async with tractor.open_nursery(
|
||||
registry_addrs=[reg_addr],
|
||||
debug_mode=debug_mode,
|
||||
) 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(
|
||||
consumer,
|
||||
subs=[sub],
|
||||
),
|
||||
an=an,
|
||||
name=f'consumer_{sub}',
|
||||
)
|
||||
await n.run_in_actor(
|
||||
consumer,
|
||||
name=f'consumer_{sub}',
|
||||
subs=[sub],
|
||||
)
|
||||
|
||||
# make one dynamic subscriber
|
||||
tn.start_soon(
|
||||
partial(
|
||||
tractor.to_actor.run,
|
||||
partial(
|
||||
consumer,
|
||||
subs=list(_registry.keys()),
|
||||
),
|
||||
an=an,
|
||||
name='consumer_dynamic',
|
||||
)
|
||||
await n.run_in_actor(
|
||||
consumer,
|
||||
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__],
|
||||
)
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
),
|
||||
spawn_and_error,
|
||||
)
|
||||
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,
|
||||
breadth=subactor_breadth,
|
||||
depth=depth,
|
||||
),
|
||||
an=an,
|
||||
spawn_and_error,
|
||||
an=nursery,
|
||||
name=f'spawner_{i}',
|
||||
breadth=subactor_breadth,
|
||||
depth=depth,
|
||||
)
|
||||
)
|
||||
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__],
|
||||
)
|
||||
|
|
|
|||
|
|
@ -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__],
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
target='sleep_and_err',
|
||||
expect_err='AssertionError',
|
||||
),
|
||||
asyncio_actor,
|
||||
an=an,
|
||||
target='sleep_and_err',
|
||||
expect_err='AssertionError',
|
||||
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'
|
||||
),
|
||||
),
|
||||
asyncio_actor,
|
||||
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,
|
||||
),
|
||||
stream_from_aio,
|
||||
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,
|
||||
),
|
||||
stream_from_aio,
|
||||
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,
|
||||
),
|
||||
stream_from_aio,
|
||||
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,
|
||||
),
|
||||
stream_from_aio,
|
||||
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,
|
||||
),
|
||||
stream_from_aio,
|
||||
an=an,
|
||||
aio_raise_err=True,
|
||||
infect_asyncio=True,
|
||||
)
|
||||
|
||||
|
|
|
|||
|
|
@ -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__],
|
||||
)
|
||||
|
|
|
|||
|
|
@ -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__],
|
||||
)
|
||||
|
|
|
|||
|
|
@ -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()
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
actor_name=subactor_requests_to,
|
||||
func_name=funcname,
|
||||
func_defined=bool(func_defined),
|
||||
exposed_mods=exposed_mods,
|
||||
reg_addr=reg_addr,
|
||||
),
|
||||
an=an,
|
||||
sleep_back_actor,
|
||||
an=n,
|
||||
actor_name=subactor_requests_to,
|
||||
|
||||
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():
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
tmp_file_path=path,
|
||||
error=error_in_child,
|
||||
),
|
||||
crash_and_clean_tmpdir,
|
||||
an=an,
|
||||
tmp_file_path=path,
|
||||
error=error_in_child,
|
||||
)
|
||||
except (
|
||||
tractor.RemoteActorError,
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
),
|
||||
spawn,
|
||||
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,
|
||||
),
|
||||
cellar_door,
|
||||
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,
|
||||
check_loglevel,
|
||||
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,
|
||||
),
|
||||
check_parent_main_inheritance,
|
||||
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,
|
||||
),
|
||||
check_parent_main_inheritance,
|
||||
an=an,
|
||||
name='isolated-parent-main',
|
||||
inherit_parent_main=False,
|
||||
expect_inherited=False,
|
||||
)
|
||||
|
||||
trio.run(main)
|
||||
|
|
|
|||
|
|
@ -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:
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
||||
|
|
|
|||
|
|
@ -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"
|
||||
|
||||
|
|
|
|||
|
|
@ -221,19 +221,16 @@ 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
|
||||
partial( # func to execute in it
|
||||
pub_service,
|
||||
topics=('clicks', 'users'),
|
||||
task_name='source1',
|
||||
)
|
||||
) as stream:
|
||||
async for value in stream:
|
||||
print(f"Subscriber received {value}")
|
||||
)
|
||||
async for value in await portal.result():
|
||||
print(f"Subscriber received {value}")
|
||||
|
||||
|
||||
Here, you don't need to provide the ``ctx`` argument since the
|
||||
|
|
|
|||
|
|
@ -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()`
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
||||
'''
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
Loading…
Reference in New Issue