Skip to content
Closed
Show file tree
Hide file tree
Changes from 1 commit
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
Prev Previous commit
Next Next commit
fix(local-runner): address review on the scheduled pipeline runner
- The Python scheduler serialized independent steps. It submitted one step's
  futures and immediately drained them before submitting the next, so only a
  fan-out's own elements overlapped — the opposite of the point, and a
  divergence from the TypeScript scheduler. Now submits every ready step before
  draining; verified with a fake API that peak concurrency is 3, not 1.
- The poll loop had no deadline and no run-id check: a run that never reached a
  terminal status hung a scheduled notebook forever, and a missing id polled
  /v2/runs/None. It now fails with a reason, and retries transient network
  errors and 429/5xx instead of abandoning a run that is still going.
- Removed a constant-false condition in _read that disabled its own null check.
  A published null is a real value and does not fall through a ?? chain, which
  is what the TypeScript resolver does.
- Conformance only compared the two implementations to each other, so a mistake
  they both made would pass. Each fixture now has a committed expected plan,
  checked by hand, and both implementations are checked against the same bad
  manifests.
- The interpreter embedded in runner.deepnote is now generated by a script and
  the test fails when the copy drifts from its source. Verified by introducing
  drift.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
  • Loading branch information
jamesbhobbs and claude committed Aug 27, 2026
commit 52b09ebe36e28d893ba76b729625ad3205d7b689
17 changes: 13 additions & 4 deletions examples/local-runner/scheduled-pipeline/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,14 @@ The notebook is self-contained: the interpreter is embedded in a code block, so
on this repository being present. Its CLI entry point is stripped, because a notebook cell runs with
`__name__ == "__main__"` and would otherwise try to parse command-line arguments.

That copy is generated, not hand-maintained. Regenerate it after changing the interpreter:

```bash
node packages/local-runner/scripts/embed-pipeline-runner.mjs
```

`pipeline-conformance.test.ts` fails if the embedded copy has drifted from its source.

## Run it locally first

```bash
Expand All @@ -45,16 +53,17 @@ DEEPNOTE_TOKEN=… python3 packages/local-runner/python/deepnote_pipeline.py --r
The interpreter exists in TypeScript (for the browser and scripts) and in Python (here). Two
implementations of one language is a standing risk that they quietly diverge, so
[`test-fixtures/pipeline-conformance`](../../../test-fixtures/pipeline-conformance) is the contract:
both planners must produce identical plans for every fixture, and `pipeline-conformance.test.ts`
fails if they do not.
both planners must produce identical plans for every fixture. Each fixture also has a committed
expected plan, because comparing the two implementations to each other cannot catch a mistake they
both make, and both are checked against the same bad manifests.

If you change `run_if`, `for_each`, `{{ }}`, or dependency derivation, change it in both and add a
fixture that would have caught the difference.

## What this does not do

Steps run concurrently within the notebook run, and the run itself is durable — but there is no
resume: if the notebook run fails halfway, rerunning starts from the beginning. Notebook runs are
Steps run concurrently within the notebook run — independent steps overlap, not just the elements of
a fan-out — and the run itself is durable. But there is no resume: if the notebook run fails halfway, rerunning starts from the beginning. Notebook runs are
not automatically idempotent, so re-running a pipeline re-runs its side effects. For replay,
per-step retries, and timers, use the Workflow SDK integration in
[`@deepnote/local-runner/workflows`](../../../packages/local-runner/README.md) instead.
83 changes: 67 additions & 16 deletions examples/local-runner/scheduled-pipeline/runner.deepnote
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@ project:
# Source of truth: packages/local-runner/python/deepnote_pipeline.py — its semantics are
# pinned against the TypeScript implementation by test-fixtures/pipeline-conformance.
# The CLI entry point is stripped: a notebook cell runs as __main__.
# Regenerate with: node packages/local-runner/scripts/embed-pipeline-runner.mjs
"""Run a ``.deepnote`` pipeline from inside a Deepnote notebook.

The TypeScript orchestrator interprets a pipeline manifest in a browser or a script. That is the
Expand Down Expand Up @@ -168,7 +169,9 @@ project:
if not isinstance(current, dict) or segment not in current:
return MISSING
current = current[segment]
return MISSING if current is None and not alternative.is_literal and False else current
# A published null is a real value and does not fall through the ?? chain; only an absent one
# does. This matches the TypeScript resolver, where `undefined` means absent and null does not.
return current


def _resolve_group(inner: str, variables: dict[str, Any]) -> Any:
Expand Down Expand Up @@ -571,13 +574,26 @@ project:
# ---------------------------------------------------------------------------


class TransientApiError(RuntimeError):
"""A failure worth retrying: a network blip, a 429, or a 5xx."""


class DeepnoteApi:
"""The three calls a step needs: start a run, poll it, read its snapshot."""

def __init__(self, token: str, base_url: str = DEFAULT_API_URL, poll_seconds: float = 3.0) -> None:
def __init__(
self,
token: str,
base_url: str = DEFAULT_API_URL,
poll_seconds: float = 3.0,
run_timeout_seconds: float = 3600.0,
max_transient_retries: int = 5,
) -> None:
self.token = token
self.base_url = base_url.rstrip("/")
self.poll_seconds = poll_seconds
self.run_timeout_seconds = run_timeout_seconds
self.max_transient_retries = max_transient_retries

def _request(self, method: str, path: str, body: dict[str, Any] | None = None) -> dict[str, Any]:
request = urllib.request.Request(
Expand All @@ -591,13 +607,39 @@ project:
return json.loads(response.read().decode() or "{}")
except urllib.error.HTTPError as error:
detail = error.read().decode(errors="replace")[:400]
raise RuntimeError(f"Deepnote API {method} {path} failed: HTTP {error.code} {detail}") from error
message = f"Deepnote API {method} {path} failed: HTTP {error.code} {detail}"
if error.code == 429 or error.code >= 500:
raise TransientApiError(message) from error
raise RuntimeError(message) from error
except (urllib.error.URLError, TimeoutError, OSError) as error:
# A network blip during a poll must not abandon a run that is still going.
raise TransientApiError(f"Deepnote API {method} {path} failed: {error}") from error

def run_notebook(self, notebook_id: str, inputs: dict[str, Any]) -> dict[str, Any]:
started = self._request("POST", "/v2/runs", {"notebookId": notebook_id, "inputs": to_run_inputs(inputs)})
run_id = started.get("runId") or started.get("id")
if not run_id:
# Without this the loop would poll /v2/runs/None until the deadline.
raise RuntimeError(f"Deepnote did not return a run id for notebook {notebook_id}: {started}")

deadline = time.monotonic() + self.run_timeout_seconds
transient_failures = 0
while True:
run = self._request("GET", f"/v2/runs/{run_id}?snapshotDelivery=inline")
if time.monotonic() >= deadline:
# A scheduled notebook should fail with a reason, not hang until someone notices.
raise RuntimeError(
f"Timed out after {self.run_timeout_seconds:.0f}s waiting for run {run_id} "
f"of notebook {notebook_id}. The run may still be going in Deepnote."
)
try:
run = self._request("GET", f"/v2/runs/{run_id}?snapshotDelivery=inline")
transient_failures = 0
except TransientApiError:
transient_failures += 1
if transient_failures > self.max_transient_retries:
raise
time.sleep(min(self.poll_seconds * 2**transient_failures, 30.0))
continue
if run.get("status") in TERMINAL_STATUSES:
return run
time.sleep(self.poll_seconds)
Expand Down Expand Up @@ -673,6 +715,11 @@ project:
if not ready:
raise RuntimeError("The pipeline stalled: no step is ready.")

# Two phases on purpose. Submitting and then immediately waiting inside one loop would
# serialize independent steps — only a fan-out's own elements would overlap — which is
# the opposite of the point and disagrees with the TypeScript scheduler.
in_flight: list[tuple[PlannedStep, list[str], list[Any]]] = []

for step in ready:
del pending[step.id]

Expand All @@ -696,32 +743,36 @@ project:
and (not step.condition or evaluate_condition(step.condition, scope))
]

if step.for_each is not None and not items:
_publish(step, [], variables)
settled.add(step.id)
continue
if not admitted:
skipped.add(step.id)
settled.add(step.id)
notify("skipped", step.id)
# A fan-out that ran nothing publishes empty lists rather than being skipped,
# whether the list was empty or every element was gated off. A plain step
# publishes a value, so "none" there really is absent.
if step.for_each is not None:
_publish(step, [], variables)
settled.add(step.id)
else:
skipped.add(step.id)
settled.add(step.id)
notify("skipped", step.id)
continue

instance_ids = []
futures = []
for instance_id, scope in admitted:
notify("started", instance_id)
futures.append(
pool.submit(api.run_notebook, step.notebook_id, resolve_value(step.inputs, scope))
)
instance_ids.append(instance_id)
futures.append(pool.submit(api.run_notebook, step.notebook_id, resolve_value(step.inputs, scope)))
in_flight.append((step, instance_ids, futures))

for step, instance_ids, futures in in_flight:
runs = []
for (instance_id, _scope), future in zip(admitted, futures):
for instance_id, future in zip(instance_ids, futures):
run = future.result()
status = run.get("status")
notify("finished", instance_id, status=status, runId=run.get("runId"))
if status != "success":
raise RuntimeError(f'Step "{instance_id}" finished with status "{status}".')
runs.append(run)

_publish(step, runs, variables)
settled.add(step.id)

Expand Down
Binary file not shown.
82 changes: 66 additions & 16 deletions packages/local-runner/python/deepnote_pipeline.py
Original file line number Diff line number Diff line change
Expand Up @@ -133,7 +133,9 @@ def _read(alternative: Alternative, variables: dict[str, Any]) -> Any:
if not isinstance(current, dict) or segment not in current:
return MISSING
current = current[segment]
return MISSING if current is None and not alternative.is_literal and False else current
# A published null is a real value and does not fall through the ?? chain; only an absent one
# does. This matches the TypeScript resolver, where `undefined` means absent and null does not.
return current


def _resolve_group(inner: str, variables: dict[str, Any]) -> Any:
Expand Down Expand Up @@ -536,13 +538,26 @@ def visit(node: str, trail: list[str]) -> None:
# ---------------------------------------------------------------------------


class TransientApiError(RuntimeError):
"""A failure worth retrying: a network blip, a 429, or a 5xx."""


class DeepnoteApi:
"""The three calls a step needs: start a run, poll it, read its snapshot."""

def __init__(self, token: str, base_url: str = DEFAULT_API_URL, poll_seconds: float = 3.0) -> None:
def __init__(
self,
token: str,
base_url: str = DEFAULT_API_URL,
poll_seconds: float = 3.0,
run_timeout_seconds: float = 3600.0,
max_transient_retries: int = 5,
) -> None:
self.token = token
self.base_url = base_url.rstrip("/")
self.poll_seconds = poll_seconds
self.run_timeout_seconds = run_timeout_seconds
self.max_transient_retries = max_transient_retries

def _request(self, method: str, path: str, body: dict[str, Any] | None = None) -> dict[str, Any]:
request = urllib.request.Request(
Expand All @@ -556,13 +571,39 @@ def _request(self, method: str, path: str, body: dict[str, Any] | None = None) -
return json.loads(response.read().decode() or "{}")
except urllib.error.HTTPError as error:
detail = error.read().decode(errors="replace")[:400]
raise RuntimeError(f"Deepnote API {method} {path} failed: HTTP {error.code} {detail}") from error
message = f"Deepnote API {method} {path} failed: HTTP {error.code} {detail}"
if error.code == 429 or error.code >= 500:
raise TransientApiError(message) from error
raise RuntimeError(message) from error
except (urllib.error.URLError, TimeoutError, OSError) as error:
# A network blip during a poll must not abandon a run that is still going.
raise TransientApiError(f"Deepnote API {method} {path} failed: {error}") from error

def run_notebook(self, notebook_id: str, inputs: dict[str, Any]) -> dict[str, Any]:
started = self._request("POST", "/v2/runs", {"notebookId": notebook_id, "inputs": to_run_inputs(inputs)})
run_id = started.get("runId") or started.get("id")
if not run_id:
# Without this the loop would poll /v2/runs/None until the deadline.
raise RuntimeError(f"Deepnote did not return a run id for notebook {notebook_id}: {started}")

deadline = time.monotonic() + self.run_timeout_seconds
transient_failures = 0
while True:
run = self._request("GET", f"/v2/runs/{run_id}?snapshotDelivery=inline")
if time.monotonic() >= deadline:
# A scheduled notebook should fail with a reason, not hang until someone notices.
raise RuntimeError(
f"Timed out after {self.run_timeout_seconds:.0f}s waiting for run {run_id} "
f"of notebook {notebook_id}. The run may still be going in Deepnote."
)
try:
run = self._request("GET", f"/v2/runs/{run_id}?snapshotDelivery=inline")
transient_failures = 0
except TransientApiError:
transient_failures += 1
if transient_failures > self.max_transient_retries:
raise
time.sleep(min(self.poll_seconds * 2**transient_failures, 30.0))
continue
if run.get("status") in TERMINAL_STATUSES:
return run
time.sleep(self.poll_seconds)
Expand Down Expand Up @@ -638,6 +679,11 @@ def notify(kind: str, step_id: str, **detail: Any) -> None:
if not ready:
raise RuntimeError("The pipeline stalled: no step is ready.")

# Two phases on purpose. Submitting and then immediately waiting inside one loop would
# serialize independent steps — only a fan-out's own elements would overlap — which is
# the opposite of the point and disagrees with the TypeScript scheduler.
in_flight: list[tuple[PlannedStep, list[str], list[Any]]] = []

for step in ready:
del pending[step.id]

Expand All @@ -661,32 +707,36 @@ def notify(kind: str, step_id: str, **detail: Any) -> None:
and (not step.condition or evaluate_condition(step.condition, scope))
]

if step.for_each is not None and not items:
_publish(step, [], variables)
settled.add(step.id)
continue
if not admitted:
skipped.add(step.id)
settled.add(step.id)
notify("skipped", step.id)
# A fan-out that ran nothing publishes empty lists rather than being skipped,
# whether the list was empty or every element was gated off. A plain step
# publishes a value, so "none" there really is absent.
if step.for_each is not None:
_publish(step, [], variables)
settled.add(step.id)
else:
skipped.add(step.id)
settled.add(step.id)
notify("skipped", step.id)
continue

instance_ids = []
futures = []
for instance_id, scope in admitted:
notify("started", instance_id)
futures.append(
pool.submit(api.run_notebook, step.notebook_id, resolve_value(step.inputs, scope))
)
instance_ids.append(instance_id)
futures.append(pool.submit(api.run_notebook, step.notebook_id, resolve_value(step.inputs, scope)))
in_flight.append((step, instance_ids, futures))

for step, instance_ids, futures in in_flight:
runs = []
for (instance_id, _scope), future in zip(admitted, futures):
for instance_id, future in zip(instance_ids, futures):
run = future.result()
status = run.get("status")
notify("finished", instance_id, status=status, runId=run.get("runId"))
if status != "success":
raise RuntimeError(f'Step "{instance_id}" finished with status "{status}".')
runs.append(run)

_publish(step, runs, variables)
settled.add(step.id)

Expand Down
6 changes: 6 additions & 0 deletions packages/local-runner/scripts/embed-pipeline-runner.d.mts
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
/** Types for the embed script, which the conformance test imports to check for drift. */
export declare const SOURCE: string
export declare const NOTEBOOK: string
export declare function embeddableSource(python: string): string
export declare function renderCell(python: string): string
export declare function replaceEmbedded(notebook: string, python: string): string
Loading
Loading