Compare commits
28 Commits
481ba00332
...
38d6688d2a
| Author | SHA1 | Date |
|---|---|---|
|
|
38d6688d2a | |
|
|
38aa92bc61 | |
|
|
7e9e39c5f8 | |
|
|
0a1431c933 | |
|
|
3146409baf | |
|
|
dd9d890921 | |
|
|
0b926210a3 | |
|
|
7c6bff50c5 | |
|
|
17d6879706 | |
|
|
2af5acde88 | |
|
|
5c374b183e | |
|
|
29ac9067ce | |
|
|
efce12e145 | |
|
|
0c3b21e9fe | |
|
|
6a7e06e39d | |
|
|
aa9c9a97ff | |
|
|
07326d2b3e | |
|
|
e876430787 | |
|
|
0402994572 | |
|
|
fa748980d5 | |
|
|
5b5235962b | |
|
|
d8962f2a8a | |
|
|
d73c3b55e0 | |
|
|
2029a154ed | |
|
|
194ce5e392 | |
|
|
64ee2eeb75 | |
|
|
62b574a3d4 | |
|
|
93227086b0 |
|
|
@ -1,43 +0,0 @@
|
|||
---
|
||||
model: openai/gpt-5.6-sol
|
||||
service: opencode
|
||||
session: pr481-review-fixes-p1-20260818
|
||||
timestamp: 2026-08-18T03:15:32Z
|
||||
git_ref: 4151b956
|
||||
scope: code
|
||||
substantive: true
|
||||
raw_file: 20260818T031532Z_4151b956_prompt_io.raw.md
|
||||
---
|
||||
|
||||
## Prompt
|
||||
|
||||
Address the approved review findings on PR #481, but work
|
||||
iteratively: implement and verify one finding at a time, prepare a
|
||||
separate `/commit-plan` after each fix, and stop for the human commit
|
||||
before starting the next finding. Begin with the P1 per-child
|
||||
lifecycle issue. Also publish the already-approved review findings
|
||||
against the reviewed PR head before editing.
|
||||
|
||||
## Response summary
|
||||
|
||||
Published the approved non-approving review at head `4151b956`, then
|
||||
implemented only the P1 lifecycle fix. Owned one-shot actors now use
|
||||
a child-specific cancellation and process-reap handshake, including
|
||||
hard escalation for unacknowledged cancellation and deterministic
|
||||
bookkeeping removal before `to_actor.run()` returns.
|
||||
|
||||
## Files changed
|
||||
|
||||
- `tractor/runtime/_supervise.py` - coordinate child-specific cancel
|
||||
and reap.
|
||||
- `tractor/spawn/_trio.py` - wait on the Trio child's reap request.
|
||||
- `tractor/spawn/_mp.py` - wait on the multiprocessing child's reap
|
||||
request.
|
||||
- `tractor/spawn/_spawn.py` - publish monitor completion centrally.
|
||||
- `tractor/to_actor/_api.py` - await owned-child process reaping.
|
||||
- `tests/test_to_actor.py` - cover cleanup, escalation, and startup
|
||||
ordering.
|
||||
|
||||
## Human edits
|
||||
|
||||
None - the generated P1 patch remains uncommitted for human review.
|
||||
|
|
@ -1,74 +0,0 @@
|
|||
---
|
||||
model: openai/gpt-5.6-sol
|
||||
service: opencode
|
||||
timestamp: 2026-08-18T03:15:32Z
|
||||
git_ref: 4151b956
|
||||
diff_cmd: git diff HEAD~1..HEAD
|
||||
---
|
||||
|
||||
Implemented only the P1 lifecycle finding from the approved PR #481
|
||||
review, preserving the requested one-fix-at-a-time commit boundary.
|
||||
|
||||
> `git diff HEAD~1..HEAD -- tractor/runtime/_supervise.py`
|
||||
|
||||
Added per-child reap request/completion events to `ActorNursery`, a
|
||||
shielded child-specific cancel-and-reap operation, late-registration
|
||||
latching for nursery teardown, and cancellation escalation that waits
|
||||
for debugger release before using non-ignorable process termination.
|
||||
The nursery-wide cancellation path snapshots child records before
|
||||
checkpointing so concurrent one-shot cleanup cannot invalidate its
|
||||
iteration.
|
||||
|
||||
> `git diff HEAD~1..HEAD -- tractor/spawn/_trio.py`
|
||||
|
||||
Changed Trio child monitors to wait on their per-child reap requests.
|
||||
|
||||
> `git diff HEAD~1..HEAD -- tractor/spawn/_mp.py`
|
||||
|
||||
Changed multiprocessing child monitors to wait on their per-child reap
|
||||
requests.
|
||||
|
||||
> `git diff HEAD~1..HEAD -- tractor/spawn/_spawn.py`
|
||||
|
||||
Ensured every backend publishes child-reap completion after its process
|
||||
monitor exits.
|
||||
|
||||
> `git diff HEAD~1..HEAD -- tractor/to_actor/_api.py`
|
||||
|
||||
Changed owned one-shot cleanup to await child-specific process joining
|
||||
and bookkeeping removal instead of treating the cancel RPC as reaping.
|
||||
|
||||
> `git diff HEAD~1..HEAD -- tests/test_to_actor.py`
|
||||
|
||||
Added regressions for immediate caller-managed nursery cleanup, failed
|
||||
cancel acknowledgement escalation, and child registration after a
|
||||
latched nursery-wide teardown request.
|
||||
|
||||
Verification:
|
||||
|
||||
`pytest -q tests/test_to_actor.py tests/test_cancellation.py tests/test_spawning.py tests/discovery/test_multi_program.py`
|
||||
|
||||
Result: `46 passed, 1 xfailed, 3 xpassed`.
|
||||
|
||||
`pytest -q tests/test_to_actor.py --tpt-proto uds`
|
||||
|
||||
Result: `13 passed`.
|
||||
|
||||
`pytest -q tests/test_to_actor.py --spawn-backend mp_spawn --tpt-proto tcp`
|
||||
|
||||
Result: `13 passed`.
|
||||
|
||||
One broad verification run was mistakenly launched in parallel with
|
||||
the UDS and `mp_spawn` actor suites. It timed out
|
||||
`test_remote_error_from_caller_nursery`; the node passed immediately
|
||||
in isolation and the complete broad selection then passed serially.
|
||||
The failure was classified as concurrent test-session interference,
|
||||
not accepted as a passing boundary result.
|
||||
|
||||
Ruff, Python compilation, and `git diff --check` passed for the changed
|
||||
boundary. Ruff's existing `_trio.py` F401 finding was reproduced at the
|
||||
unmodified PR head and excluded from attribution to this patch.
|
||||
|
||||
No source files were staged, committed, pushed, or used for review
|
||||
replies. The previously approved top-level review was published before
|
||||
the fix at reviewed head `4151b956`.
|
||||
|
|
@ -1,36 +0,0 @@
|
|||
---
|
||||
model: openai/gpt-5.6-sol
|
||||
service: opencode
|
||||
session: ses_3e4c90d3eafeqHEtRYSIHgHhpA
|
||||
timestamp: 2026-08-18T19:30:01Z
|
||||
git_ref: bf06b4f8
|
||||
scope: code
|
||||
substantive: true
|
||||
raw_file: 20260818T193001Z_bf06b4f8_prompt_io.raw.md
|
||||
---
|
||||
|
||||
## Prompt
|
||||
|
||||
Fix cancellation during an indeterminate transport write without
|
||||
allowing a second RPC msg to corrupt the stream. Keep the transport
|
||||
fix separate from the higher-level `to_actor` remediation and explain
|
||||
the length-prefixed protocol rather than using unexplained "framed
|
||||
send" terminology.
|
||||
|
||||
## Response summary
|
||||
|
||||
Close a msgpack transport stream when cancellation interrupts its
|
||||
length-prefixed `send_all()` operation. The stream can no longer be
|
||||
safely reused because Trio cannot report how many bytes were written.
|
||||
|
||||
## Files changed
|
||||
|
||||
- `tractor/ipc/_transport.py` - close an interrupted send stream.
|
||||
- `tests/ipc/test_each_tpt.py` - cover cancellation during the write.
|
||||
|
||||
## Human edits
|
||||
|
||||
The human required this transport edge-case fix to land as its own
|
||||
behavioral commit with a detailed message. During staged review, the
|
||||
human also rejected the unexplained "framed send" wording and asked
|
||||
for terminology tied directly to the actual transport operation.
|
||||
|
|
@ -1,19 +0,0 @@
|
|||
---
|
||||
model: openai/gpt-5.6-sol
|
||||
service: opencode
|
||||
timestamp: 2026-08-18T19:30:01Z
|
||||
git_ref: bf06b4f8
|
||||
diff_cmd: git diff HEAD~1..HEAD
|
||||
---
|
||||
|
||||
Prospective review found that cancellation can interrupt
|
||||
`MsgpackTransport.send()` after `send_all()` writes only part of its
|
||||
length-prefixed msg. Sending a cancellation request afterward can
|
||||
append another msg to the indeterminate stream and desynchronize the
|
||||
peer decoder.
|
||||
|
||||
> `git diff HEAD~1..HEAD -- tractor/ipc/_transport.py tests/ipc/test_each_tpt.py`
|
||||
|
||||
Close the stream under a cancellation shield when `send_all()` is
|
||||
cancelled. Cover the behavior with a fake stream that checkpoints
|
||||
inside the write and records forced closure.
|
||||
|
|
@ -1,37 +0,0 @@
|
|||
---
|
||||
model: openai/gpt-5.6-sol
|
||||
service: opencode
|
||||
session: ses_3e4c90d3eafeqHEtRYSIHgHhpA
|
||||
timestamp: 2026-08-18T19:30:02Z
|
||||
git_ref: bf06b4f8
|
||||
scope: code
|
||||
substantive: true
|
||||
raw_file: 20260818T193002Z_bf06b4f8_prompt_io.raw.md
|
||||
---
|
||||
|
||||
## Prompt
|
||||
|
||||
Distill repeated `Actor._contexts.pop()` machinery into a wrapper like
|
||||
the RPC-task registration helper so future teardown sites do not keep
|
||||
reconstructing the context-registry key independently. Preserve the
|
||||
existing lifecycle-specific cleanup behavior.
|
||||
|
||||
## Response summary
|
||||
|
||||
Add idempotent `Actor._drop_context()` registry removal keyed from the
|
||||
context's own channel and CID. Use it for caller context teardown and
|
||||
the strict callee-side RPC deregistration path.
|
||||
|
||||
## Files changed
|
||||
|
||||
- `tractor/runtime/_runtime.py` - own context-registry removal.
|
||||
- `tractor/runtime/_rpc.py` - use the helper for callee teardown.
|
||||
- `tractor/_context.py` - use the helper after caller teardown.
|
||||
|
||||
## Human edits
|
||||
|
||||
The human identified the repeated registry-pop code and requested a
|
||||
central primitive analogous to `_register_rpc_task()`. The agent first
|
||||
suggested an async helper that also closed receive channels; the final
|
||||
design was narrowed to registry removal only so each lifecycle owner
|
||||
retains its existing closure, debugger, shielding, and error policy.
|
||||
|
|
@ -1,18 +0,0 @@
|
|||
---
|
||||
model: openai/gpt-5.6-sol
|
||||
service: opencode
|
||||
timestamp: 2026-08-18T19:30:02Z
|
||||
git_ref: bf06b4f8
|
||||
diff_cmd: git diff HEAD~1..HEAD
|
||||
---
|
||||
|
||||
Repeated teardown sites reconstruct the `Actor._contexts` registry
|
||||
key from a portal channel and context ID before popping it. Add an
|
||||
idempotent actor-owned helper deriving the key from the context itself,
|
||||
then route caller and callee context teardown through that helper.
|
||||
|
||||
> `git diff HEAD~1..HEAD -- tractor/runtime/_runtime.py tractor/runtime/_rpc.py tractor/_context.py`
|
||||
|
||||
Keep receive-channel closure and cancellation shielding in each
|
||||
lifecycle owner so the helper centralizes registry machinery without
|
||||
changing their teardown ordering.
|
||||
|
|
@ -1,38 +0,0 @@
|
|||
---
|
||||
model: openai/gpt-5.6-sol
|
||||
service: opencode
|
||||
session: ses_3e4c90d3eafeqHEtRYSIHgHhpA
|
||||
timestamp: 2026-08-18T19:30:03Z
|
||||
git_ref: bf06b4f8
|
||||
scope: code
|
||||
substantive: true
|
||||
raw_file: 20260818T193003Z_bf06b4f8_prompt_io.raw.md
|
||||
---
|
||||
|
||||
## Prompt
|
||||
|
||||
Cancel a remote task when its caller is cancelled after `Start`
|
||||
publication but before startup acknowledgement. Keep cancellation
|
||||
bounded, prevent its private `_cancel_task` RPC from recursively
|
||||
cancelling itself and preserve public target kwargs unchanged.
|
||||
|
||||
## Response summary
|
||||
|
||||
Add private portal startup policy, use it for non-recursive context
|
||||
cancellation and clean caller-side startup state under a shield.
|
||||
|
||||
## Files changed
|
||||
|
||||
- `tractor/runtime/_portal.py` - separate private startup policy.
|
||||
- `tractor/_context.py` - disable recursion for cancellation RPCs.
|
||||
- `tractor/runtime/_runtime.py` - clean cancelled task startup.
|
||||
- `tests/test_context_stream_semantics.py` - control cancellation
|
||||
between `Start` publication and acknowledgement.
|
||||
|
||||
## Human edits
|
||||
|
||||
The human required this cancellation behavior to remain a distinct
|
||||
commit from general startup failures and from the public `to_actor`
|
||||
API. The human also requested that its runtime comment describe the
|
||||
actual length-prefixed transport guarantee and concrete `_cancel_task`
|
||||
operation rather than referring to an unnamed wrapper.
|
||||
|
|
@ -1,20 +0,0 @@
|
|||
---
|
||||
model: openai/gpt-5.6-sol
|
||||
service: opencode
|
||||
timestamp: 2026-08-18T19:30:03Z
|
||||
git_ref: bf06b4f8
|
||||
diff_cmd: git diff HEAD~1..HEAD
|
||||
---
|
||||
|
||||
Cancellation while `Actor.start_remote_task()` waits for `StartAck`
|
||||
can strand its caller-side context and leave the remote task running.
|
||||
Make one bounded cleanup request, remove local startup state and close
|
||||
its receive channel.
|
||||
|
||||
> `git diff HEAD~1..HEAD -- tractor/runtime/_runtime.py tractor/runtime/_portal.py tractor/_context.py tests/test_context_stream_semantics.py`
|
||||
|
||||
Separate private startup-cancellation policy from public target kwargs
|
||||
using `Portal._run_from_ns()`. Have `Context.cancel()` disable recursive
|
||||
startup cancellation for its own `_cancel_task` RPC. Exercise
|
||||
cancellation after `Start` publication and prove the caller-owned actor
|
||||
remains reusable without leaked contexts.
|
||||
|
|
@ -1,36 +0,0 @@
|
|||
---
|
||||
model: openai/gpt-5.6-sol
|
||||
service: opencode
|
||||
session: ses_3e4c90d3eafeqHEtRYSIHgHhpA
|
||||
timestamp: 2026-08-18T19:30:04Z
|
||||
git_ref: bf06b4f8
|
||||
scope: code
|
||||
substantive: true
|
||||
raw_file: 20260818T193004Z_bf06b4f8_prompt_io.raw.md
|
||||
---
|
||||
|
||||
## Prompt
|
||||
|
||||
Release caller-side context state for every remote-task startup failure,
|
||||
not only local cancellation. Preserve the remote error, avoid unsafe
|
||||
follow-up sends and prove pre-publication serialization failures leave
|
||||
a reused portal healthy.
|
||||
|
||||
## Response summary
|
||||
|
||||
Extend remote-task startup cleanup across send, acknowledgement and
|
||||
validation errors. Track completed publication, perform only safe
|
||||
best-effort cancellation and deterministically remove local state.
|
||||
|
||||
## Files changed
|
||||
|
||||
- `tractor/runtime/_runtime.py` - clean every startup failure path.
|
||||
- `tests/test_context_stream_semantics.py` - cover authorization and
|
||||
serialization failures before context entry.
|
||||
|
||||
## Human edits
|
||||
|
||||
The human accepted the discovered edge-case fixes but required general
|
||||
startup cleanup to land separately from cancellation cleanup, transport
|
||||
integrity and the public API. This boundary preserves that behavioral
|
||||
distinction and its dedicated commit-message rationale.
|
||||
|
|
@ -1,19 +0,0 @@
|
|||
---
|
||||
model: openai/gpt-5.6-sol
|
||||
service: opencode
|
||||
timestamp: 2026-08-18T19:30:04Z
|
||||
git_ref: bf06b4f8
|
||||
diff_cmd: git diff HEAD~1..HEAD
|
||||
---
|
||||
|
||||
`Actor.start_remote_task()` inserts a context before sending `Start`,
|
||||
but startup errors other than cancellation escape without removing or
|
||||
closing that caller state. Serialization errors, acknowledgement
|
||||
timeouts, malformed acknowledgements and remote authorization errors
|
||||
can therefore leak context-registry entries.
|
||||
|
||||
> `git diff HEAD~1..HEAD -- tractor/runtime/_runtime.py tests/test_context_stream_semantics.py`
|
||||
|
||||
Cover the complete send, acknowledgement and validation phase with
|
||||
exceptional cleanup. Attempt remote cancellation only when publication
|
||||
is known complete or protocol-safe, and always release local state.
|
||||
|
|
@ -1,52 +0,0 @@
|
|||
---
|
||||
model: openai/gpt-5.6-sol
|
||||
service: opencode
|
||||
session: ses_3e4c90d3eafeqHEtRYSIHgHhpA
|
||||
timestamp: 2026-08-18T19:30:05Z
|
||||
git_ref: bf06b4f8
|
||||
scope: code
|
||||
substantive: true
|
||||
raw_file: 20260818T193005Z_bf06b4f8_prompt_io.raw.md
|
||||
---
|
||||
|
||||
## Prompt
|
||||
|
||||
Replace abandoned `Portal.run()` one-shots with a static linked-context
|
||||
endpoint. Follow Trio positional-call semantics, use partials for target
|
||||
keywords, preserve Python 3.14 Placeholder behavior, keep target lookup
|
||||
behind the RPC allowlist and support private, nursery and portal
|
||||
placement.
|
||||
|
||||
## Response summary
|
||||
|
||||
Use `Portal.open_context()` and `Context.wait_for_result()` for one-shot
|
||||
tasks. Normalize every partial layer, validate signatures locally and
|
||||
send target namespace/function components separately to the authorized
|
||||
remote resolver. Retain the client-side function in its `NamespacePath`
|
||||
so `to_tuple()` does not re-import it. Owned actors enable the declaring
|
||||
`_api.__name__` directly; caller-owned portals opt in through the public
|
||||
`to_actor.MODULE` alias.
|
||||
|
||||
## Files changed
|
||||
|
||||
- `tractor/to_actor/_api.py` - implement linked one-shot calls.
|
||||
- `tractor/to_actor/__init__.py` - export `MODULE`.
|
||||
- `tractor/msg/ptr.py` - retain refs created by `from_ref()`.
|
||||
- `tests/test_to_actor.py` - cover the public API and authorization.
|
||||
- `examples/parallelism/to_actor_one_shots.py` - use positional inputs.
|
||||
|
||||
## Human edits
|
||||
|
||||
The human rejected nested target-kwargs configuration and selected
|
||||
Trio-style positional inputs plus `functools.partial()`. During staged
|
||||
review the human required a Python 3.14 compatibility comment rather
|
||||
than removing Placeholder support, requested separate namespace and
|
||||
function inputs, preserved `_get_rpc_func(ns: str, funcname: str)`
|
||||
authorization, renamed `RPC_MODULE` to `MODULE`, rejected global module
|
||||
exposure and deferred speculative nursery/module-list helpers to the
|
||||
`open_taskman()` design line. The human also required this public API
|
||||
to land only after its lower-level safety dependencies. In final staged
|
||||
review, the human required `_invoke_from_portal()` to use
|
||||
`NamespacePath.to_tuple()` with the already-held function ref and
|
||||
required internal actor setup to use `_api.__name__` directly, keeping
|
||||
`to_actor.MODULE` solely as the public importer-facing alias.
|
||||
|
|
@ -1,24 +0,0 @@
|
|||
---
|
||||
model: openai/gpt-5.6-sol
|
||||
service: opencode
|
||||
timestamp: 2026-08-18T19:30:05Z
|
||||
git_ref: bf06b4f8
|
||||
diff_cmd: git diff HEAD~1..HEAD
|
||||
---
|
||||
|
||||
Implement `to_actor.run()` with Trio-style positional target arguments,
|
||||
`functools.partial` keyword and Python 3.14 Placeholder binding, and a
|
||||
static context endpoint that links remote results, errors and caller
|
||||
cancellation.
|
||||
|
||||
> `git diff HEAD~1..HEAD -- tractor/to_actor/_api.py tractor/to_actor/__init__.py`
|
||||
|
||||
Resolve target functions through `Actor._get_rpc_func()` so module
|
||||
authorization remains authoritative. Automatically expose the helper
|
||||
module for actors owned by `to_actor.run()` and document explicit
|
||||
exposure for a caller-owned portal.
|
||||
|
||||
> `git diff HEAD~1..HEAD -- tests/test_to_actor.py examples/parallelism/to_actor_one_shots.py`
|
||||
|
||||
Cover placement modes, argument binding, nested partials, caller-linked
|
||||
cancellation, remote errors and module authorization.
|
||||
|
|
@ -1,34 +0,0 @@
|
|||
---
|
||||
model: openai/gpt-5.6-sol
|
||||
service: opencode
|
||||
session: ses_3e4c90d3eafeqHEtRYSIHgHhpA
|
||||
timestamp: 2026-08-19T02:07:57Z
|
||||
git_ref: b38efed7
|
||||
scope: code
|
||||
substantive: true
|
||||
raw_file: 20260819T020757Z_b38efed7_prompt_io.raw.md
|
||||
---
|
||||
|
||||
## Prompt
|
||||
|
||||
Resolve the remaining P3 review finding before landing PR #481:
|
||||
`runtime_kwargs={}` must not be silently accepted alongside either
|
||||
`an=` or `portal=` merely because the dict is falsey. Keep this as its
|
||||
own final review-remediation commit.
|
||||
|
||||
## Response summary
|
||||
|
||||
Treat any non-`None` `runtime_kwargs` value as provided when validating
|
||||
placement. Cover both placement APIs with empty and configured dicts,
|
||||
proving the error is raised locally before actor startup.
|
||||
|
||||
## Files changed
|
||||
|
||||
- `tractor/to_actor/_api.py` - validate option presence explicitly.
|
||||
- `tests/test_to_actor.py` - cover four invalid option combinations.
|
||||
|
||||
## Human edits
|
||||
|
||||
No direct line edits. The human accepted the P3 finding, required it to
|
||||
remain separate from the five P2 behavioral commits and prioritized it
|
||||
before the final PR #484 integration rebase and PR #481 landing steps.
|
||||
|
|
@ -1,25 +0,0 @@
|
|||
---
|
||||
model: openai/gpt-5.6-sol
|
||||
service: opencode
|
||||
timestamp: 2026-08-19T02:07:57Z
|
||||
git_ref: b38efed7
|
||||
diff_cmd: git diff HEAD~1..HEAD
|
||||
---
|
||||
|
||||
Fix the final PR #481 review finding: `runtime_kwargs` is mutually
|
||||
exclusive with both caller placement options whenever it is provided,
|
||||
including an empty dict.
|
||||
|
||||
> `git diff HEAD~1..HEAD -- tractor/to_actor/_api.py tests/test_to_actor.py`
|
||||
|
||||
Use an explicit `is not None` check rather than dict truthiness. Expand
|
||||
the validation regression across `an=` and `portal=`, each with empty
|
||||
and configured runtime kwargs, so every invalid combination fails
|
||||
before actor runtime startup.
|
||||
|
||||
Verification:
|
||||
|
||||
- Trio/TCP: `23 passed`
|
||||
- Trio/UDS: `23 passed`
|
||||
- `mp_spawn`/TCP: `23 passed`
|
||||
- Ruff and `git diff --check`: clean
|
||||
|
|
@ -41,7 +41,7 @@ async def main() -> None:
|
|||
# subactor, tears the runtime back down.
|
||||
assert await tractor.to_actor.run(
|
||||
is_prime,
|
||||
2,
|
||||
n=2,
|
||||
)
|
||||
|
||||
# the "worker-pool-ish" pattern from the original
|
||||
|
|
@ -57,9 +57,9 @@ async def main() -> None:
|
|||
) -> None:
|
||||
results[n] = await tractor.to_actor.run(
|
||||
is_prime,
|
||||
n,
|
||||
an=an,
|
||||
name=f'prime_checker_{i}',
|
||||
n=n,
|
||||
)
|
||||
|
||||
inputs: list[int] = [
|
||||
|
|
|
|||
|
|
@ -17,65 +17,9 @@ import trio
|
|||
import tractor
|
||||
from tractor import Actor
|
||||
from tractor.discovery import _addr
|
||||
from tractor.ipc._transport import MsgpackTransport
|
||||
from tractor.runtime import _state
|
||||
|
||||
|
||||
|
||||
def test_cancelled_transport_send_closes_stream():
|
||||
'''
|
||||
Discard a transport after cancellation interrupts a framed send.
|
||||
|
||||
Trio's `SendStream.send_all()` may write an arbitrary frame prefix
|
||||
before raising `Cancelled`. Sending another IPC msg afterward
|
||||
would append a second frame and desynchronize the peer decoder.
|
||||
The fake stream checkpoints after recording send entry; cancelling
|
||||
its nursery deterministically interrupts that unknown-publication
|
||||
window. Its close assertion proves the transport is made unusable
|
||||
before another framed msg can be attempted.
|
||||
|
||||
'''
|
||||
class PartialSendStream:
|
||||
def __init__(self) -> None:
|
||||
self.send_entered = trio.Event()
|
||||
self.closed = False
|
||||
|
||||
async def send_all(
|
||||
self,
|
||||
data: bytes,
|
||||
) -> None:
|
||||
assert data
|
||||
self.send_entered.set()
|
||||
await trio.sleep_forever()
|
||||
|
||||
async def aclose(self) -> None:
|
||||
self.closed = True
|
||||
|
||||
async def main() -> None:
|
||||
stream = PartialSendStream()
|
||||
transport = object.__new__(MsgpackTransport)
|
||||
transport.stream = stream
|
||||
transport._send_lock = trio.StrictFIFOLock()
|
||||
|
||||
async with trio.open_nursery() as tn:
|
||||
tn.start_soon(
|
||||
transport.send,
|
||||
tractor.msg.Start(
|
||||
ns=__name__,
|
||||
func='add_one',
|
||||
kwargs={'n': 1},
|
||||
uid=('root', 'test'),
|
||||
cid='partial-send',
|
||||
),
|
||||
)
|
||||
await stream.send_entered.wait()
|
||||
tn.cancel_scope.cancel()
|
||||
|
||||
assert stream.closed
|
||||
|
||||
trio.run(main)
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def bindspace_dir_str() -> str:
|
||||
|
||||
|
|
|
|||
|
|
@ -7,7 +7,6 @@ sync-opening a ``tractor.Context`` beforehand.
|
|||
'''
|
||||
from itertools import count
|
||||
import math
|
||||
from pathlib import Path
|
||||
import platform
|
||||
from pprint import pformat
|
||||
import sys
|
||||
|
|
@ -76,37 +75,6 @@ from tractor._testing import (
|
|||
_state: bool = False
|
||||
|
||||
|
||||
def _non_registration_contexts(
|
||||
actor: Actor,
|
||||
) -> dict[tuple, str]:
|
||||
return {
|
||||
key: str(ctx._nsf)
|
||||
for key, ctx in actor._contexts.items()
|
||||
if str(ctx._nsf) != (
|
||||
'tractor.discovery._registry:'
|
||||
'Registrar.register_actor'
|
||||
)
|
||||
}
|
||||
|
||||
|
||||
@tractor.context
|
||||
async def startup_cancel_target(
|
||||
ctx: Context,
|
||||
started_path: str,
|
||||
cancelled_path: str,
|
||||
) -> None:
|
||||
Path(started_path).touch()
|
||||
try:
|
||||
await ctx.started()
|
||||
await trio.sleep_forever()
|
||||
finally:
|
||||
Path(cancelled_path).touch()
|
||||
|
||||
|
||||
async def return_one() -> int:
|
||||
return 1
|
||||
|
||||
|
||||
@tractor.context
|
||||
async def too_many_starteds(
|
||||
ctx: Context,
|
||||
|
|
@ -200,168 +168,6 @@ async def assert_state(value: bool):
|
|||
assert _state == value
|
||||
|
||||
|
||||
@tractor_test
|
||||
async def test_cancel_during_context_startup(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
tmp_path: Path,
|
||||
start_method: str,
|
||||
debug_mode: bool,
|
||||
):
|
||||
'''
|
||||
Cancel a context after sending `Start` but before its ack.
|
||||
|
||||
`Portal.open_context()` allocates its caller-side `Context` while
|
||||
entering the async context manager. Cancellation used to strand
|
||||
that local context and leave the remote target running. The patched
|
||||
`Channel.send()` publishes `Start`, then blocks before
|
||||
`Actor.start_remote_task()` can await `StartAck`. Cancelling the
|
||||
caller proves cleanup issues one bounded, non-recursive cancel RPC,
|
||||
stops the target and removes both helper contexts. A subsequent
|
||||
RPC proves the caller-owned actor remains usable.
|
||||
|
||||
'''
|
||||
started_path = tmp_path / 'startup_started'
|
||||
cancelled_path = tmp_path / 'startup_cancelled'
|
||||
start_sent = trio.Event()
|
||||
original_send = tractor.Channel.send
|
||||
|
||||
async def delay_after_start(
|
||||
chan: tractor.Channel,
|
||||
payload: object,
|
||||
hide_tb: bool = False,
|
||||
) -> None:
|
||||
await original_send(
|
||||
chan,
|
||||
payload,
|
||||
hide_tb=hide_tb,
|
||||
)
|
||||
if isinstance(payload, tractor.msg.Start):
|
||||
if payload.func == 'startup_cancel_target':
|
||||
start_sent.set()
|
||||
await trio.sleep_forever()
|
||||
|
||||
async def open_target(
|
||||
portal: tractor.Portal,
|
||||
) -> None:
|
||||
async with portal.open_context(
|
||||
startup_cancel_target,
|
||||
started_path=str(started_path),
|
||||
cancelled_path=str(cancelled_path),
|
||||
):
|
||||
raise AssertionError('context startup should be cancelled')
|
||||
|
||||
async with tractor.open_nursery() as an:
|
||||
actor = tractor.current_actor()
|
||||
portal: tractor.Portal = await an.start_actor(
|
||||
'startup_cancel_worker',
|
||||
enable_modules=[__name__],
|
||||
)
|
||||
contexts_before = _non_registration_contexts(actor)
|
||||
monkeypatch.setattr(
|
||||
tractor.Channel,
|
||||
'send',
|
||||
delay_after_start,
|
||||
)
|
||||
|
||||
async with trio.open_nursery() as tn:
|
||||
tn.start_soon(open_target, portal)
|
||||
with trio.fail_after(5):
|
||||
await start_sent.wait()
|
||||
while not started_path.exists():
|
||||
await trio.sleep(0.01)
|
||||
tn.cancel_scope.cancel()
|
||||
|
||||
monkeypatch.setattr(
|
||||
tractor.Channel,
|
||||
'send',
|
||||
original_send,
|
||||
)
|
||||
assert cancelled_path.exists()
|
||||
assert _non_registration_contexts(actor) == contexts_before
|
||||
assert await portal.run_from_ns(
|
||||
__name__,
|
||||
'return_one',
|
||||
) == 1
|
||||
assert _non_registration_contexts(actor) == contexts_before
|
||||
await portal.cancel_actor()
|
||||
|
||||
|
||||
@tractor_test
|
||||
async def test_start_serialization_error_cleans_context(
|
||||
start_method: str,
|
||||
debug_mode: bool,
|
||||
):
|
||||
'''
|
||||
Deallocate caller state when `Start` can not be serialized.
|
||||
|
||||
`Actor.start_remote_task()` registers its caller-side `Context`
|
||||
before encoding the request. An unsupported argument used to raise
|
||||
`MsgTypeError` before publication while leaking that registry
|
||||
entry. Comparing the context registry around the failed start
|
||||
proves cleanup, and a following valid context proves no bytes
|
||||
reached or damaged the reused portal's transport.
|
||||
|
||||
'''
|
||||
async with tractor.open_nursery() as an:
|
||||
actor = tractor.current_actor()
|
||||
portal: tractor.Portal = await an.start_actor(
|
||||
'serialization_error_worker',
|
||||
enable_modules=[__name__],
|
||||
)
|
||||
contexts_before = _non_registration_contexts(actor)
|
||||
with pytest.raises(tractor.MsgTypeError):
|
||||
async with portal.open_context(
|
||||
simple_setup_teardown,
|
||||
data=object(),
|
||||
):
|
||||
raise AssertionError('invalid `Start` was accepted')
|
||||
|
||||
assert _non_registration_contexts(actor) == contexts_before
|
||||
async with portal.open_context(
|
||||
simple_setup_teardown,
|
||||
data=1,
|
||||
) as (ctx, started):
|
||||
assert started == 2
|
||||
assert await ctx.wait_for_result() == 'yo'
|
||||
|
||||
assert _non_registration_contexts(actor) == contexts_before
|
||||
await portal.cancel_actor()
|
||||
|
||||
|
||||
@tractor_test
|
||||
async def test_start_module_error_cleans_context(
|
||||
start_method: str,
|
||||
debug_mode: bool,
|
||||
):
|
||||
'''
|
||||
Deallocate caller state after a remote startup rejection.
|
||||
|
||||
A target actor without this test module rejects the requested
|
||||
context before sending `StartAck`. That remote
|
||||
`ModuleNotExposed` used to escape startup validation while leaving
|
||||
the caller context registered. The boxed error and before/after
|
||||
registry comparison prove the remote failure remains visible and
|
||||
local startup state is released.
|
||||
|
||||
'''
|
||||
async with tractor.open_nursery() as an:
|
||||
actor = tractor.current_actor()
|
||||
portal: tractor.Portal = await an.start_actor(
|
||||
'module_error_worker',
|
||||
)
|
||||
contexts_before = _non_registration_contexts(actor)
|
||||
with pytest.raises(tractor.RemoteActorError) as excinfo:
|
||||
async with portal.open_context(
|
||||
simple_setup_teardown,
|
||||
data=1,
|
||||
):
|
||||
raise AssertionError('unexposed context was started')
|
||||
|
||||
assert excinfo.value.boxed_type is tractor.ModuleNotExposed
|
||||
assert _non_registration_contexts(actor) == contexts_before
|
||||
await portal.cancel_actor()
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
'error_parent',
|
||||
[False, ValueError, KeyboardInterrupt],
|
||||
|
|
|
|||
|
|
@ -7,7 +7,6 @@ https://github.com/goodboy/tractor/issues/477
|
|||
|
||||
'''
|
||||
from functools import partial
|
||||
from pathlib import Path
|
||||
|
||||
import pytest
|
||||
import trio
|
||||
|
|
@ -17,9 +16,6 @@ from tractor import (
|
|||
to_actor,
|
||||
)
|
||||
from tractor._testing import tractor_test
|
||||
from tractor.msg import ptr as msgptr
|
||||
from tractor.msg.ptr import NamespacePath
|
||||
from tractor.to_actor import _api as to_actor_api
|
||||
|
||||
|
||||
async def add_one(
|
||||
|
|
@ -32,98 +28,6 @@ async def raise_value_error() -> None:
|
|||
raise ValueError('kaboom')
|
||||
|
||||
|
||||
async def echo_control_names(
|
||||
value: int,
|
||||
/,
|
||||
*,
|
||||
name: str,
|
||||
portal: str,
|
||||
an: str,
|
||||
runtime_kwargs: str,
|
||||
) -> dict[str, int|str]:
|
||||
return {
|
||||
'value': value,
|
||||
'name': name,
|
||||
'portal': portal,
|
||||
'an': an,
|
||||
'runtime_kwargs': runtime_kwargs,
|
||||
}
|
||||
|
||||
|
||||
async def mark_task_cancellation(
|
||||
started_path: str,
|
||||
cancelled_path: str,
|
||||
) -> None:
|
||||
Path(started_path).touch()
|
||||
try:
|
||||
await trio.sleep_forever()
|
||||
finally:
|
||||
Path(cancelled_path).touch()
|
||||
|
||||
|
||||
async def echo_startup_control(
|
||||
_cancel_on_startup: str,
|
||||
) -> str:
|
||||
return _cancel_on_startup
|
||||
|
||||
|
||||
async def collect_args(
|
||||
*args: object,
|
||||
) -> tuple[object, ...]:
|
||||
return args
|
||||
|
||||
|
||||
async def collect_call(
|
||||
*args: object,
|
||||
**kwargs: object,
|
||||
) -> tuple[tuple[object, ...], dict[str, object]]:
|
||||
return args, kwargs
|
||||
|
||||
|
||||
def _non_registration_contexts(
|
||||
actor: tractor.Actor,
|
||||
) -> dict[tuple, str]:
|
||||
return {
|
||||
key: str(ctx._nsf)
|
||||
for key, ctx in actor._contexts.items()
|
||||
if str(ctx._nsf) != (
|
||||
'tractor.discovery._registry:'
|
||||
'Registrar.register_actor'
|
||||
)
|
||||
}
|
||||
|
||||
|
||||
def test_namespace_path_retains_target_ref(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
):
|
||||
'''
|
||||
Reuse the client-side target ref when splitting its namespace path.
|
||||
|
||||
`NamespacePath.from_ref()` previously discarded `add_one`, so
|
||||
`to_tuple()` imported and resolved the just-created string again.
|
||||
Replacing `resolve_name()` with a failure proves the retained ref
|
||||
supplies the tuple without a redundant lookup. The public module
|
||||
alias assertion also keeps internal `_api.__name__` authoritative.
|
||||
|
||||
'''
|
||||
target = NamespacePath.from_ref(add_one)
|
||||
|
||||
def fail_resolve(name: str) -> object:
|
||||
raise AssertionError(f'unexpected lookup for {name!r}')
|
||||
|
||||
monkeypatch.setattr(
|
||||
msgptr,
|
||||
'resolve_name',
|
||||
fail_resolve,
|
||||
)
|
||||
assert target.to_tuple() == (
|
||||
add_one.__module__,
|
||||
add_one.__name__,
|
||||
)
|
||||
assert to_actor.MODULE == to_actor_api.__name__
|
||||
assert not hasattr(to_actor_api, 'MODULE')
|
||||
|
||||
|
||||
@tractor_test
|
||||
async def test_one_shot_in_private_nursery(
|
||||
start_method: str,
|
||||
|
|
@ -136,7 +40,7 @@ async def test_one_shot_in_private_nursery(
|
|||
'''
|
||||
assert await to_actor.run(
|
||||
add_one,
|
||||
1,
|
||||
n=1,
|
||||
) == 2
|
||||
|
||||
|
||||
|
|
@ -157,7 +61,7 @@ def test_one_shot_boots_implicit_runtime(
|
|||
) is None
|
||||
result = await to_actor.run(
|
||||
add_one,
|
||||
41,
|
||||
n=41,
|
||||
runtime_kwargs=dict(
|
||||
registry_addrs=[reg_addr],
|
||||
start_method=start_method,
|
||||
|
|
@ -193,90 +97,17 @@ async def test_spawn_from_caller_nursery(
|
|||
debug_mode: bool,
|
||||
):
|
||||
'''
|
||||
Pass a caller-managed `an: ActorNursery` for the spawn.
|
||||
|
||||
Previously `to_actor.run()` treated an actor-runtime cancel ack
|
||||
as process reaping, so the call returned while the child monitor
|
||||
and its `ActorNursery._children` record remained alive until the
|
||||
entire nursery exited. The assertion inside the still-open
|
||||
nursery proves child-process joining and record removal now
|
||||
complete before the one-shot call returns.
|
||||
Pass a caller-managed `an: ActorNursery` for the
|
||||
spawn; the subactor is still one-shot reaped by the
|
||||
time the call returns.
|
||||
|
||||
'''
|
||||
async with tractor.open_nursery() as an:
|
||||
assert await to_actor.run(
|
||||
add_one,
|
||||
10,
|
||||
an=an,
|
||||
n=10,
|
||||
) == 11
|
||||
assert not an._children
|
||||
|
||||
|
||||
@tractor_test
|
||||
async def test_cancel_ack_failure_hard_reaps_child(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
start_method: str,
|
||||
debug_mode: bool,
|
||||
):
|
||||
'''
|
||||
Escalate a failed cancel acknowledgement and reap the child.
|
||||
|
||||
`Portal.cancel_actor()` can return `False` when its transport is
|
||||
already closed without confirming runtime cancellation. The old
|
||||
one-shot path ignored that result, released the nursery-wide join
|
||||
gate and then waited forever for a still-running process. This
|
||||
test forces that exact result without cancelling the actor, caps
|
||||
the call to detect the former hang and verifies the child monitor
|
||||
removes its `ActorNursery._children` record before returning.
|
||||
|
||||
'''
|
||||
async def cancel_without_ack(
|
||||
portal: tractor.Portal,
|
||||
timeout: float|None = None,
|
||||
raise_on_timeout: bool = False,
|
||||
) -> bool:
|
||||
assert raise_on_timeout
|
||||
return False
|
||||
|
||||
monkeypatch.setattr(
|
||||
tractor.Portal,
|
||||
'cancel_actor',
|
||||
cancel_without_ack,
|
||||
)
|
||||
|
||||
async with tractor.open_nursery() as an:
|
||||
with trio.fail_after(5):
|
||||
assert await to_actor.run(
|
||||
add_one,
|
||||
20,
|
||||
an=an,
|
||||
) == 21
|
||||
assert not an._children
|
||||
|
||||
|
||||
def test_late_child_reap_registration_is_released():
|
||||
'''
|
||||
Preserve a nursery-wide reap request across child startup.
|
||||
|
||||
A child monitor can checkpoint while connecting to its parent as
|
||||
the surrounding `ActorNursery` begins teardown. Previously the
|
||||
nursery signalled only already-registered child events, so a
|
||||
monitor registering afterward waited forever. This models that
|
||||
ordering by publishing the nursery-wide request first and proves
|
||||
the later per-child event inherits its set state immediately.
|
||||
|
||||
'''
|
||||
an = object.__new__(tractor.ActorNursery)
|
||||
an._join_procs = trio.Event()
|
||||
an._child_reap_requests = {}
|
||||
an._child_reaped = {}
|
||||
|
||||
an._join_procs.set()
|
||||
reap_request, _ = an._register_child_reap(
|
||||
('late_child', 'uid'),
|
||||
)
|
||||
|
||||
assert reap_request.is_set()
|
||||
|
||||
|
||||
@tractor_test
|
||||
|
|
@ -309,35 +140,19 @@ async def test_reuse_existing_actor_via_portal(
|
|||
Pass `portal=` to schedule the one-shot task in an
|
||||
already-running actor; no spawn, no implicit reap.
|
||||
|
||||
The low-level `Portal.run_from_ns()` assertion also proves its
|
||||
target kwargs remain separate from the private startup-cancel
|
||||
policy used by context cleanup.
|
||||
|
||||
'''
|
||||
async with tractor.open_nursery() as an:
|
||||
actor = tractor.current_actor()
|
||||
portal: tractor.Portal = await an.start_actor(
|
||||
'one_shot_worker',
|
||||
enable_modules=[
|
||||
__name__,
|
||||
to_actor.MODULE,
|
||||
],
|
||||
enable_modules=[__name__],
|
||||
)
|
||||
contexts_before = _non_registration_contexts(actor)
|
||||
for i in range(3):
|
||||
assert await to_actor.run(
|
||||
add_one,
|
||||
i,
|
||||
portal=portal,
|
||||
n=i,
|
||||
) == i + 1
|
||||
|
||||
assert await portal.run_from_ns(
|
||||
__name__,
|
||||
'echo_startup_control',
|
||||
_cancel_on_startup='target_value',
|
||||
) == 'target_value'
|
||||
assert _non_registration_contexts(actor) == contexts_before
|
||||
|
||||
# still alive: caller owns the actor's lifetime.
|
||||
await portal.cancel_actor()
|
||||
|
||||
|
|
@ -363,9 +178,9 @@ async def test_concurrent_one_shots_from_task_nursery(
|
|||
) -> None:
|
||||
results[i] = await to_actor.run(
|
||||
add_one,
|
||||
i,
|
||||
an=an,
|
||||
name=f'one_shot_{i}',
|
||||
n=i,
|
||||
)
|
||||
|
||||
async with (
|
||||
|
|
@ -416,83 +231,6 @@ def test_rejects_streaming_fn():
|
|||
)
|
||||
|
||||
|
||||
def test_partial_placeholder_normalization(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
):
|
||||
'''
|
||||
Preserve Python 3.14 `functools.partial` placeholder semantics.
|
||||
|
||||
The test environment runs Python 3.13, so this installs an identity
|
||||
sentinel matching Python 3.14's `functools.Placeholder` API.
|
||||
Interleaved placeholders prove call-time positional arguments are
|
||||
merged in order. Undersupply and a mismatched final target
|
||||
signature both fail locally before actor runtime startup.
|
||||
|
||||
'''
|
||||
placeholder = object()
|
||||
monkeypatch.setattr(
|
||||
to_actor_api.functools,
|
||||
'Placeholder',
|
||||
placeholder,
|
||||
raising=False,
|
||||
)
|
||||
fn = partial(
|
||||
collect_args,
|
||||
placeholder,
|
||||
2,
|
||||
placeholder,
|
||||
)
|
||||
normalized_fn, args, kwargs = to_actor_api._normalize_call(
|
||||
fn,
|
||||
(1, 3, 4),
|
||||
)
|
||||
assert normalized_fn is collect_args
|
||||
assert args == (1, 2, 3, 4)
|
||||
assert kwargs == {}
|
||||
|
||||
with pytest.raises(TypeError, match='Not enough positional'):
|
||||
to_actor_api._normalize_call(fn, (1,))
|
||||
|
||||
with pytest.raises(TypeError, match='too many positional'):
|
||||
to_actor_api._normalize_call(
|
||||
partial(add_one, 1),
|
||||
(2,),
|
||||
)
|
||||
|
||||
|
||||
def test_nested_partial_normalization():
|
||||
'''
|
||||
Flatten every retained `functools.partial` layer before RPC.
|
||||
|
||||
CPython normally combines nested partials, but preserves the inner
|
||||
object when it has instance attributes. Unwrapping only the outer
|
||||
layer left a non-namespace-addressable partial as the RPC target.
|
||||
The custom attribute triggers that retained shape; the assertions
|
||||
prove positional ordering and outer-keyword precedence match a
|
||||
direct nested-partial call.
|
||||
|
||||
'''
|
||||
inner = partial(
|
||||
collect_call,
|
||||
1,
|
||||
label='inner',
|
||||
)
|
||||
inner.note = 'retain this partial layer'
|
||||
outer = partial(
|
||||
inner,
|
||||
2,
|
||||
label='outer',
|
||||
)
|
||||
|
||||
fn, args, kwargs = to_actor_api._normalize_call(
|
||||
outer,
|
||||
(3,),
|
||||
)
|
||||
assert fn is collect_call
|
||||
assert args == (1, 2, 3)
|
||||
assert kwargs == {'label': 'outer'}
|
||||
|
||||
|
||||
def test_rejects_portal_and_an_combo():
|
||||
'''
|
||||
`portal=` and `an=` are mutually exclusive
|
||||
|
|
@ -504,37 +242,19 @@ def test_rejects_portal_and_an_combo():
|
|||
partial(
|
||||
to_actor.run,
|
||||
add_one,
|
||||
1,
|
||||
portal=object(),
|
||||
an=object(),
|
||||
n=1,
|
||||
)
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
'placement',
|
||||
['an', 'portal'],
|
||||
)
|
||||
@pytest.mark.parametrize(
|
||||
'runtime_kwargs',
|
||||
[
|
||||
{},
|
||||
{'loglevel': 'cancel'},
|
||||
],
|
||||
ids=['empty', 'configured'],
|
||||
)
|
||||
def test_rejects_runtime_kwargs_with_placement(
|
||||
placement: str,
|
||||
runtime_kwargs: dict,
|
||||
):
|
||||
def test_rejects_runtime_kwargs_with_placement():
|
||||
'''
|
||||
`runtime_kwargs` only applies when the call opens
|
||||
its own private actor-nursery; passing it alongside
|
||||
a placement opt is an error, never silently
|
||||
ignored. In particular, an empty dict still means the
|
||||
caller provided this mutually exclusive option; testing
|
||||
both placement modes prevents truthiness checks from
|
||||
accepting it before any actor runtime is started.
|
||||
ignored.
|
||||
|
||||
'''
|
||||
with pytest.raises(ValueError):
|
||||
|
|
@ -542,179 +262,10 @@ def test_rejects_runtime_kwargs_with_placement(
|
|||
partial(
|
||||
to_actor.run,
|
||||
add_one,
|
||||
1,
|
||||
**{
|
||||
placement: object(),
|
||||
'runtime_kwargs': runtime_kwargs,
|
||||
},
|
||||
)
|
||||
)
|
||||
|
||||
|
||||
@tractor_test
|
||||
async def test_trio_style_args_and_partial_kwargs(
|
||||
start_method: str,
|
||||
debug_mode: bool,
|
||||
):
|
||||
'''
|
||||
Forward positional args and partial-bound keyword arguments.
|
||||
|
||||
The original API captured every keyword matching an actor
|
||||
control, so ordinary target parameters such as `name`, `portal`,
|
||||
`an` and `runtime_kwargs` could not be called. This test uses a
|
||||
positional-only target argument plus all colliding keyword names.
|
||||
Binding the target keywords with `functools.partial()` proves the
|
||||
Trio-style calling convention keeps target inputs separate from
|
||||
actor controls.
|
||||
|
||||
'''
|
||||
fn = partial(
|
||||
echo_control_names,
|
||||
name='target_name',
|
||||
portal='target_portal',
|
||||
an='target_an',
|
||||
runtime_kwargs='target_runtime_kwargs',
|
||||
)
|
||||
async with tractor.open_nursery() as an:
|
||||
result = await to_actor.run(
|
||||
fn,
|
||||
42,
|
||||
an=an,
|
||||
name='actor_name',
|
||||
)
|
||||
|
||||
assert result == {
|
||||
'value': 42,
|
||||
'name': 'target_name',
|
||||
'portal': 'target_portal',
|
||||
'an': 'target_an',
|
||||
'runtime_kwargs': 'target_runtime_kwargs',
|
||||
}
|
||||
|
||||
|
||||
@tractor_test
|
||||
async def test_portal_task_cancelled_with_local_caller(
|
||||
tmp_path: Path,
|
||||
start_method: str,
|
||||
debug_mode: bool,
|
||||
):
|
||||
'''
|
||||
Couple a reused portal's remote task to its local caller.
|
||||
|
||||
The former `Portal.run()` path abandoned its remote task when the
|
||||
local `to_actor.run()` caller was cancelled. The target writes
|
||||
one file after starting and another from its cancellation
|
||||
`finally`. Cancelling the local task nursery and observing the
|
||||
second file proves `Portal.open_context()` propagated
|
||||
cancellation before the caller exited. A subsequent call proves
|
||||
the caller-owned actor was not cancelled with that task.
|
||||
|
||||
'''
|
||||
started_path = tmp_path / 'started'
|
||||
cancelled_path = tmp_path / 'cancelled'
|
||||
|
||||
async with tractor.open_nursery() as an:
|
||||
actor = tractor.current_actor()
|
||||
portal: tractor.Portal = await an.start_actor(
|
||||
'context_worker',
|
||||
enable_modules=[
|
||||
__name__,
|
||||
to_actor.MODULE,
|
||||
],
|
||||
)
|
||||
contexts_before = _non_registration_contexts(actor)
|
||||
|
||||
async with trio.open_nursery() as tn:
|
||||
tn.start_soon(
|
||||
partial(
|
||||
to_actor.run,
|
||||
mark_task_cancellation,
|
||||
str(started_path),
|
||||
str(cancelled_path),
|
||||
portal=portal,
|
||||
an=object(),
|
||||
runtime_kwargs=dict(
|
||||
loglevel='cancel',
|
||||
),
|
||||
n=1,
|
||||
)
|
||||
with trio.fail_after(5):
|
||||
while not started_path.exists():
|
||||
await trio.sleep(0.01)
|
||||
tn.cancel_scope.cancel()
|
||||
|
||||
assert cancelled_path.exists()
|
||||
assert _non_registration_contexts(actor) == contexts_before
|
||||
assert await to_actor.run(
|
||||
add_one,
|
||||
1,
|
||||
portal=portal,
|
||||
) == 2
|
||||
assert _non_registration_contexts(actor) == contexts_before
|
||||
|
||||
await portal.cancel_actor()
|
||||
|
||||
|
||||
@tractor_test
|
||||
async def test_context_trampoline_preserves_module_allowlist(
|
||||
start_method: str,
|
||||
debug_mode: bool,
|
||||
):
|
||||
'''
|
||||
Keep target resolution behind the actor's RPC module allowlist.
|
||||
|
||||
Loading the target with `NamespacePath.load_ref()` would silently
|
||||
bypass the actor's existing module-exposure boundary. This actor
|
||||
exposes only the trusted trampoline, not the test module; the
|
||||
boxed `ModuleNotExposed` proves the trampoline delegates target
|
||||
resolution to `Actor._get_rpc_func()`.
|
||||
|
||||
'''
|
||||
async with tractor.open_nursery() as an:
|
||||
actor = tractor.current_actor()
|
||||
portal: tractor.Portal = await an.start_actor(
|
||||
'restricted_context_worker',
|
||||
enable_modules=[to_actor.MODULE],
|
||||
)
|
||||
contexts_before = _non_registration_contexts(actor)
|
||||
with pytest.raises(RemoteActorError) as excinfo:
|
||||
await to_actor.run(
|
||||
add_one,
|
||||
1,
|
||||
portal=portal,
|
||||
)
|
||||
|
||||
assert excinfo.value.boxed_type is tractor.ModuleNotExposed
|
||||
assert _non_registration_contexts(actor) == contexts_before
|
||||
await portal.cancel_actor()
|
||||
|
||||
|
||||
@tractor_test
|
||||
async def test_portal_requires_context_trampoline(
|
||||
start_method: str,
|
||||
debug_mode: bool,
|
||||
):
|
||||
'''
|
||||
Require explicit trampoline exposure on a caller-owned actor.
|
||||
|
||||
Automatically exposing the module in every actor weakens the RPC
|
||||
allowlist for actors that never use `to_actor.run()`. A portal to
|
||||
such an actor instead fails with the usual `ModuleNotExposed`,
|
||||
naming the module callers must opt into.
|
||||
|
||||
'''
|
||||
async with tractor.open_nursery() as an:
|
||||
actor = tractor.current_actor()
|
||||
portal: tractor.Portal = await an.start_actor(
|
||||
'no_context_trampoline_worker',
|
||||
enable_modules=[__name__],
|
||||
)
|
||||
contexts_before = _non_registration_contexts(actor)
|
||||
with pytest.raises(RemoteActorError) as excinfo:
|
||||
await to_actor.run(
|
||||
add_one,
|
||||
1,
|
||||
portal=portal,
|
||||
)
|
||||
|
||||
err = excinfo.value
|
||||
assert err.boxed_type is tractor.ModuleNotExposed
|
||||
assert to_actor.MODULE in str(err)
|
||||
assert _non_registration_contexts(actor) == contexts_before
|
||||
await portal.cancel_actor()
|
||||
|
|
|
|||
|
|
@ -1108,11 +1108,10 @@ class Context:
|
|||
# NOTE: we're telling the far end actor to cancel a task
|
||||
# corresponding to *this actor*. The far end local channel
|
||||
# instance is passed to `Actor._cancel_task()` implicitly.
|
||||
await self._portal._run_from_ns(
|
||||
await self._portal.run_from_ns(
|
||||
'self',
|
||||
'_cancel_task',
|
||||
kwargs={'cid': cid},
|
||||
cancel_on_startup=False,
|
||||
cid=cid,
|
||||
)
|
||||
|
||||
if cs.cancelled_caught:
|
||||
|
|
@ -2626,7 +2625,10 @@ async def open_context_from_portal(
|
|||
f'uid: {uid}\n'
|
||||
f'cid: {ctx.cid}\n'
|
||||
)
|
||||
portal.actor._drop_context(ctx)
|
||||
portal.actor._contexts.pop(
|
||||
(uid, ctx.cid),
|
||||
None,
|
||||
)
|
||||
|
||||
# XXX revert to prior IPC-task-ctx scope
|
||||
_ctxvar_Context.reset(prior_ctx_tok)
|
||||
|
|
|
|||
|
|
@ -499,12 +499,6 @@ class MsgpackTransport(MsgTransport):
|
|||
size: bytes = struct.pack("<I", len(bytes_data))
|
||||
try:
|
||||
return await self.stream.send_all(size + bytes_data)
|
||||
except trio.Cancelled:
|
||||
# `send_all()` may have written a partial frame. The
|
||||
# stream can not safely carry another framed msg.
|
||||
with trio.CancelScope(shield=True):
|
||||
await self.stream.aclose()
|
||||
raise
|
||||
except (
|
||||
trio.BrokenResourceError,
|
||||
trio.ClosedResourceError,
|
||||
|
|
|
|||
|
|
@ -125,9 +125,7 @@ class NamespacePath(str):
|
|||
) -> NamespacePath:
|
||||
|
||||
fqnp: tuple[str, str] = cls._mk_fqnp(ref)
|
||||
nsp = cls(':'.join(fqnp))
|
||||
nsp._ref = ref
|
||||
return nsp
|
||||
return cls(':'.join(fqnp))
|
||||
|
||||
def to_tuple(
|
||||
self,
|
||||
|
|
|
|||
|
|
@ -273,38 +273,6 @@ class Portal:
|
|||
|
||||
return False
|
||||
|
||||
async def _run_from_ns(
|
||||
self,
|
||||
namespace_path: str,
|
||||
function_name: str,
|
||||
kwargs: dict[str, Any],
|
||||
cancel_on_startup: bool = True,
|
||||
) -> Any:
|
||||
'''
|
||||
Run a namespace target with local startup policy controls.
|
||||
|
||||
'''
|
||||
nsf = NamespacePath(
|
||||
f'{namespace_path}:{function_name}'
|
||||
)
|
||||
ctx: Context = await self.actor.start_remote_task(
|
||||
chan=self.channel,
|
||||
nsf=nsf,
|
||||
kwargs=kwargs,
|
||||
portal=self,
|
||||
cancel_on_startup=cancel_on_startup,
|
||||
)
|
||||
try:
|
||||
return await ctx._pld_rx.recv_pld(
|
||||
ipc=ctx,
|
||||
expect_msg=Return,
|
||||
)
|
||||
finally:
|
||||
self.actor._drop_context(ctx)
|
||||
if not ctx._rx_chan._closed:
|
||||
with trio.CancelScope(shield=True):
|
||||
await ctx._rx_chan.aclose()
|
||||
|
||||
# TODO: do we still need this for low level `Actor`-runtime
|
||||
# method calls or can we also remove it?
|
||||
async def run_from_ns(
|
||||
|
|
@ -330,10 +298,18 @@ class Portal:
|
|||
|
||||
'''
|
||||
__runtimeframe__: int = 1 # noqa
|
||||
return await self._run_from_ns(
|
||||
namespace_path,
|
||||
function_name,
|
||||
nsf = NamespacePath(
|
||||
f'{namespace_path}:{function_name}'
|
||||
)
|
||||
ctx: Context = await self.actor.start_remote_task(
|
||||
chan=self.channel,
|
||||
nsf=nsf,
|
||||
kwargs=kwargs,
|
||||
portal=self,
|
||||
)
|
||||
return await ctx._pld_rx.recv_pld(
|
||||
ipc=ctx,
|
||||
expect_msg=Return,
|
||||
)
|
||||
|
||||
# TODO: factor this out into a `.highlevel` API-wrapper that uses
|
||||
|
|
|
|||
|
|
@ -879,8 +879,10 @@ async def _invoke(
|
|||
# don't pop the local context until we know the
|
||||
# associated child isn't in debug any more
|
||||
await debug.maybe_wait_for_debugger()
|
||||
dropped_ctx: Context|None = actor._drop_context(ctx)
|
||||
assert dropped_ctx is ctx
|
||||
ctx: Context = actor._contexts.pop((
|
||||
chan.aid.uid,
|
||||
cid,
|
||||
))
|
||||
|
||||
logmeth: Callable = log.runtime
|
||||
merr: Exception|None = ctx.maybe_error
|
||||
|
|
|
|||
|
|
@ -600,21 +600,14 @@ class Actor:
|
|||
# - cancel_rpc_tasks(),
|
||||
# - _cancel_task(),
|
||||
#
|
||||
def _get_rpc_func(
|
||||
self,
|
||||
ns: str,
|
||||
funcname: str,
|
||||
):
|
||||
def _get_rpc_func(self, ns, funcname):
|
||||
'''
|
||||
Try to lookup and return a target RPC func from the
|
||||
post-fork enabled module set.
|
||||
|
||||
'''
|
||||
try:
|
||||
return getattr(
|
||||
self._mods[ns],
|
||||
funcname,
|
||||
)
|
||||
return getattr(self._mods[ns], funcname)
|
||||
except KeyError as err:
|
||||
mne = ModuleNotExposed(*err.args)
|
||||
|
||||
|
|
@ -760,23 +753,6 @@ class Actor:
|
|||
|
||||
return ctx
|
||||
|
||||
def _drop_context(
|
||||
self,
|
||||
ctx: Context,
|
||||
) -> Context|None:
|
||||
'''
|
||||
Remove `ctx` from this actor's IPC context registry.
|
||||
|
||||
Teardown paths can converge after normal return, cancellation
|
||||
or startup failure, so registry removal is idempotent.
|
||||
|
||||
'''
|
||||
peer_uid: tuple[str, str] = ctx.chan.aid.uid
|
||||
return self._contexts.pop(
|
||||
(peer_uid, ctx.cid),
|
||||
None,
|
||||
)
|
||||
|
||||
async def start_remote_task(
|
||||
self,
|
||||
chan: Channel,
|
||||
|
|
@ -791,7 +767,6 @@ class Actor:
|
|||
allow_overruns: bool = False,
|
||||
load_nsf: bool = False,
|
||||
ack_timeout: float = float('inf'),
|
||||
cancel_on_startup: bool = True,
|
||||
|
||||
) -> Context:
|
||||
'''
|
||||
|
|
@ -843,17 +818,13 @@ class Actor:
|
|||
|
||||
f'{pretty_struct.pformat(msg)}'
|
||||
)
|
||||
start_published: bool = False
|
||||
try:
|
||||
await chan.send(msg)
|
||||
start_published = True
|
||||
|
||||
# NOTE wait on first `StartAck` response msg and validate;
|
||||
# this should be immediate and does not (yet) wait for the
|
||||
# remote child task to sync via `Context.started()`.
|
||||
with trio.fail_after(ack_timeout):
|
||||
first_msg: msgtypes.StartAck = await ctx._rx_chan.receive()
|
||||
|
||||
try:
|
||||
functype: str = first_msg.functype
|
||||
except AttributeError:
|
||||
|
|
@ -865,45 +836,9 @@ class Actor:
|
|||
'context',
|
||||
):
|
||||
raise ValueError(
|
||||
f'Invalid `StartAck.functype: str = '
|
||||
f'{first_msg!r}` ??'
|
||||
f'Invalid `StartAck.functype: str = {first_msg!r}` ??'
|
||||
)
|
||||
|
||||
except BaseException as startup_err:
|
||||
with trio.CancelScope(shield=True):
|
||||
# `MsgpackTransport.send()` closes its stream when
|
||||
# cancellation interrupts the length-prefixed write
|
||||
# because an unknown prefix may already be sent. A
|
||||
# connected channel means cancellation happened before
|
||||
# that write or after it completed, so `_cancel_task`
|
||||
# is protocol-safe (and a no-op if `Start` was unsent).
|
||||
if (
|
||||
cancel_on_startup
|
||||
and
|
||||
(
|
||||
start_published
|
||||
or (
|
||||
isinstance(startup_err, trio.Cancelled)
|
||||
and
|
||||
chan.connected()
|
||||
)
|
||||
)
|
||||
):
|
||||
try:
|
||||
await ctx.cancel()
|
||||
except BaseException as cancel_err:
|
||||
log.warning(
|
||||
'Failed to cancel RPC task during '
|
||||
'startup?\n'
|
||||
f'{cancel_err!r}\n'
|
||||
)
|
||||
|
||||
self._drop_context(ctx)
|
||||
if not ctx._rx_chan._closed:
|
||||
await ctx._rx_chan.aclose()
|
||||
|
||||
raise
|
||||
|
||||
ctx._remote_func_type = functype
|
||||
return ctx
|
||||
|
||||
|
|
|
|||
|
|
@ -520,7 +520,7 @@ async def _open_and_supervise_one_cancels_all_nursery(
|
|||
'Waiting on subactors to complete:\n'
|
||||
f'>}} {len(an._children)}\n'
|
||||
)
|
||||
an._request_reap_all()
|
||||
an._join_procs.set()
|
||||
|
||||
# Single one-cancels-all handler for the (now single)
|
||||
# daemon nursery. Pre-#477 a 2ndary `._ria_nursery`
|
||||
|
|
@ -549,7 +549,7 @@ async def _open_and_supervise_one_cancels_all_nursery(
|
|||
# if the caller's scope errored then we activate our
|
||||
# one-cancels-all supervisor strategy (don't
|
||||
# worry more are coming).
|
||||
an._request_reap_all()
|
||||
an._join_procs.set()
|
||||
|
||||
# XXX NOTE XXX: hypothetically an error could
|
||||
# be raised and then a cancel signal shows up
|
||||
|
|
|
|||
|
|
@ -169,16 +169,13 @@ async def mp_proc(
|
|||
# any process we may have started.
|
||||
|
||||
portal = Portal(chan)
|
||||
reap_request, _ = actor_nursery._register_child_reap(
|
||||
subactor.aid.uid,
|
||||
)
|
||||
actor_nursery._children[subactor.aid.uid] = (subactor, proc, portal)
|
||||
|
||||
# unblock parent task
|
||||
task_status.started(portal)
|
||||
|
||||
# wait for this child or its `ActorNursery` to signal that
|
||||
# the subprocess can be joined.
|
||||
# wait for ``ActorNursery`` block to signal that
|
||||
# subprocesses can be waited upon.
|
||||
# This is required to ensure synchronization
|
||||
# with user code that may want to manually await results
|
||||
# from nursery spawned sub-actors. We don't want the
|
||||
|
|
@ -187,7 +184,7 @@ async def mp_proc(
|
|||
# nursery block closes do we allow subactor results to be
|
||||
# awaited and reported upwards to the supervisor.
|
||||
with trio.CancelScope(shield=True):
|
||||
await reap_request.wait()
|
||||
await actor_nursery._join_procs.wait()
|
||||
|
||||
# This is a "soft" (cancellable) join/reap which
|
||||
# will remote cancel the actor on a ``trio.Cancelled``
|
||||
|
|
|
|||
|
|
@ -354,7 +354,6 @@ async def new_proc(
|
|||
# mark the new actor with the global spawn method
|
||||
subactor._spawn_method = _spawn_method
|
||||
|
||||
try:
|
||||
await target(
|
||||
name,
|
||||
actor_nursery,
|
||||
|
|
@ -367,8 +366,6 @@ async def new_proc(
|
|||
task_status=task_status,
|
||||
proc_kwargs=proc_kwargs
|
||||
)
|
||||
finally:
|
||||
actor_nursery._mark_child_reaped(subactor.aid.uid)
|
||||
|
||||
|
||||
# NOTE: bottom-of-module to avoid a circular import since the
|
||||
|
|
|
|||
|
|
@ -160,9 +160,6 @@ async def trio_proc(
|
|||
assert proc
|
||||
|
||||
portal = Portal(chan)
|
||||
reap_request, _ = actor_nursery._register_child_reap(
|
||||
subactor.aid.uid,
|
||||
)
|
||||
actor_nursery._children[subactor.aid.uid] = (
|
||||
subactor,
|
||||
proc,
|
||||
|
|
@ -193,10 +190,9 @@ async def trio_proc(
|
|||
# resume caller at next checkpoint now that child is up
|
||||
task_status.started(portal)
|
||||
|
||||
# wait for this child or its `ActorNursery` to request
|
||||
# process joining.
|
||||
# wait for ActorNursery.wait() to be called
|
||||
with trio.CancelScope(shield=True):
|
||||
await reap_request.wait()
|
||||
await actor_nursery._join_procs.wait()
|
||||
|
||||
# This is a "soft" (cancellable) join/reap which
|
||||
# will remote cancel the actor on a ``trio.Cancelled``
|
||||
|
|
|
|||
|
|
@ -22,18 +22,12 @@ Adopts the "run it over there" parlance from analogous
|
|||
`anyio.to_process` but for SC-supervised actors: spawn (or
|
||||
reuse) a subactor, schedule a single remote task, wait on
|
||||
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
|
||||
https://github.com/goodboy/tractor/issues/477
|
||||
|
||||
'''
|
||||
from . import _api as _api
|
||||
from ._api import (
|
||||
run as run,
|
||||
)
|
||||
|
||||
|
||||
MODULE: str = _api.__name__
|
||||
|
|
|
|||
|
|
@ -23,8 +23,8 @@ the lower level daemon-actor spawn + portal APIs,
|
|||
|
||||
- `ActorNursery.start_actor()` for (daemon-style) subactor
|
||||
spawning,
|
||||
- `Portal.open_context()` for scheduling the lone remote
|
||||
task with linked cancellation and waiting on its result,
|
||||
- `Portal.run()` for scheduling the lone remote task and
|
||||
waiting on its result,
|
||||
- `Portal.cancel_actor()` for reaping the subactor once
|
||||
that result (or error) arrives,
|
||||
|
||||
|
|
@ -36,24 +36,13 @@ spawn-machinery nurseries as with the (now removed) legacy
|
|||
|
||||
'''
|
||||
from __future__ import annotations
|
||||
import functools
|
||||
import inspect
|
||||
from typing import (
|
||||
Any,
|
||||
Awaitable,
|
||||
Callable,
|
||||
TYPE_CHECKING,
|
||||
TypeVar,
|
||||
TypeVarTuple,
|
||||
Unpack,
|
||||
)
|
||||
|
||||
from .._context import (
|
||||
Context,
|
||||
context,
|
||||
)
|
||||
from ..msg.ptr import NamespacePath
|
||||
from ..runtime._state import current_actor
|
||||
from ..runtime._supervise import (
|
||||
ActorNursery,
|
||||
open_nursery,
|
||||
|
|
@ -64,10 +53,6 @@ if TYPE_CHECKING:
|
|||
from ..runtime._portal import Portal
|
||||
|
||||
|
||||
ArgsT = TypeVarTuple('ArgsT')
|
||||
RetT = TypeVar('RetT')
|
||||
|
||||
|
||||
def _validate_one_shot_fn(
|
||||
fn: Callable,
|
||||
) -> None:
|
||||
|
|
@ -75,7 +60,7 @@ def _validate_one_shot_fn(
|
|||
Ensure `fn` is a non-streaming async function, raise
|
||||
a `TypeError` otherwise.
|
||||
|
||||
The same constraint enforced by `Portal.open_context()` but
|
||||
The same constraint enforced by `Portal.run()` but
|
||||
checked up-front, BEFORE any subactor is spawned.
|
||||
|
||||
'''
|
||||
|
|
@ -94,127 +79,18 @@ def _validate_one_shot_fn(
|
|||
)
|
||||
|
||||
|
||||
def _normalize_call(
|
||||
fn: Callable,
|
||||
args: tuple[Any, ...],
|
||||
) -> tuple[
|
||||
Callable,
|
||||
tuple[Any, ...],
|
||||
dict[str, Any],
|
||||
]:
|
||||
'''
|
||||
Normalize Trio-style positional and partial-bound arguments.
|
||||
|
||||
Actor calls must send a namespace-addressable base function and
|
||||
serializable inputs to another process, so decompose partials and
|
||||
validate their complete call signature before runtime startup.
|
||||
|
||||
'''
|
||||
kwargs: dict[str, Any] = {}
|
||||
while isinstance(fn, functools.partial):
|
||||
partial_args: tuple[Any, ...] = fn.args
|
||||
|
||||
# `functools.Placeholder` was added in Python 3.14. Drop
|
||||
# this `getattr()` guard once 3.14 is the minimum version.
|
||||
placeholder = getattr(
|
||||
functools,
|
||||
'Placeholder',
|
||||
None,
|
||||
)
|
||||
if (
|
||||
placeholder is not None
|
||||
and
|
||||
any(
|
||||
arg is placeholder
|
||||
for arg in partial_args
|
||||
)
|
||||
):
|
||||
call_args = iter(args)
|
||||
merged_args: list[Any] = []
|
||||
for arg in partial_args:
|
||||
if arg is placeholder:
|
||||
try:
|
||||
arg = next(call_args)
|
||||
except StopIteration:
|
||||
raise TypeError(
|
||||
'Not enough positional arguments to '
|
||||
'fill `functools.Placeholder`s'
|
||||
) from None
|
||||
|
||||
merged_args.append(arg)
|
||||
|
||||
merged_args.extend(call_args)
|
||||
args = tuple(merged_args)
|
||||
else:
|
||||
args = partial_args + args
|
||||
|
||||
partial_kwargs = dict(fn.keywords or {})
|
||||
partial_kwargs.update(kwargs)
|
||||
kwargs = partial_kwargs
|
||||
fn = fn.func
|
||||
|
||||
_validate_one_shot_fn(fn)
|
||||
inspect.signature(fn).bind(*args, **kwargs)
|
||||
return fn, args, kwargs
|
||||
|
||||
|
||||
@context
|
||||
async def _invoke_one_shot(
|
||||
ctx: Context,
|
||||
namespace: str,
|
||||
funcname: str,
|
||||
args: list[Any],
|
||||
kwargs: dict[str, Any],
|
||||
) -> Any:
|
||||
'''
|
||||
Invoke an ordinary async function inside a linked IPC context.
|
||||
|
||||
'''
|
||||
# Do not use `NamespacePath.load_ref()` here: target resolution
|
||||
# must remain behind the actor's RPC module allowlist.
|
||||
fn: Callable = current_actor()._get_rpc_func(
|
||||
namespace,
|
||||
funcname,
|
||||
)
|
||||
_validate_one_shot_fn(fn)
|
||||
await ctx.started()
|
||||
return await fn(*args, **kwargs)
|
||||
|
||||
|
||||
async def _invoke_from_portal(
|
||||
portal: Portal,
|
||||
fn: Callable,
|
||||
args: tuple[Any, ...],
|
||||
kwargs: dict[str, Any],
|
||||
) -> Any:
|
||||
'''
|
||||
Run `fn` through the context-linked one-shot endpoint.
|
||||
|
||||
'''
|
||||
namespace, funcname = NamespacePath.from_ref(fn).to_tuple()
|
||||
async with portal.open_context(
|
||||
_invoke_one_shot,
|
||||
namespace=namespace,
|
||||
funcname=funcname,
|
||||
args=list(args),
|
||||
kwargs=kwargs,
|
||||
) as (ctx, _):
|
||||
return await ctx.wait_for_result()
|
||||
|
||||
|
||||
async def _invoke_in_subactor(
|
||||
an: ActorNursery,
|
||||
fn: Callable,
|
||||
args: tuple[Any, ...],
|
||||
kwargs: dict[str, Any],
|
||||
name: str,
|
||||
spawn_kwargs: dict[str, Any],
|
||||
fn_kwargs: dict[str, Any],
|
||||
) -> Any:
|
||||
'''
|
||||
Spawn a (daemon) subactor via `an.start_actor()`,
|
||||
schedule `fn` as its context-linked lone remote task and,
|
||||
ALWAYS, reap the subactor once that task's result (or error)
|
||||
has been delivered.
|
||||
schedule `fn` as its lone remote task via
|
||||
`Portal.run()` and, ALWAYS, reap the subactor once
|
||||
that task's result (or error) has been delivered.
|
||||
|
||||
'''
|
||||
portal: Portal = await an.start_actor(
|
||||
|
|
@ -222,22 +98,23 @@ async def _invoke_in_subactor(
|
|||
**spawn_kwargs,
|
||||
)
|
||||
try:
|
||||
return await _invoke_from_portal(
|
||||
portal,
|
||||
return await portal.run(
|
||||
fn,
|
||||
args,
|
||||
kwargs,
|
||||
**fn_kwargs,
|
||||
)
|
||||
finally:
|
||||
# Cancel and join this child before returning. The nursery
|
||||
# helper shields teardown, escalates a missed cancel ack and
|
||||
# waits for the child monitor to remove its process record.
|
||||
await an._cancel_and_reap_child(portal)
|
||||
# one-shot semantics: the subactor's lifetime is
|
||||
# bound to its lone task's completion; the
|
||||
# cancel-req's bounded wait is shielded
|
||||
# internally (see `Portal.cancel_actor()`) so
|
||||
# this reap also runs when the caller's scope
|
||||
# was itself cancelled.
|
||||
await portal.cancel_actor()
|
||||
|
||||
|
||||
async def run(
|
||||
fn: Callable[[Unpack[ArgsT]], Awaitable[RetT]],
|
||||
*args: Unpack[ArgsT],
|
||||
fn: Callable,
|
||||
*,
|
||||
|
||||
# actor "placement": reuse an already-running peer
|
||||
# via its `portal`, spawn a fresh subactor from
|
||||
|
|
@ -265,24 +142,18 @@ async def run(
|
|||
# when NO `an`/`portal` is provided.
|
||||
runtime_kwargs: dict[str, Any]|None = None,
|
||||
|
||||
) -> RetT:
|
||||
**fn_kwargs, # explicit (keyword) args to `fn`
|
||||
|
||||
) -> Any:
|
||||
'''
|
||||
Run the async `fn(*args)` as the lone task in a (new)
|
||||
subactor, block waiting on its result and return it; the
|
||||
distributed-parallelism equivalent of
|
||||
Run the async `fn` as the lone task in a (new)
|
||||
subactor, block waiting on its result and return it;
|
||||
the distributed-parallelism equivalent of
|
||||
`trio.to_thread.run_sync()`.
|
||||
|
||||
As with Trio's API, target arguments are positional. Use
|
||||
`functools.partial()` to bind target keyword arguments; all
|
||||
keyword arguments accepted here configure actor placement or
|
||||
spawning. A caller-supplied `portal` must address an actor started
|
||||
with both `tractor.to_actor.MODULE` and the target function's
|
||||
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
|
||||
actor-nursery teardown) this is a plain "call and
|
||||
Unlike the removed legacy `.run_in_actor()` (which
|
||||
returned a `Portal` whose result was 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
|
||||
`trio` way by scheduling multiple `run()` calls in
|
||||
|
|
@ -292,10 +163,10 @@ async def run(
|
|||
|
||||
'''
|
||||
__runtimeframe__: int = 1 # noqa
|
||||
fn, args, kwargs = _normalize_call(fn, args)
|
||||
_validate_one_shot_fn(fn)
|
||||
|
||||
if (
|
||||
runtime_kwargs is not None
|
||||
runtime_kwargs
|
||||
and
|
||||
(
|
||||
an is not None
|
||||
|
|
@ -315,22 +186,15 @@ async def run(
|
|||
'Pass at most ONE of `portal` or `an`, '
|
||||
'not both!'
|
||||
)
|
||||
return await _invoke_from_portal(
|
||||
portal,
|
||||
return await portal.run(
|
||||
fn,
|
||||
args,
|
||||
kwargs,
|
||||
**fn_kwargs,
|
||||
)
|
||||
|
||||
name: str = name or fn.__name__
|
||||
spawn_kwargs: dict[str, Any] = dict(
|
||||
enable_modules=(
|
||||
[
|
||||
# The public `to_actor.MODULE` alias is only for
|
||||
# callers configuring an existing actor.
|
||||
__name__,
|
||||
fn.__module__,
|
||||
]
|
||||
[fn.__module__]
|
||||
+
|
||||
(enable_modules or [])
|
||||
),
|
||||
|
|
@ -345,21 +209,18 @@ async def run(
|
|||
return await _invoke_in_subactor(
|
||||
an,
|
||||
fn,
|
||||
args,
|
||||
kwargs,
|
||||
name,
|
||||
spawn_kwargs,
|
||||
fn_kwargs,
|
||||
)
|
||||
|
||||
an: ActorNursery
|
||||
async with open_nursery(
|
||||
**(runtime_kwargs or {}),
|
||||
) as an:
|
||||
return await _invoke_in_subactor(
|
||||
an,
|
||||
fn,
|
||||
args,
|
||||
kwargs,
|
||||
name,
|
||||
spawn_kwargs,
|
||||
fn_kwargs,
|
||||
)
|
||||
|
|
|
|||
Loading…
Reference in New Issue