Skip to content

Commit b587099

Browse files
jamesbhobbsclaude
andcommitted
feat(local-runner): orchestrate notebook pipelines
Run several Deepnote notebooks as one pipeline: fan out, gate on the results, decide. The callback's own control flow is the pipeline — await, Promise.all, loops, branches — and the library records what happened: the graph, the events, the normalized results. It needs no server and no local kernel. Every step is an HTTP call to Deepnote, so nothing reachable from orchestrate() imports node:*, and the same pipeline runs in a script, in CI, and in a browser page. Notebooks are addressed by id and must already exist: running a pipeline needs permission to run a notebook, not to create one, which is what lets a page do it with a viewer's short-lived token. - orchestrate(workflow, { token }) — the pipeline API. - runOrchestration(workflow, options, executor) — the same engine with the runner left open, for callers that want to run steps somewhere else. - control() records a local decision as a graph node, so a gate is visible rather than happening invisibly between steps. - outputs.lastJson / lastAgentText read results without depending on block ids, which Deepnote reassigns when it creates a notebook. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
1 parent 726cbb5 commit b587099

11 files changed

Lines changed: 1379 additions & 31 deletions

File tree

Lines changed: 24 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,24 @@
1+
# orchestration
2+
3+
A pipeline as a script: fan out across regional notebooks, gate on their structured results, and
4+
report. About 40 lines.
5+
6+
```bash
7+
DEEPNOTE_TOKEN=… NA_NOTEBOOK_ID=… EU_NOTEBOOK_ID=… APAC_NOTEBOOK_ID=… pnpm example:orchestration
8+
```
9+
10+
The point of the example is that none of it is Node-specific. `orchestrate` runs every step as an
11+
HTTP call to Deepnote, so the same pipeline runs unchanged in a browser page with no server behind
12+
it — see [`run-app`](../run-app). What Node adds here is a shell, not a capability.
13+
14+
Notebooks are named by id and must already exist: running a pipeline needs permission to run a
15+
notebook, not to create one.
16+
17+
Each regional notebook should end by emitting a JSON object the gate can read:
18+
19+
```python
20+
print(json.dumps({"region": region, "revenueK": 812.4, "qualityScore": 0.97}))
21+
```
22+
23+
`outputs.lastJson(step)` reads that back without depending on block ids, which Deepnote reassigns
24+
when it creates a notebook.
Lines changed: 70 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,70 @@
1+
// A pipeline in ~40 lines: fan out, gate on the results, then decide.
2+
//
3+
// Nothing here is browser- or Node-specific. The same code runs in a page — every step is an HTTP
4+
// call to Deepnote, so there is no server, no local kernel, and no orchestration daemon.
5+
//
6+
// In a real project, import from '@deepnote/local-runner' after installing it.
7+
import { orchestrate } from '../../../packages/local-runner/dist/index.js'
8+
9+
try {
10+
process.loadEnvFile()
11+
} catch {}
12+
13+
const token = process.env.DEEPNOTE_TOKEN
14+
if (!token) {
15+
console.error('Set DEEPNOTE_TOKEN (or put it in a .env) to run this example.')
16+
process.exit(1)
17+
}
18+
19+
// Notebooks already in Deepnote. Pass their ids in rather than editing this file.
20+
const REGIONS = [
21+
{ name: 'North America', notebookId: process.env.NA_NOTEBOOK_ID },
22+
{ name: 'Europe', notebookId: process.env.EU_NOTEBOOK_ID },
23+
{ name: 'Asia Pacific', notebookId: process.env.APAC_NOTEBOOK_ID },
24+
].filter(region => region.notebookId)
25+
26+
if (REGIONS.length === 0) {
27+
console.error('Set NA_NOTEBOOK_ID / EU_NOTEBOOK_ID / APAC_NOTEBOOK_ID to the notebooks to run.')
28+
process.exit(1)
29+
}
30+
31+
const QUALITY_THRESHOLD = 0.95
32+
33+
const { value, graph, durationMs } = await orchestrate(
34+
async ({ run, control, outputs }) => {
35+
// Ordinary control flow is the pipeline: these fan out concurrently because Promise.all does.
36+
const analyses = await Promise.all(
37+
REGIONS.map(region =>
38+
run({
39+
id: `analyze-${slug(region.name)}`,
40+
label: `${region.name} analysis`,
41+
notebookId: region.notebookId,
42+
inputs: { region: region.name, trailing_months: 6 },
43+
})
44+
)
45+
)
46+
const readings = analyses.map(step => outputs.lastJson(step))
47+
48+
// A control node is a local decision that should still show up in the graph.
49+
const belowThreshold = await control(
50+
{
51+
id: 'quality-gate',
52+
kind: 'gate',
53+
label: `${QUALITY_THRESHOLD * 100}% quality gate`,
54+
dependsOn: analyses.map(step => step.id),
55+
},
56+
() => readings.filter(reading => reading.qualityScore < QUALITY_THRESHOLD).map(reading => reading.region)
57+
)
58+
59+
return { checked: readings.length, belowThreshold }
60+
},
61+
{ token, onEvent: event => event.type === 'step_started' && console.log(` → ${event.stepId}`) }
62+
)
63+
64+
console.log(`\n ${value.checked} regions in ${(durationMs / 1000).toFixed(1)}s`)
65+
console.log(` below threshold: ${value.belowThreshold.join(', ') || 'none'}`)
66+
console.log(` graph: ${graph.nodes.length} nodes, ${graph.edges.length} edges\n`)
67+
68+
function slug(value) {
69+
return value.toLowerCase().replace(/[^a-z0-9]+/g, '-')
70+
}

‎package.json‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@
1717
"build": "pnpm -r run build",
1818
"example:gallery": "pnpm --filter @deepnote/local-runner... build && node examples/local-runner/gallery/serve.mjs",
1919
"example:local-runner": "pnpm --filter @deepnote/local-runner... build && node examples/local-runner/run-app/serve.mjs",
20+
"example:orchestration": "pnpm --filter @deepnote/local-runner... build && node examples/local-runner/orchestration/run.mjs",
2021
"example:schedule-cloud": "pnpm --filter @deepnote/local-runner... build && node examples/local-runner/schedule-cloud.mjs",
2122
"example:snapshot-viewer": "pnpm --filter @deepnote/local-runner... build && node examples/local-runner/snapshot-viewer/serve.mjs",
2223
"license-check": "license-checker-rseidelsohn --json --onlyAllow \"MIT;Apache-2.0;BSD-2-Clause;BSD-3-Clause;ISC\" --excludePackages \"deepnote\"",

‎packages/local-runner/README.md‎

Lines changed: 59 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -183,6 +183,65 @@ The server binds to `127.0.0.1` and provides no WebSocket, watch, or rendering.
183183
— or, to _view_ an existing snapshot rather than run one, read it directly (below); that needs no
184184
server at all.
185185

186+
### Orchestrate notebook pipelines
187+
188+
Run several notebooks as one pipeline — fan out, gate on the results, decide:
189+
190+
```ts
191+
import { orchestrate } from "@deepnote/local-runner";
192+
193+
const { value, graph } = await orchestrate(
194+
async ({ run, control, outputs }) => {
195+
const analyses = await Promise.all(
196+
REGIONS.map((region) =>
197+
run({
198+
id: region.name,
199+
notebookId: region.notebookId,
200+
inputs: { region: region.name },
201+
}),
202+
),
203+
);
204+
const readings = analyses.map((step) => outputs.lastJson(step));
205+
206+
const failing = await control(
207+
{
208+
id: "quality-gate",
209+
kind: "gate",
210+
dependsOn: analyses.map((s) => s.id),
211+
},
212+
() => readings.filter((r) => r.qualityScore < 0.95).map((r) => r.region),
213+
);
214+
215+
return { checked: readings.length, failing };
216+
},
217+
{ token, onEvent: (event) => render(event) },
218+
);
219+
```
220+
221+
This is deliberately an imperative API, not a workflow language: `await`, `Promise.all`, loops and
222+
branches in the callback provide sequencing, concurrency, and conditionals. The library records what
223+
happened — the graph, the events, the normalized results.
224+
225+
**It needs no server and no local kernel.** Every step is an HTTP call to Deepnote, so the same
226+
pipeline runs in a script, in CI, and in a browser page. Nothing reachable from `orchestrate`
227+
imports `node:*`.
228+
229+
Notebooks are addressed by id and must already exist. Running a pipeline needs permission to run a
230+
notebook, not to create one — which is what lets a page do it with a viewer's short-lived token.
231+
232+
`control` records a local decision as a node, so a gate or an aggregation shows up in the graph
233+
instead of happening invisibly between steps. `outputs.lastJson(step)` and
234+
`outputs.lastAgentText(step)` read a step's results without depending on block ids, which Deepnote
235+
reassigns when it creates a notebook.
236+
237+
A failed notebook throws `OrchestrationStepError` carrying the result, so a caller can still show
238+
how far the run got; `allowFailure: true` returns it instead.
239+
240+
`runOrchestration(workflow, options, executor)` is the same engine with the runner left open, for
241+
callers that want to run steps somewhere else.
242+
243+
See [`examples/local-runner/orchestration`](../../examples/local-runner/orchestration).
244+
186245
### Read a snapshot — no Python, no kernel
187246

188247
A snapshot is a `.deepnote` file with the outputs stored inline, so reading one is parsing, not

‎packages/local-runner/src/cloud-common.ts‎

Lines changed: 2 additions & 31 deletions
Original file line numberDiff line numberDiff line change
@@ -1,10 +1,10 @@
11
import type { DeepnoteFile } from '@deepnote/blocks'
22
import { deepnoteFileSchema, deepnoteSnapshotSchema, parseYaml } from '@deepnote/blocks'
33
import { describeRunError, findNotebook, getWorkspace, type NormalizedRun, notebookUrl } from '@deepnote/cloud'
4-
import type { RunBlockOutput } from './run-with-inputs'
5-
import type { SnapshotView } from './snapshot-view'
64
import { parseSnapshot } from './snapshot-view'
75

6+
export { extractOutputs } from './extract-outputs'
7+
88
/**
99
* The plumbing every cloud entry point needs — `run-in-cloud.ts` and `cloud-runs.ts` both reach for
1010
* all of it. Internal: none of this is exported from `index.ts`.
@@ -150,32 +150,3 @@ function describeFailedBlocks(snapshotYaml: string): string | undefined {
150150
}
151151
return undefined
152152
}
153-
154-
/**
155-
* Parse the per-block outputs out of a cloud snapshot's YAML, in document order. Any executable
156-
* block type carries outputs — code, SQL, visualization, big-number — so read them off whatever
157-
* block has them (via {@link parseSnapshot}) rather than special-casing `code`.
158-
*
159-
* A snapshot that won't parse throws rather than returning nothing: the run succeeded, so "no
160-
* outputs" is a claim about the notebook, and it would be a false one. The caller still has the raw
161-
* YAML to inspect.
162-
*/
163-
export function extractOutputs(snapshotYaml: string): RunBlockOutput[] {
164-
let view: SnapshotView
165-
try {
166-
view = parseSnapshot(snapshotYaml)
167-
} catch (error) {
168-
throw new Error(
169-
`Deepnote returned a snapshot that could not be parsed: ${error instanceof Error ? error.message : String(error)}`
170-
)
171-
}
172-
const outputs: RunBlockOutput[] = []
173-
for (const notebook of view.notebooks) {
174-
for (const block of notebook.blocks) {
175-
if (block.outputs.length > 0) {
176-
outputs.push({ blockId: block.id, outputs: block.outputs, executionCount: block.executionCount })
177-
}
178-
}
179-
}
180-
return outputs
181-
}
Lines changed: 124 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,124 @@
1+
import type { NormalizedRun, PollOptions, RunInputValue, WaitForRunSnapshotOptions } from '@deepnote/cloud'
2+
import {
3+
describeRunError,
4+
isSuccessStatus,
5+
pollRunUntilComplete,
6+
triggerNotebookRun,
7+
waitForRunSnapshot,
8+
} from '@deepnote/cloud'
9+
import { extractOutputs } from './extract-outputs'
10+
import { finishResult, type OrchestrationStepExecutor } from './orchestrate'
11+
import { parseSnapshot } from './snapshot-view'
12+
13+
/**
14+
* The step executor every pipeline uses by default: three plain HTTP calls per step — start a run,
15+
* poll it, read its snapshot.
16+
*
17+
* It is `fetch` and nothing else, which is what makes a pipeline portable. The same executor runs
18+
* in a Node script, in CI, and in a browser page; there is no server in the middle and no local
19+
* kernel involved.
20+
*/
21+
22+
/** The API origin Deepnote Cloud serves from, used when a caller does not name one. */
23+
export const DEFAULT_CLOUD_API_URL = 'https://api.deepnote.com'
24+
25+
export interface CloudExecutorOptions {
26+
/**
27+
* Deepnote API token.
28+
*
29+
* In a published app this is the short-lived, viewer-scoped token the Deepnote shell issues, and
30+
* it arrives with the origin to send it to. In a script it is an ordinary API token.
31+
*/
32+
token: string
33+
/** API origin. Pair this with the token that was issued for it. */
34+
baseUrl?: string
35+
/** Poll tuning shared by every step. */
36+
poll?: Omit<PollOptions, 'onStatus'>
37+
/** Snapshot-settling tuning shared by every step. */
38+
snapshot?: WaitForRunSnapshotOptions
39+
/** Abort every in-flight request for this pipeline. */
40+
signal?: AbortSignal
41+
}
42+
43+
/**
44+
* Coerce a pipeline's input value to what `POST /v2/runs` accepts.
45+
*
46+
* The API takes exactly `string | boolean | string[]`, so numbers are stringified here rather than
47+
* rejected — a pipeline computing `monthly_target_k` should not have to remember to call
48+
* `String()`. Anything without an unambiguous textual form is refused instead of guessed at.
49+
*/
50+
export function toRunInputs(inputs: Record<string, unknown>): Record<string, RunInputValue> {
51+
const coerced: Record<string, RunInputValue> = {}
52+
for (const [name, value] of Object.entries(inputs)) {
53+
if (typeof value === 'string' || typeof value === 'boolean') {
54+
coerced[name] = value
55+
} else if (typeof value === 'number') {
56+
if (!Number.isFinite(value)) {
57+
throw new Error(`Input "${name}" is ${String(value)}, which Deepnote cannot accept.`)
58+
}
59+
coerced[name] = String(value)
60+
} else if (Array.isArray(value) && value.every(part => typeof part === 'string')) {
61+
coerced[name] = value as string[]
62+
} else if (value !== undefined && value !== null) {
63+
throw new Error(
64+
`Input "${name}" is a ${Array.isArray(value) ? 'mixed array' : typeof value}. Deepnote inputs accept a string, boolean, or array of strings.`
65+
)
66+
}
67+
}
68+
return coerced
69+
}
70+
71+
/** Build the executor {@link orchestrate} uses. Exported for callers composing their own engine. */
72+
export function createCloudStepExecutor(options: CloudExecutorOptions): OrchestrationStepExecutor {
73+
const baseUrl = options.baseUrl ?? DEFAULT_CLOUD_API_URL
74+
const { token } = options
75+
76+
return async ({ id, step, startedMs, startedAt, emit }) => {
77+
if (!step.notebookId) {
78+
throw new Error(`Orchestration step "${id}" has no notebookId to run.`)
79+
}
80+
81+
const started = await triggerNotebookRun(
82+
baseUrl,
83+
token,
84+
{ notebookId: step.notebookId, inputs: toRunInputs(step.inputs ?? {}) },
85+
{ signal: options.signal }
86+
)
87+
emit({ type: 'step_status', stepId: id, status: started.status })
88+
89+
const completed: NormalizedRun = await pollRunUntilComplete(baseUrl, token, started.runId, {
90+
...options.poll,
91+
signal: options.signal,
92+
onStatus: status => {
93+
emit({ type: 'step_status', stepId: id, status })
94+
},
95+
} as PollOptions)
96+
97+
// Read the snapshot even when the run failed: it is usually the only place the failing block's
98+
// error is recorded, and a page that shows nothing is worse than one that shows why.
99+
const settled = await waitForRunSnapshot(baseUrl, token, completed, {
100+
...options.snapshot,
101+
signal: options.signal,
102+
})
103+
const snapshotYaml = settled.content
104+
const success = isSuccessStatus(completed.status)
105+
106+
return finishResult(
107+
{
108+
id,
109+
target: 'cloud',
110+
success,
111+
status: completed.status,
112+
outputs: snapshotYaml ? extractOutputs(snapshotYaml) : [],
113+
snapshotYaml,
114+
snapshot: snapshotYaml ? parseSnapshot(snapshotYaml) : null,
115+
runId: completed.runId,
116+
error: success
117+
? undefined
118+
: (describeRunError(completed) ?? `the run finished with status "${completed.status}"`),
119+
},
120+
startedMs,
121+
startedAt
122+
)
123+
}
124+
}
Lines changed: 37 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,37 @@
1+
import type { RunBlockOutput } from './run-with-inputs'
2+
import type { SnapshotView } from './snapshot-view'
3+
import { parseSnapshot } from './snapshot-view'
4+
5+
/**
6+
* Read the per-block outputs out of a cloud snapshot's YAML, in document order.
7+
*
8+
* Split out of `cloud-common.ts` so it can bundle for the browser: everything else there resolves
9+
* tokens from `process.env`, which a page has no business carrying. Only `parseSnapshot` and the
10+
* block schemas are reachable from here.
11+
*
12+
* Any executable block type carries outputs — code, SQL, visualization, big-number — so read them
13+
* off whatever block has them rather than special-casing `code`.
14+
*
15+
* A snapshot that won't parse throws rather than returning nothing: the run succeeded, so "no
16+
* outputs" is a claim about the notebook, and it would be a false one. The caller still has the raw
17+
* YAML to inspect.
18+
*/
19+
export function extractOutputs(snapshotYaml: string): RunBlockOutput[] {
20+
let view: SnapshotView
21+
try {
22+
view = parseSnapshot(snapshotYaml)
23+
} catch (error) {
24+
throw new Error(
25+
`Deepnote returned a snapshot that could not be parsed: ${error instanceof Error ? error.message : String(error)}`
26+
)
27+
}
28+
const outputs: RunBlockOutput[] = []
29+
for (const notebook of view.notebooks) {
30+
for (const block of notebook.blocks) {
31+
if (block.outputs.length > 0) {
32+
outputs.push({ blockId: block.id, outputs: block.outputs, executionCount: block.executionCount })
33+
}
34+
}
35+
}
36+
return outputs
37+
}

0 commit comments

Comments
 (0)