Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
22 changes: 22 additions & 0 deletions .github/workflows/tests.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -74,6 +74,28 @@ jobs:
./tests.sh medium --ignore tests/old --ignore tests/dev -x -vv --timeout=180 --timeout-method=thread
fi

notebook-dispatch:
name: Notebook dispatch (Ubuntu, Python 3.12)
runs-on: ubuntu-latest
timeout-minutes: 10

steps:
- uses: actions/checkout@v5

- uses: actions/setup-python@v6
with:
python-version: "3.12"

- name: Install notebook test dependencies
run: |
python -m pip install --upgrade pip
python -m pip install -r test_requirements.txt "ipykernel==7.2.0" "jupyter-client==8.8.0"
python -m pip install .

- name: Run real-kernel Dispatch regression
run: |
./tests.sh tests/dispatch/test_notebook.py --no-cov -x --timeout=180 --timeout-method=thread --basetemp="$RUNNER_TEMP/dryml-notebook-tests"

local-filesystem-publication:
name: Good-enough (${{ matrix.os }}, Python ${{ matrix.python-version }})
runs-on: ${{ matrix.os }}
Expand Down
33 changes: 24 additions & 9 deletions docs/dispatch.md
Original file line number Diff line number Diff line change
Expand Up @@ -83,17 +83,18 @@ finalization retains its existing publication behavior.

Reports contain only bounded diagnostic categories, redacted backend identifiers,
coverage and requirement results. They never retain call data, handles,
credentials, source, selectors, Stores, or reservations. Valid incomplete static
coverage is reported; `run` and `submit` emit `DispatchCoverageWarning`, while
`explain` does not warn.
credentials, source, selectors, Stores, or reservations. Incomplete static
coverage remains in `explain` diagnostics even when a normal unresolved call
does not warn. `explain` never emits `DispatchCoverageWarning`.

## Probing And Static Coverage

`ProbeOptions` is an immutable, inert policy with defaults
`placement="auto"`, `execution_timeout=30.0`, `max_targets=256`, and
`max_depth=32`. `placement="auto"` uses an explicit probe backend when supplied;
otherwise it probes inline only with compatible current-process evidence and uses
an owned local subprocess when isolation is required. `placement="execute"`
`max_depth=32`, and `coverage_policy="default"`. `placement="auto"` uses an
explicit probe backend when supplied; otherwise it probes inline only with
compatible current-process evidence and uses an owned local subprocess when
isolation is required. `placement="execute"`
uses the supplied backend or that local subprocess default. `placement="in_process"`
requires compatible current-process evidence and rejects a contradictory backend.
Probe placement is independent of the selected workload route. It never falls
Expand All @@ -117,8 +118,14 @@ at 4 MiB. They retain at most 64 diagnostic entries with 512 characters per
field. Capture separately limits individual source reads to 1 MiB, aggregate
source reads to 8 MiB, candidate targets to 4,096, binding/call facts to 16,384,
and raw annotation occurrences to 4,096. These ceilings do not make static
analysis a full call-graph proof. A valid incomplete result warns and can proceed
when known requirements pass; malformed projections/results, conflicts, crashes,
analysis a full call-graph proof. By default, accepted incomplete coverage only
warns when a traversal limit or other diagnostic accompanies `static.unresolved`;
an unresolved-only call (common for notebook-defined functions and callbacks)
proceeds quietly with `coverage="incomplete"` and its diagnostics retained in
`explain`. `ProbeOptions(coverage_policy="warn")` restores warnings for every
accepted incomplete probe. `ProbeOptions(coverage_policy="strict")` instead
returns an ineligible report and rejects `run`/`submit` before workload acceptance
for any incomplete probe. Malformed projections/results, conflicts, crashes,
timeouts, or cleanup failures stop submission rather than becoming empty results.

An `EnvironmentSpec` pin selects an existing interpreter, venv executable, or
Expand All @@ -140,7 +147,15 @@ result recovery, and cleanup retain the core Execute contracts.
result, and reconciles its owned cleanup. If both the workload result and
cleanup fail, the workload failure remains primary and the cleanup failure is
its cause. Dispatch does not retry work, adapt results, or choose another
backend after an unavailable or rejected selected backend.
backend after an unavailable or rejected selected backend. A `KeyboardInterrupt`
while the future is still pending first attempts confirmed pre-GO cancellation,
then requests best-effort running cancellation if GO has won. It propagates the
interrupt without prematurely cleaning up running work; core's one-off
completion owns eventual cleanup. Blocking `run` occupies a notebook shell until
it completes or is interrupted. In an async notebook cell, use `submit` and
`await future` to avoid blocking the shell task; the caller remains responsible
for observing and cleaning up the returned future. Captured live asyncio tasks,
futures, and event loops cannot be transported into workers.

The selected worker environment and world requirements are the combined
configured and discovered requirements. A resolved worker Python selector is
Expand Down
11 changes: 8 additions & 3 deletions docs/release_notes.md
Original file line number Diff line number Diff line change
Expand Up @@ -39,9 +39,14 @@ and async-generator roots reject before probing; lazy values returned from a
synchronous root are data and are not driven by Dispatch.

Discovery uses shared generic static dependency, environment, and world kernels
for inline and isolated probes. Incomplete but valid static coverage issues a
visible warning for `run`/`submit`; conflicts, malformed results, timeout/crash,
cleanup failure, target drift, and hard admission failures stop work. Reports and
for inline and isolated probes. Unresolved-only static coverage remains visible
in `explain` without warning on routine accepted calls; traversal limits still
warn, and explicit probe policy can warn for or reject all incomplete coverage.
Notebook-owned asyncio tasks, futures, and loops reject at the transport boundary
rather than crossing worker serialization, and interruption of pending blocking
dispatch no longer attempts premature cleanup. Conflicts, malformed results,
timeout/crash, cleanup failure, target drift, and hard admission failures stop
work. Reports and
probe envelopes are bounded, redacted, ephemeral observations rather than Store
records or admission tickets. Existing `EnvironmentSpec` selection now flows
through generic/core Execute and Dispatch as an exact point-in-time pin; an
Expand Down
7 changes: 7 additions & 0 deletions docs/testing.md
Original file line number Diff line number Diff line change
Expand Up @@ -364,6 +364,13 @@ already included in their routine jobs, avoiding duplicate runners. Package
tests run only for manually requested exhaustive/coverage verification, through
`medium` on Ubuntu/Windows and a package-only step on macOS.

A separate Ubuntu Python 3.12 notebook job installs ipykernel 7.2.0 and
jupyter-client 8.8.0 and runs the real-kernel Dispatch regression on every push
and pull request. The general lightweight matrix does not install those optional
dependencies, so its skipped notebook case is not the notebook verification gate.
The test uses a job-owned kernel, temporary Store, and worker spool; it checks
repeated publication/query, shell liveness, task ownership, and interrupt recovery.

The heavy matrix runs on Ubuntu for Python 3.10 through 3.13 only when a user
manually dispatches `exhaustive` or `coverage`. It
installs and preflights TensorFlow, Torch, JAX/JAXlib, and pinned
Expand Down
5 changes: 5 additions & 0 deletions src/dryml/core/execute_codec.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@

import builtins
import ast
import asyncio
import dis
import inspect
import pathlib
Expand Down Expand Up @@ -95,6 +96,8 @@ def persistent_id(self, value: object) -> object | None:
from dryml.managed.config import ManagedConfig
from .template import TemplateBundle

if isinstance(value, (asyncio.Future, asyncio.AbstractEventLoop)):
raise _DillLeafError("live asyncio resource")
if isinstance(value, (Repo, Store)):
raise _DillLeafError("live core resource")
if isinstance(
Expand Down Expand Up @@ -427,6 +430,8 @@ def _node(self, value: Any, path: str, depth: int) -> dict[str, Any]:
"receiver": self.value(value._instance, f"{path}.receiver", depth + 1),
"composite": self.value(value._composite, f"{path}.composite", depth + 1),
}
if isinstance(value, (asyncio.Future, asyncio.AbstractEventLoop)):
_fail("live asyncio resource", path)
if _is_resource(value):
_fail("live core resource", path)
imported_capture = self.imported_captures.get(id(value))
Expand Down
15 changes: 14 additions & 1 deletion src/dryml/dispatch/_preflight.py
Original file line number Diff line number Diff line change
Expand Up @@ -116,7 +116,16 @@ def _report(
if probe.placement == "in_process"
else "selected Execute probe configuration"
)
warnings = (_COVERAGE_WARNING,) if not probe.coverage.complete else ()
coverage_policy = options.probe.coverage_policy
warn = (
eligible
and not probe.coverage.complete
and (coverage_policy == "warn" or (
coverage_policy == "default"
and set(probe.coverage.diagnostics) != {"static.unresolved"}
))
)
warnings = (_COVERAGE_WARNING,) if warn else ()
return DispatchReport(
workload_placement,
workload_backend,
Expand Down Expand Up @@ -313,6 +322,10 @@ def preflight(
if eligible
else "configured and discovered requirements are incompatible"
)
if eligible and options.probe.coverage_policy == "strict" and not result.coverage.complete:
eligible = False
reason = "static requirement coverage is incomplete under strict policy"
outcome_diagnostics = (*outcome_diagnostics, _COVERAGE_WARNING)
report = _report(
options,
probe=result,
Expand Down
17 changes: 14 additions & 3 deletions src/dryml/dispatch/api.py
Original file line number Diff line number Diff line change
Expand Up @@ -44,9 +44,9 @@ def _checked_preflight(


def _warn_coverage(prepared: _Preflight) -> None:
"""Emit the valid-incomplete warning for accepted operations."""
"""Emit only coverage categories selected by the probe policy."""

if prepared.report.coverage == "incomplete":
if "dispatch.coverage_incomplete" in prepared.report.warnings:
warnings.warn(
"Dispatch static requirement coverage is incomplete",
DispatchCoverageWarning,
Expand Down Expand Up @@ -293,6 +293,14 @@ def _run(
try:
result = future.result()
except BaseException as error:
if isinstance(error, KeyboardInterrupt) and not future.done():
try:
if not future.cancel():
future.request_cancel()
except Exception:
pass
# Core's one-off completion callback owns cleanup after terminality.
raise
try:
future.cleanup()
except CleanupError as cleanup_error:
Expand All @@ -319,10 +327,13 @@ def run(fn: Any, /, *args: Any, **kwargs: Any) -> Any:
workload before acceptance.
BaseException: Existing core/backend result and cleanup failures, with
the primary execution failure retained when cleanup also fails.
KeyboardInterrupt: Propagates interruption; pending backend work is
first cancelled before GO or requested to stop if already running.

Side Effects:
Runs bounded preflight and one backend submission for Execute routes.
It neither retries work nor changes the selected route.
It neither retries work nor changes the selected route. A pending
interruption leaves terminal cleanup to the core one-off owner.
"""

return _run(fn, args, kwargs, None)
Expand Down
13 changes: 10 additions & 3 deletions src/dryml/dispatch/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -41,11 +41,12 @@ class InProcess:


class DispatchCoverageWarning(RuntimeWarning):
"""Warn that valid static requirement collection was incomplete.
"""Warn that accepted static requirement collection was incomplete.

The warning carries no workload value, source, backend credentials, or
reservation. It is emitted by ``run`` and ``submit`` only; ``explain``
keeps the same fact in its immutable report.
reservation. ``run`` and ``submit`` emit it for bounded-analysis limits or
when the caller opts into all incomplete-coverage warnings. ``explain``
retains diagnostics without emitting warnings.
"""


Expand Down Expand Up @@ -263,6 +264,9 @@ class ProbeOptions:
checked cooperatively at inline analysis boundaries.
max_targets: Positive maximum static traversal targets, including root.
max_depth: Positive maximum static traversal depth.
coverage_policy: ``"default"`` warns for bounded-analysis limits but
not unresolved-only calls; ``"warn"`` warns for all incomplete
coverage; ``"strict"`` rejects incomplete coverage before execution.

Raises:
TypeError: If a field has an unsupported type.
Expand All @@ -281,12 +285,15 @@ class ProbeOptions:
execution_timeout: float = 30.0
max_targets: int = 256
max_depth: int = 32
coverage_policy: Literal["default", "warn", "strict"] = "default"

def __post_init__(self) -> None:
"""Validate the inert policy without resolving external state."""

if self.placement not in ("auto", "in_process", "execute"):
raise ValueError("probe placement is invalid")
if self.coverage_policy not in ("default", "warn", "strict"):
raise ValueError("probe coverage_policy is invalid")
if self.backend is not None and not isinstance(
self.backend, (BackendConfig, str)
):
Expand Down
3 changes: 2 additions & 1 deletion src/dryml/execute/_spooling.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@

from __future__ import annotations

import asyncio
import concurrent.futures
import hashlib
import hmac
Expand Down Expand Up @@ -233,7 +234,7 @@ def _reject_live_resource(value: object) -> None:
value_type = type(value)
if any(cls.__module__ == "dryml.core" or cls.__module__.startswith("dryml.core.") for cls in value_type.__mro__):
raise _UnsupportedTransportResource("live resource: core semantic values are unsupported by Execute transport")
if isinstance(value, (io.IOBase, _LOCK_TYPES, GeneratorType, CoroutineType, socket.socket, threading.Thread, concurrent.futures.Executor, concurrent.futures.Future)):
if isinstance(value, (io.IOBase, _LOCK_TYPES, GeneratorType, CoroutineType, socket.socket, threading.Thread, concurrent.futures.Executor, concurrent.futures.Future, asyncio.Future, asyncio.AbstractEventLoop)):
raise _UnsupportedTransportResource("live resource: streams, locks, sockets, threads, and futures are unsupported by Execute transport")
if isinstance(value, _ConnectionBase):
raise _UnsupportedTransportResource("live resource: Connection is unsupported by Execute transport")
Expand Down
26 changes: 26 additions & 0 deletions tests/core/test_execute_callables.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,18 +3,21 @@
from __future__ import annotations

from dataclasses import dataclass
import asyncio
import functools
import os
from pathlib import Path

from dryml.core import ObjectRef, Repo, Serializable, function
from dryml.core.execute import ExecutionContext, SharedDirStoreStrategy, invoke_prepared_call, worker_context
from dryml.core.execute_codec import CoreCallCodecError, encode_invocation
from dryml.core.signatures import Ref
from dryml.core.store.dir import DirStore
from dryml.execute import Executor
from dryml.execute.subprocess import SubProcessConfig
from dryml.managed import managed_operation
from dryml.methods import Method
import pytest


def _invoke(fn, args, repo):
Expand Down Expand Up @@ -258,6 +261,29 @@ def test_fresh_subprocess_uses_dill_for_ordinary_leaves(tmp_path):
assert result == FrozenOrdinaryValue(7)


def test_core_rejects_captured_kernel_task_before_dill_traversal():
"""Never serialize an active task owned by the caller's event loop."""

async def exercise():
task = asyncio.create_task(asyncio.Event().wait())

def captured():
return task

try:
with pytest.raises(CoreCallCodecError, match="live asyncio resource"):
encode_invocation(captured, (), {}, repo=Repo())
with pytest.raises(CoreCallCodecError, match="live asyncio resource"):
encode_invocation(lambda value: value, (FrozenOrdinaryValue(task),), {}, repo=Repo())
assert not task.done()
finally:
task.cancel()
with pytest.raises(asyncio.CancelledError):
await task

asyncio.run(exercise())


def test_copied_function_wrapper_metadata_rebuilds_its_established_owner(
tmp_path):
"""
Expand Down
Loading
Loading