Skip to content

Commit 8161e8a

Browse files
authored
Merge pull request #53 from ncsa/fix/notebook-dispatch-lifecycle
fix(dispatch): keep notebook tasks owned and coverage warnings actionable
2 parents 7d49f70 + 2c0c183 commit 8161e8a

17 files changed

Lines changed: 418 additions & 33 deletions

‎.github/workflows/tests.yaml‎

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -74,6 +74,28 @@ jobs:
7474
./tests.sh medium --ignore tests/old --ignore tests/dev -x -vv --timeout=180 --timeout-method=thread
7575
fi
7676
77+
notebook-dispatch:
78+
name: Notebook dispatch (Ubuntu, Python 3.12)
79+
runs-on: ubuntu-latest
80+
timeout-minutes: 10
81+
82+
steps:
83+
- uses: actions/checkout@v5
84+
85+
- uses: actions/setup-python@v6
86+
with:
87+
python-version: "3.12"
88+
89+
- name: Install notebook test dependencies
90+
run: |
91+
python -m pip install --upgrade pip
92+
python -m pip install -r test_requirements.txt "ipykernel==7.2.0" "jupyter-client==8.8.0"
93+
python -m pip install .
94+
95+
- name: Run real-kernel Dispatch regression
96+
run: |
97+
./tests.sh tests/dispatch/test_notebook.py --no-cov -x --timeout=180 --timeout-method=thread --basetemp="$RUNNER_TEMP/dryml-notebook-tests"
98+
7799
local-filesystem-publication:
78100
name: Good-enough (${{ matrix.os }}, Python ${{ matrix.python-version }})
79101
runs-on: ${{ matrix.os }}

‎docs/dispatch.md‎

Lines changed: 24 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -83,17 +83,18 @@ finalization retains its existing publication behavior.
8383

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

9090
## Probing And Static Coverage
9191

9292
`ProbeOptions` is an immutable, inert policy with defaults
9393
`placement="auto"`, `execution_timeout=30.0`, `max_targets=256`, and
94-
`max_depth=32`. `placement="auto"` uses an explicit probe backend when supplied;
95-
otherwise it probes inline only with compatible current-process evidence and uses
96-
an owned local subprocess when isolation is required. `placement="execute"`
94+
`max_depth=32`, and `coverage_policy="default"`. `placement="auto"` uses an
95+
explicit probe backend when supplied; otherwise it probes inline only with
96+
compatible current-process evidence and uses an owned local subprocess when
97+
isolation is required. `placement="execute"`
9798
uses the supplied backend or that local subprocess default. `placement="in_process"`
9899
requires compatible current-process evidence and rejects a contradictory backend.
99100
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
117118
field. Capture separately limits individual source reads to 1 MiB, aggregate
118119
source reads to 8 MiB, candidate targets to 4,096, binding/call facts to 16,384,
119120
and raw annotation occurrences to 4,096. These ceilings do not make static
120-
analysis a full call-graph proof. A valid incomplete result warns and can proceed
121-
when known requirements pass; malformed projections/results, conflicts, crashes,
121+
analysis a full call-graph proof. By default, accepted incomplete coverage only
122+
warns when a traversal limit or other diagnostic accompanies `static.unresolved`;
123+
an unresolved-only call (common for notebook-defined functions and callbacks)
124+
proceeds quietly with `coverage="incomplete"` and its diagnostics retained in
125+
`explain`. `ProbeOptions(coverage_policy="warn")` restores warnings for every
126+
accepted incomplete probe. `ProbeOptions(coverage_policy="strict")` instead
127+
returns an ineligible report and rejects `run`/`submit` before workload acceptance
128+
for any incomplete probe. Malformed projections/results, conflicts, crashes,
122129
timeouts, or cleanup failures stop submission rather than becoming empty results.
123130

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

145160
The selected worker environment and world requirements are the combined
146161
configured and discovered requirements. A resolved worker Python selector is

‎docs/release_notes.md‎

Lines changed: 8 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -39,9 +39,14 @@ and async-generator roots reject before probing; lazy values returned from a
3939
synchronous root are data and are not driven by Dispatch.
4040

4141
Discovery uses shared generic static dependency, environment, and world kernels
42-
for inline and isolated probes. Incomplete but valid static coverage issues a
43-
visible warning for `run`/`submit`; conflicts, malformed results, timeout/crash,
44-
cleanup failure, target drift, and hard admission failures stop work. Reports and
42+
for inline and isolated probes. Unresolved-only static coverage remains visible
43+
in `explain` without warning on routine accepted calls; traversal limits still
44+
warn, and explicit probe policy can warn for or reject all incomplete coverage.
45+
Notebook-owned asyncio tasks, futures, and loops reject at the transport boundary
46+
rather than crossing worker serialization, and interruption of pending blocking
47+
dispatch no longer attempts premature cleanup. Conflicts, malformed results,
48+
timeout/crash, cleanup failure, target drift, and hard admission failures stop
49+
work. Reports and
4550
probe envelopes are bounded, redacted, ephemeral observations rather than Store
4651
records or admission tickets. Existing `EnvironmentSpec` selection now flows
4752
through generic/core Execute and Dispatch as an exact point-in-time pin; an

‎docs/testing.md‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -364,6 +364,13 @@ already included in their routine jobs, avoiding duplicate runners. Package
364364
tests run only for manually requested exhaustive/coverage verification, through
365365
`medium` on Ubuntu/Windows and a package-only step on macOS.
366366

367+
A separate Ubuntu Python 3.12 notebook job installs ipykernel 7.2.0 and
368+
jupyter-client 8.8.0 and runs the real-kernel Dispatch regression on every push
369+
and pull request. The general lightweight matrix does not install those optional
370+
dependencies, so its skipped notebook case is not the notebook verification gate.
371+
The test uses a job-owned kernel, temporary Store, and worker spool; it checks
372+
repeated publication/query, shell liveness, task ownership, and interrupt recovery.
373+
367374
The heavy matrix runs on Ubuntu for Python 3.10 through 3.13 only when a user
368375
manually dispatches `exhaustive` or `coverage`. It
369376
installs and preflights TensorFlow, Torch, JAX/JAXlib, and pinned

‎src/dryml/core/execute_codec.py‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@
99

1010
import builtins
1111
import ast
12+
import asyncio
1213
import dis
1314
import inspect
1415
import pathlib
@@ -95,6 +96,8 @@ def persistent_id(self, value: object) -> object | None:
9596
from dryml.managed.config import ManagedConfig
9697
from .template import TemplateBundle
9798

99+
if isinstance(value, (asyncio.Future, asyncio.AbstractEventLoop)):
100+
raise _DillLeafError("live asyncio resource")
98101
if isinstance(value, (Repo, Store)):
99102
raise _DillLeafError("live core resource")
100103
if isinstance(
@@ -427,6 +430,8 @@ def _node(self, value: Any, path: str, depth: int) -> dict[str, Any]:
427430
"receiver": self.value(value._instance, f"{path}.receiver", depth + 1),
428431
"composite": self.value(value._composite, f"{path}.composite", depth + 1),
429432
}
433+
if isinstance(value, (asyncio.Future, asyncio.AbstractEventLoop)):
434+
_fail("live asyncio resource", path)
430435
if _is_resource(value):
431436
_fail("live core resource", path)
432437
imported_capture = self.imported_captures.get(id(value))

‎src/dryml/dispatch/_preflight.py‎

Lines changed: 14 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -116,7 +116,16 @@ def _report(
116116
if probe.placement == "in_process"
117117
else "selected Execute probe configuration"
118118
)
119-
warnings = (_COVERAGE_WARNING,) if not probe.coverage.complete else ()
119+
coverage_policy = options.probe.coverage_policy
120+
warn = (
121+
eligible
122+
and not probe.coverage.complete
123+
and (coverage_policy == "warn" or (
124+
coverage_policy == "default"
125+
and set(probe.coverage.diagnostics) != {"static.unresolved"}
126+
))
127+
)
128+
warnings = (_COVERAGE_WARNING,) if warn else ()
120129
return DispatchReport(
121130
workload_placement,
122131
workload_backend,
@@ -313,6 +322,10 @@ def preflight(
313322
if eligible
314323
else "configured and discovered requirements are incompatible"
315324
)
325+
if eligible and options.probe.coverage_policy == "strict" and not result.coverage.complete:
326+
eligible = False
327+
reason = "static requirement coverage is incomplete under strict policy"
328+
outcome_diagnostics = (*outcome_diagnostics, _COVERAGE_WARNING)
316329
report = _report(
317330
options,
318331
probe=result,

‎src/dryml/dispatch/api.py‎

Lines changed: 14 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -44,9 +44,9 @@ def _checked_preflight(
4444

4545

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

49-
if prepared.report.coverage == "incomplete":
49+
if "dispatch.coverage_incomplete" in prepared.report.warnings:
5050
warnings.warn(
5151
"Dispatch static requirement coverage is incomplete",
5252
DispatchCoverageWarning,
@@ -293,6 +293,14 @@ def _run(
293293
try:
294294
result = future.result()
295295
except BaseException as error:
296+
if isinstance(error, KeyboardInterrupt) and not future.done():
297+
try:
298+
if not future.cancel():
299+
future.request_cancel()
300+
except Exception:
301+
pass
302+
# Core's one-off completion callback owns cleanup after terminality.
303+
raise
296304
try:
297305
future.cleanup()
298306
except CleanupError as cleanup_error:
@@ -319,10 +327,13 @@ def run(fn: Any, /, *args: Any, **kwargs: Any) -> Any:
319327
workload before acceptance.
320328
BaseException: Existing core/backend result and cleanup failures, with
321329
the primary execution failure retained when cleanup also fails.
330+
KeyboardInterrupt: Propagates interruption; pending backend work is
331+
first cancelled before GO or requested to stop if already running.
322332
323333
Side Effects:
324334
Runs bounded preflight and one backend submission for Execute routes.
325-
It neither retries work nor changes the selected route.
335+
It neither retries work nor changes the selected route. A pending
336+
interruption leaves terminal cleanup to the core one-off owner.
326337
"""
327338

328339
return _run(fn, args, kwargs, None)

‎src/dryml/dispatch/models.py‎

Lines changed: 10 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -41,11 +41,12 @@ class InProcess:
4141

4242

4343
class DispatchCoverageWarning(RuntimeWarning):
44-
"""Warn that valid static requirement collection was incomplete.
44+
"""Warn that accepted static requirement collection was incomplete.
4545
4646
The warning carries no workload value, source, backend credentials, or
47-
reservation. It is emitted by ``run`` and ``submit`` only; ``explain``
48-
keeps the same fact in its immutable report.
47+
reservation. ``run`` and ``submit`` emit it for bounded-analysis limits or
48+
when the caller opts into all incomplete-coverage warnings. ``explain``
49+
retains diagnostics without emitting warnings.
4950
"""
5051

5152

@@ -263,6 +264,9 @@ class ProbeOptions:
263264
checked cooperatively at inline analysis boundaries.
264265
max_targets: Positive maximum static traversal targets, including root.
265266
max_depth: Positive maximum static traversal depth.
267+
coverage_policy: ``"default"`` warns for bounded-analysis limits but
268+
not unresolved-only calls; ``"warn"`` warns for all incomplete
269+
coverage; ``"strict"`` rejects incomplete coverage before execution.
266270
267271
Raises:
268272
TypeError: If a field has an unsupported type.
@@ -281,12 +285,15 @@ class ProbeOptions:
281285
execution_timeout: float = 30.0
282286
max_targets: int = 256
283287
max_depth: int = 32
288+
coverage_policy: Literal["default", "warn", "strict"] = "default"
284289

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

288293
if self.placement not in ("auto", "in_process", "execute"):
289294
raise ValueError("probe placement is invalid")
295+
if self.coverage_policy not in ("default", "warn", "strict"):
296+
raise ValueError("probe coverage_policy is invalid")
290297
if self.backend is not None and not isinstance(
291298
self.backend, (BackendConfig, str)
292299
):

‎src/dryml/execute/_spooling.py‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@
66

77
from __future__ import annotations
88

9+
import asyncio
910
import concurrent.futures
1011
import hashlib
1112
import hmac
@@ -233,7 +234,7 @@ def _reject_live_resource(value: object) -> None:
233234
value_type = type(value)
234235
if any(cls.__module__ == "dryml.core" or cls.__module__.startswith("dryml.core.") for cls in value_type.__mro__):
235236
raise _UnsupportedTransportResource("live resource: core semantic values are unsupported by Execute transport")
236-
if isinstance(value, (io.IOBase, _LOCK_TYPES, GeneratorType, CoroutineType, socket.socket, threading.Thread, concurrent.futures.Executor, concurrent.futures.Future)):
237+
if isinstance(value, (io.IOBase, _LOCK_TYPES, GeneratorType, CoroutineType, socket.socket, threading.Thread, concurrent.futures.Executor, concurrent.futures.Future, asyncio.Future, asyncio.AbstractEventLoop)):
237238
raise _UnsupportedTransportResource("live resource: streams, locks, sockets, threads, and futures are unsupported by Execute transport")
238239
if isinstance(value, _ConnectionBase):
239240
raise _UnsupportedTransportResource("live resource: Connection is unsupported by Execute transport")

‎tests/core/test_execute_callables.py‎

Lines changed: 26 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,18 +3,21 @@
33
from __future__ import annotations
44

55
from dataclasses import dataclass
6+
import asyncio
67
import functools
78
import os
89
from pathlib import Path
910

1011
from dryml.core import ObjectRef, Repo, Serializable, function
1112
from dryml.core.execute import ExecutionContext, SharedDirStoreStrategy, invoke_prepared_call, worker_context
13+
from dryml.core.execute_codec import CoreCallCodecError, encode_invocation
1214
from dryml.core.signatures import Ref
1315
from dryml.core.store.dir import DirStore
1416
from dryml.execute import Executor
1517
from dryml.execute.subprocess import SubProcessConfig
1618
from dryml.managed import managed_operation
1719
from dryml.methods import Method
20+
import pytest
1821

1922

2023
def _invoke(fn, args, repo):
@@ -258,6 +261,29 @@ def test_fresh_subprocess_uses_dill_for_ordinary_leaves(tmp_path):
258261
assert result == FrozenOrdinaryValue(7)
259262

260263

264+
def test_core_rejects_captured_kernel_task_before_dill_traversal():
265+
"""Never serialize an active task owned by the caller's event loop."""
266+
267+
async def exercise():
268+
task = asyncio.create_task(asyncio.Event().wait())
269+
270+
def captured():
271+
return task
272+
273+
try:
274+
with pytest.raises(CoreCallCodecError, match="live asyncio resource"):
275+
encode_invocation(captured, (), {}, repo=Repo())
276+
with pytest.raises(CoreCallCodecError, match="live asyncio resource"):
277+
encode_invocation(lambda value: value, (FrozenOrdinaryValue(task),), {}, repo=Repo())
278+
assert not task.done()
279+
finally:
280+
task.cancel()
281+
with pytest.raises(asyncio.CancelledError):
282+
await task
283+
284+
asyncio.run(exercise())
285+
286+
261287
def test_copied_function_wrapper_metadata_rebuilds_its_established_owner(
262288
tmp_path):
263289
"""

0 commit comments

Comments
 (0)