diff --git a/.github/workflows/tests.yaml b/.github/workflows/tests.yaml index d9005b55..3262e467 100644 --- a/.github/workflows/tests.yaml +++ b/.github/workflows/tests.yaml @@ -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 }} diff --git a/docs/dispatch.md b/docs/dispatch.md index a235feb8..8494f670 100644 --- a/docs/dispatch.md +++ b/docs/dispatch.md @@ -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 @@ -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 @@ -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 diff --git a/docs/release_notes.md b/docs/release_notes.md index 7fcf8d12..81ffa68d 100644 --- a/docs/release_notes.md +++ b/docs/release_notes.md @@ -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 diff --git a/docs/testing.md b/docs/testing.md index cbef2fea..00345698 100644 --- a/docs/testing.md +++ b/docs/testing.md @@ -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 diff --git a/src/dryml/core/execute_codec.py b/src/dryml/core/execute_codec.py index 395383df..b1b08f0b 100644 --- a/src/dryml/core/execute_codec.py +++ b/src/dryml/core/execute_codec.py @@ -9,6 +9,7 @@ import builtins import ast +import asyncio import dis import inspect import pathlib @@ -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( @@ -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)) diff --git a/src/dryml/dispatch/_preflight.py b/src/dryml/dispatch/_preflight.py index 19ad9fdd..a7089183 100644 --- a/src/dryml/dispatch/_preflight.py +++ b/src/dryml/dispatch/_preflight.py @@ -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, @@ -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, diff --git a/src/dryml/dispatch/api.py b/src/dryml/dispatch/api.py index 69e4a533..c5e78d64 100644 --- a/src/dryml/dispatch/api.py +++ b/src/dryml/dispatch/api.py @@ -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, @@ -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: @@ -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) diff --git a/src/dryml/dispatch/models.py b/src/dryml/dispatch/models.py index 1a65c0dd..d5c7d39c 100644 --- a/src/dryml/dispatch/models.py +++ b/src/dryml/dispatch/models.py @@ -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. """ @@ -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. @@ -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) ): diff --git a/src/dryml/execute/_spooling.py b/src/dryml/execute/_spooling.py index 88afe772..1ef8b939 100644 --- a/src/dryml/execute/_spooling.py +++ b/src/dryml/execute/_spooling.py @@ -6,6 +6,7 @@ from __future__ import annotations +import asyncio import concurrent.futures import hashlib import hmac @@ -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") diff --git a/tests/core/test_execute_callables.py b/tests/core/test_execute_callables.py index 57e9d701..5c1867a4 100644 --- a/tests/core/test_execute_callables.py +++ b/tests/core/test_execute_callables.py @@ -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): @@ -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): """ diff --git a/tests/dispatch/test_backend_execution.py b/tests/dispatch/test_backend_execution.py index a7799d80..40422543 100644 --- a/tests/dispatch/test_backend_execution.py +++ b/tests/dispatch/test_backend_execution.py @@ -5,6 +5,7 @@ from dataclasses import dataclass, field from functools import wraps from pathlib import Path +import warnings import pytest @@ -435,3 +436,72 @@ def cleanup(self): assert raised.value is primary assert raised.value.__cause__ is cleanup + + +@pytest.mark.parametrize("prestart", (True, False)) +def test_interrupt_during_pending_dispatch_does_not_cleanup_live_future( + monkeypatch: pytest.MonkeyPatch, + prestart: bool, +) -> None: + """Interrupt preserves the shell failure and leaves terminal cleanup to core.""" + + dispatch.set_execute_backend_default(SubProcessConfig()) + interrupted = KeyboardInterrupt() + calls = [] + + class _Future: + def result(self): + raise interrupted + + def done(self): + return False + + def cancel(self): + calls.append("cancel-prestart") + return prestart + + def request_cancel(self): + calls.append("cancel-running") + return True + + def cleanup(self): + calls.append("cleanup") + raise RuntimeError("cannot clean pending future") + + monkeypatch.setattr(dispatch.api, "_submit_backend", lambda *_args: _Future()) + + with pytest.raises(KeyboardInterrupt) as raised: + dispatch.run(lambda: None) + + assert raised.value is interrupted + assert calls == (["cancel-prestart"] if prestart else ["cancel-prestart", "cancel-running"]) + + +def test_backend_submit_warning_policy_preserves_explain_diagnostics( + monkeypatch: pytest.MonkeyPatch, +) -> None: + """Warn only when explicitly requested for an unresolved backend call.""" + + def notebook_style(callback): + return callback() + + dispatch.set_execute_backend_default(SubProcessConfig()) + future = object() + monkeypatch.setattr(dispatch.api, "_submit_backend", lambda *_args: future) + report = dispatch.explain(notebook_style) + assert report.eligible and report.diagnostics == ("static.unresolved",) + assert report.warnings == () + with warnings.catch_warnings(record=True) as captured: + warnings.simplefilter("always") + assert dispatch.submit(notebook_style, lambda: None) is future + assert not any(item.category is dispatch.DispatchCoverageWarning for item in captured) + + warn = dispatch.with_options(probe=dispatch.ProbeOptions(coverage_policy="warn")) + with warnings.catch_warnings(record=True) as captured: + warnings.simplefilter("always") + assert warn.submit(notebook_style, lambda: None) is future + assert sum(item.category is dispatch.DispatchCoverageWarning for item in captured) == 1 + + strict = dispatch.with_options(probe=dispatch.ProbeOptions(coverage_policy="strict")) + with pytest.raises(dispatch.DispatchError): + strict.submit(notebook_style, lambda: None) diff --git a/tests/dispatch/test_in_process.py b/tests/dispatch/test_in_process.py index b474fede..c611ca3c 100644 --- a/tests/dispatch/test_in_process.py +++ b/tests/dispatch/test_in_process.py @@ -390,23 +390,73 @@ def workload() -> str: assert isinstance(result[0], PublicationBusyError) -def test_incomplete_coverage_warns_but_known_conflict_is_ineligible() -> None: - """Preserve advisory incomplete coverage while stopping hard conflicts.""" +def test_unresolved_coverage_is_quiet_by_default_but_explicitly_warns_or_rejects() -> None: + """Keep notebook-style unresolved calls observable without warning by default.""" def incomplete(callback): """Keep a parameter call unresolved while returning valid data.""" return callback() + view = _local_view() + report = view.explain(incomplete) + assert report.eligible + assert report.coverage == "incomplete" + assert report.diagnostics == ("static.unresolved",) + assert report.warnings == () + with warnings.catch_warnings(record=True) as captured: warnings.simplefilter("always") - assert ( - _local_view().run(incomplete, lambda: "complete enough") - == "complete enough" - ) - assert any( - item.category is dispatch.DispatchCoverageWarning for item in captured - ) + assert view.run(incomplete, lambda: "complete enough") == "complete enough" + assert not any(item.category is dispatch.DispatchCoverageWarning for item in captured) + + warn = view.with_options(probe=dispatch.ProbeOptions(coverage_policy="warn")) + assert warn.explain(incomplete).warnings == ("dispatch.coverage_incomplete",) + with warnings.catch_warnings(record=True) as captured: + warnings.simplefilter("always") + assert warn.run(incomplete, lambda: "warned") == "warned" + assert sum(item.category is dispatch.DispatchCoverageWarning for item in captured) == 1 + + strict = view.with_options(probe=dispatch.ProbeOptions(coverage_policy="strict")) + report = strict.explain(incomplete) + assert not report.eligible + assert report.coverage == "incomplete" + assert "static.unresolved" in report.diagnostics + assert "dispatch.coverage_incomplete" in report.diagnostics + with pytest.raises(dispatch.DispatchError) as raised: + strict.run(incomplete, lambda: pytest.fail("strict must reject")) + assert raised.value.report.diagnostics == report.diagnostics + + +def test_bounded_static_traversal_warns_by_default() -> None: + """Exhausting the selected probe bound still alerts accepted callers.""" + + def helper(): + return "ready" + + def workload(): + return helper() + + view = _local_view().with_options(probe=dispatch.ProbeOptions(max_targets=1)) + report = view.explain(workload) + assert report.eligible + assert "static.target_limit" in report.diagnostics + assert report.warnings == ("dispatch.coverage_incomplete",) + with warnings.catch_warnings(record=True) as captured: + warnings.simplefilter("always") + assert view.run(workload) == "ready" + assert sum(item.category is dispatch.DispatchCoverageWarning for item in captured) == 1 + + +def test_probe_coverage_policy_rejects_invalid_values() -> None: + """Coverage handling uses a closed coordinator-local policy.""" + + with pytest.raises(ValueError, match="coverage_policy"): + dispatch.ProbeOptions(coverage_policy="unknown") + + +def test_known_conflict_is_ineligible() -> None: + """Known incompatible requirements always stop execution.""" @environment_req(python=">=4") @environment_req(python="<3") diff --git a/tests/dispatch/test_notebook.py b/tests/dispatch/test_notebook.py new file mode 100644 index 00000000..1bbd3549 --- /dev/null +++ b/tests/dispatch/test_notebook.py @@ -0,0 +1,124 @@ +"""Real-kernel checks for notebook Dispatch and shell task ownership.""" + +from __future__ import annotations + +import asyncio +import os + +import pytest + + +def test_repeated_notebook_dispatch_preserves_shell_and_kernel_task(tmp_path): + """Successful publications leave later cells and unrelated kernel tasks alive.""" + + pytest.importorskip("ipykernel") + client_module = pytest.importorskip("jupyter_client") + + async def exercise(): + manager = client_module.AsyncKernelManager(kernel_name="python3") + await manager.start_kernel(env={**os.environ, "DRYML_NOTEBOOK_TEST_ROOT": str(tmp_path)}) + client = manager.client() + client.start_channels() + try: + await client.wait_for_ready(timeout=30) + + async def execute(source, *, expected="ok", interrupt=False): + message_id = client.execute(source) + if interrupt: + while True: + busy = await asyncio.wait_for(client.get_iopub_msg(), timeout=30) + if busy["parent_header"].get("msg_id") == message_id and busy["msg_type"] == "status" and busy["content"]["execution_state"] == "busy": + break + marker = tmp_path / "worker-started" + deadline = asyncio.get_running_loop().time() + 20 + while not marker.exists(): + if asyncio.get_running_loop().time() >= deadline: + raise AssertionError("notebook dispatch worker never started") + await asyncio.sleep(0.05) + await manager.interrupt_kernel() + while True: + reply = await asyncio.wait_for(client.get_shell_msg(), timeout=60) + if reply["parent_header"].get("msg_id") == message_id: + break + output = [] + while True: + message = await asyncio.wait_for(client.get_iopub_msg(), timeout=60) + if message["parent_header"].get("msg_id") != message_id: + continue + if message["msg_type"] == "stream": + output.append(message["content"]["text"]) + elif message["msg_type"] == "error": + output.extend(message["content"]["traceback"]) + elif message["msg_type"] == "status" and message["content"]["execution_state"] == "idle": + status = reply["content"]["status"] + # ipykernel may signal its loop thread while a blocking shell call finishes. + assert status == expected or (interrupt and status == "ok"), "\n".join(output) + if status == "error": + assert reply["content"]["ename"] == "KeyboardInterrupt", "\n".join(output) + return (status, "\n".join(output)) if interrupt else "\n".join(output) + + await execute( + "import asyncio, os\n" + "from pathlib import Path\n" + "from dryml import dispatch\n" + "from dryml.core import Repo, StateRef\n" + "from dryml.core.execute import CoreOptions\n" + "from dryml.core.execute_codec import CoreCallCodecError\n" + "from dryml.core.store.dir import DirStore\n" + "from dryml.execute.subprocess import SubProcessConfig\n" + "from tests.dispatch.test_backend_execution import _StatefulResult\n" + "root = Path(os.environ['DRYML_NOTEBOOK_TEST_ROOT'])\n" + "(root / 'spool').mkdir()\n" + "repo = Repo(DirStore(root / 'store'))\n" + "view = dispatch.with_options(backend=SubProcessConfig(spool_directory=root / 'spool'), core=CoreOptions(repo=repo, return_objects=False))\n" + "def notebook_fn(value):\n return _StatefulResult(value)\n" + "kernel_task = asyncio.create_task(asyncio.Event().wait())\n" + "initial_shell_task = asyncio.current_task()\n" + ) + shell_status = ( + "print('shell-alive', not kernel_task.done(), " + "kernel_task in asyncio.all_tasks(), initial_shell_task.done(), " + "sum(t.get_coro().__qualname__ == 'Kernel.shell_main' " + "for t in asyncio.all_tasks()))" + ) + for count in range(1, 4): + result = await execute( + f"result = view.run(notebook_fn, {count})\n" + f"print('published', isinstance(result, StateRef), len(list(repo.find_defs(None, refresh=True))))\n" + ) + assert f"published True {count}" in result + assert "DispatchCoverageWarning" not in result + assert "Task was destroyed" not in result + alive = await execute(shell_status) + assert "shell-alive True True True 1" in alive + assert "Kernel.shell_main" not in alive + rejected = await execute( + "def captured_kernel_task():\n return kernel_task\n" + "try:\n view.run(captured_kernel_task)\n" + "except CoreCallCodecError:\n print('task-transport-rejected', not kernel_task.done())\n" + ) + assert "task-transport-rejected True" in rejected + assert "Task was destroyed" not in rejected + assert "shell-alive True True True 1" in await execute(shell_status) + status, interrupted = await execute( + "def slow_notebook_fn():\n import time\n" + f" with open({str(tmp_path / 'worker-started')!r}, 'w') as started:\n" + " started.write('ready')\n" + " time.sleep(30)\n return _StatefulResult(99)\n" + "view.run(slow_notebook_fn)", + expected="error", + interrupt=True, + ) + if status == "error": + assert "KeyboardInterrupt" in interrupted + assert "Task was destroyed" not in interrupted + after = await execute( + "print('after-interrupt', len(list(repo.find_defs(None, refresh=True))), not kernel_task.done())" + ) + assert f"after-interrupt {3 if status == 'error' else 4} True" in after + await execute("kernel_task.cancel()") + finally: + client.stop_channels() + await manager.shutdown_kernel(now=True) + + asyncio.run(exercise()) diff --git a/tests/docs/test_dispatch_documentation.py b/tests/docs/test_dispatch_documentation.py index 9c0ba317..c5391942 100644 --- a/tests/docs/test_dispatch_documentation.py +++ b/tests/docs/test_dispatch_documentation.py @@ -32,6 +32,7 @@ def test_dispatch_guide_documents_actual_defaults_and_boundaries() -> None: ), ("max_targets", f"max_targets={defaults.max_targets}"), ("max_depth", f"max_depth={defaults.max_depth}"), + ("coverage_policy", f'coverage_policy="{defaults.coverage_policy}"'), ): assert getattr(defaults, name) is not None assert rendered in guide @@ -75,7 +76,7 @@ def test_dispatch_public_surface_and_dataclass_fields_are_closed() -> None: } assert [field.name for field in fields(ProbeOptions)] == [ "placement", "backend", "environment", "world", "environment_spec", - "execution_timeout", "max_targets", "max_depth", + "execution_timeout", "max_targets", "max_depth", "coverage_policy", ] assert [field.name for field in fields(DispatchReport)] == [ "workload_placement", "workload_backend", "supported_methods", diff --git a/tests/execute/test_transport.py b/tests/execute/test_transport.py index 9edbff80..2067105a 100644 --- a/tests/execute/test_transport.py +++ b/tests/execute/test_transport.py @@ -1,6 +1,7 @@ from __future__ import annotations from dataclasses import replace +import asyncio import concurrent.futures import errno import os @@ -194,6 +195,29 @@ def closed(): resources[2].shutdown() +def test_serializer_does_not_serialize_active_asyncio_task(): + """Generic Execute rejects event-loop-owned state before Dill traverses it.""" + + async def exercise(): + task = asyncio.create_task(asyncio.Event().wait()) + + def captured(): + return task + + try: + with pytest.raises(TypeError, match="live resource"): + serialize_call(captured, (), {}, limit_bytes=1_000_000) + with pytest.raises(TypeError, match="live resource"): + serialize_call(lambda value: value, (task,), {}, limit_bytes=1_000_000) + assert not task.done() + finally: + task.cancel() + with pytest.raises(asyncio.CancelledError): + await task + + asyncio.run(exercise()) + + def test_cleanup_retry_preserves_unrelated_child_contents_and_payload_matching(tmp_path: Path, monkeypatch): """Only known files are removed; failed cleanup stays recoverable and charged.""" config = FakeConfig(spool_directory=tmp_path, spool_limit_bytes=2_000, invocation_limit_bytes=1_000, result_limit_bytes=1_000, spool_file_limit=2) diff --git a/tests/package/test_public_imports.py b/tests/package/test_public_imports.py index 93dc0d22..0809620e 100644 --- a/tests/package/test_public_imports.py +++ b/tests/package/test_public_imports.py @@ -554,7 +554,7 @@ _EXPECTED_DISPATCH_DATACLASS_FIELDS = { "ProbeOptions": [ "placement", "backend", "environment", "world", "environment_spec", - "execution_timeout", "max_targets", "max_depth", + "execution_timeout", "max_targets", "max_depth", "coverage_policy", ], "DispatchReport": [ "workload_placement", "workload_backend", "supported_methods", diff --git a/tests/test_profiles.json b/tests/test_profiles.json index e125e2eb..5ccbf475 100644 --- a/tests/test_profiles.json +++ b/tests/test_profiles.json @@ -536,7 +536,9 @@ "test_dispatch_forwards_one_frozen_exact_selector_without_reresolution", "test_target_drift_during_core_preparation_reclaims_owned_resources", "test_closed_borrowed_repo_rejects_final_acceptance_without_substitution", - "test_dispatch_run_keeps_execution_failure_primary_when_cleanup_fails" + "test_dispatch_run_keeps_execution_failure_primary_when_cleanup_fails", + "test_interrupt_during_pending_dispatch_does_not_cleanup_live_future", + "test_backend_submit_warning_policy_preserves_explain_diagnostics" ] }, "tests/dispatch/test_in_process.py": { @@ -548,7 +550,9 @@ "test_in_process_final_target_guard_runs_under_held_lease", "test_in_process_releases_failed_admission_and_rejects_local_core_controls", "test_incompatible_publication_fails_while_direct_call_holds_lease", - "test_incomplete_coverage_warns_but_known_conflict_is_ineligible" + "test_unresolved_coverage_is_quiet_by_default_but_explicitly_warns_or_rejects", + "test_bounded_static_traversal_warns_by_default", + "test_known_conflict_is_ineligible" ] }, "tests/dispatch/test_managed_execution.py": {