Repository navigation
feat(pipelines): an ergonomic client over the v2 runs API #504
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: feat/orchestration-app
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,22 @@ | ||
| # A pipeline with no pipeline API | ||
|
|
||
| The same fan-out-and-gate as [`../script`](../script), written with the client instead of | ||
| `runPipeline`. About 30 lines of ordinary JavaScript. | ||
|
|
||
| ```bash | ||
| DEEPNOTE_TOKEN=… NA_NOTEBOOK_ID=… EU_NOTEBOOK_ID=… APAC_NOTEBOOK_ID=… pnpm example:pipeline-sdk | ||
| ``` | ||
|
|
||
| The point of this example is what is missing from it. There is no workflow object, no step | ||
| registry, no graph to declare: `Promise.all` fans out, `filter` gates, `await` sequences. JavaScript | ||
| executes the pipeline, and the SDK only makes each remote operation awaitable, typed, and named. | ||
|
|
||
| That is also the trade. Because the coordination lives in this process, it is not durable: kill the | ||
| script mid-run and the notebook runs continue in Deepnote — they are detached, and their ids are | ||
| printed — but nothing aggregates them and no gate fires. Pick this up again with | ||
| `deepnote.getRun(id)`, or reach for one of the durable options in the | ||
| [package README](../../../packages/pipelines/README.md). | ||
|
|
||
| Compared with [`../script`](../script), what you give up is the execution graph and the event | ||
| stream. `runPipeline` records both; a plain function records neither, because nothing is watching. | ||
| Callers who want the graph and event stream import `runPipeline` from `@deepnote/pipelines` directly. |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,60 @@ | ||
| // The same fan-out-and-gate pipeline as ../script, written with no pipeline API at all. | ||
| // | ||
| // Every remote operation is awaitable, so the pipeline is the language: `Promise.all` fans out, | ||
| // `filter` gates, `await` sequences. Nothing here schedules, persists, replays, or supervises — | ||
| // which is why the same file runs from cron, CI, a Lambda, or another Deepnote notebook. | ||
| // | ||
| // In a real project, import from '@deepnote/pipelines' after installing it. | ||
| import { Deepnote, outputs } from '../../../packages/pipelines/dist/index.js' | ||
|
|
||
| try { | ||
| process.loadEnvFile() | ||
| } catch {} | ||
|
|
||
| const REGIONS = [ | ||
| { name: 'North America', notebookId: process.env.NA_NOTEBOOK_ID }, | ||
| { name: 'Europe', notebookId: process.env.EU_NOTEBOOK_ID }, | ||
| { name: 'Asia Pacific', notebookId: process.env.APAC_NOTEBOOK_ID }, | ||
| ].filter(region => region.notebookId) | ||
|
|
||
| if (REGIONS.length === 0) { | ||
| console.error('Set NA_NOTEBOOK_ID / EU_NOTEBOOK_ID / APAC_NOTEBOOK_ID to the notebooks to run.') | ||
| process.exit(1) | ||
| } | ||
|
|
||
| const QUALITY_THRESHOLD = 0.95 | ||
|
|
||
| // Reads DEEPNOTE_TOKEN, and DEEPNOTE_API_URL when you are not pointing at Deepnote Cloud. | ||
| const deepnote = Deepnote.fromEnv() | ||
|
|
||
| // A notebook plus the names of the values it publishes. Declared once, reused per region. | ||
| const analysis = notebookId => | ||
| deepnote.notebooks.define({ | ||
| id: notebookId, | ||
| outputs: { | ||
| // lastJson survives Deepnote reassigning block ids when it creates a notebook. | ||
| reading: outputs.lastJson(), | ||
| }, | ||
| }) | ||
|
|
||
| const started = Date.now() | ||
|
|
||
| // Fan out. Independent work is concurrent because Promise.all is, not because a framework said so. | ||
| const analyses = await Promise.all( | ||
| REGIONS.map(region => | ||
| analysis(region.notebookId).runAndWait({ | ||
| inputs: { region: region.name, trailing_months: 6 }, | ||
| onStatus: status => console.log(` ${region.name}: ${status}`), | ||
| }) | ||
| ) | ||
| ) | ||
|
|
||
| // Gate. An ordinary filter over values that are already typed and named. | ||
| const belowThreshold = analyses | ||
| .map(result => result.values.reading) | ||
| .filter(reading => reading.qualityScore < QUALITY_THRESHOLD) | ||
| .map(reading => reading.region) | ||
|
|
||
| console.log(`\n ${analyses.length} regions in ${((Date.now() - started) / 1000).toFixed(1)}s`) | ||
| console.log(` below threshold: ${belowThreshold.join(', ') || 'none'}`) | ||
| console.log(` runs: ${analyses.map(result => result.runId).join(', ')}\n`) | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -6,6 +6,78 @@ Compose Deepnote notebook runs into pipelines. No server, no kernel, no orchestr | |
| pnpm add @deepnote/pipelines | ||
| ``` | ||
|
|
||
| ## Start a run, wait for it, use its result | ||
|
|
||
| ```ts | ||
| import { Deepnote, outputs } from "@deepnote/pipelines"; | ||
|
|
||
| const deepnote = Deepnote.fromEnv(); // DEEPNOTE_TOKEN, optionally DEEPNOTE_API_URL | ||
|
|
||
| const extract = deepnote.notebooks.define({ | ||
| id: "nb-extract", | ||
| outputs: { | ||
| datasetUri: outputs.text("uri-block"), | ||
| rowCount: outputs.json<number>("stats-block", "row_count"), | ||
| }, | ||
| }); | ||
|
|
||
| // Starting and waiting are separate, because they are separate in the API. | ||
| const run = await extract.run({ inputs: { region: "eu", months: 6 } }); | ||
| console.log(run.id); // the run continues in Deepnote whether or not this process does | ||
|
|
||
| const result = await run.wait({ onStatus: (status) => console.log(status) }); | ||
| result.values.datasetUri; // string | ||
| result.values.rowCount; // number | ||
| ``` | ||
|
|
||
| `runAndWait()` is those two steps when you want them together. A failed run throws | ||
| `DeepnoteRunError` carrying the result — the snapshot is usually the only record of what the failing | ||
| block actually said — and `allowFailure: true` returns it instead. `wait({ timeoutMs })` throws | ||
| `DeepnoteRunTimeout` when the deadline passes first; only the watching stopped, the run continues in | ||
| Deepnote, and `deepnote.getRun(id)` picks it up again. | ||
|
|
||
| `extract.runs({ pageSize, pageToken })` is one page of the notebook's run history, newest first, | ||
| including runs started from Deepnote's UI; the page's `nextPageToken` fetches the next one. | ||
|
|
||
| ### A pipeline is just a function | ||
|
|
||
| There is no workflow API to learn for the common case. `await` sequences, `Promise.all` fans out, | ||
| `if` branches, `try/catch` handles failure: | ||
|
|
||
| ```ts | ||
| const [customers, products] = await Promise.all([ | ||
| deepnote.notebooks.ref("nb-customers").runAndWait({ inputs: { date } }), | ||
| deepnote.notebooks.ref("nb-products").runAndWait({ inputs: { date } }), | ||
| ]); | ||
|
|
||
| if (customers.values.rowCount > 1_000_000) { | ||
|
coderabbitai[bot] marked this conversation as resolved.
|
||
| await deepnote.notebooks.ref("nb-partition").runAndWait(); | ||
| } | ||
| ``` | ||
|
|
||
| Deepnote does not interpret that function; JavaScript does. Which means the same code runs from | ||
| cron, GitHub Actions, a Lambda, a FastAPI route, a CLI, Temporal, Airflow, or another Deepnote | ||
| notebook, and Deepnote competes with none of them. What the SDK adds is the remote operation being | ||
| awaitable, typed, and named — not a runtime that owns your control flow. | ||
|
|
||
| Reach for `runPipeline` (below) when you want what a plain function does not give you: the execution | ||
| graph, an event stream, and control nodes that make a local gate visible. | ||
|
|
||
| ### Named outputs are a client-side contract | ||
|
|
||
| Inputs are symmetrical already: `POST /v2/runs` takes values keyed by the notebook's input-block | ||
| names. Outputs are not — a finished run is a snapshot of blocks, and only the author knows which | ||
| block holds the answer. So `outputs.text()`, `outputs.json()` and `outputs.lastJson()` declare that | ||
| mapping on the client, and the error names your binding rather than a block id you never typed: | ||
|
|
||
| ``` | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 📐 Maintainability & Code Quality | 🟡 Minor | ⚡ Quick win Add a language identifier to this fenced block. Markdownlint reports MD040 for this fence. Use Proposed fix-```
+```text
Output "euShare" could not be read from totals.eu of block "stats-block" of run run-42: …🧰 Tools🪛 markdownlint-cli2 (0.23.2)[warning] 68-68: Fenced code blocks should have a language specified (MD040, fenced-code-language) 🤖 Prompt for AI AgentsSource: Linters/SAST tools |
||
| Output "euShare" could not be read from totals.eu of block "stats-block" of run run-42: … | ||
| ``` | ||
|
|
||
| `outputs.lastJson()` is the one to prefer for a notebook Deepnote created from a file, since | ||
| Deepnote reassigns block ids on creation. If Deepnote later grows a server-side notion of named | ||
| outputs, this surface does not change — only the resolver behind it does. | ||
|
|
||
| ## Run several notebooks as one pipeline | ||
|
|
||
| Fan out, gate on the results, decide: | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,140 @@ | ||
| import type { PipelineStepResult } from '../pipeline' | ||
| import { allOutputText, lastAgentText, lastOutputJson, outputJson, outputText } from '../pipeline' | ||
|
|
||
| /** | ||
| * Named outputs, declared by the caller. | ||
| * | ||
| * Deepnote's API has a clean contract for a notebook's *inputs* — named input blocks, and | ||
| * `POST /v2/runs` takes values keyed by those names. Outputs are not symmetrical: a finished run | ||
| * gives you a snapshot of blocks and Jupyter-style outputs, and which block holds "the answer" is | ||
| * something only the author knows. | ||
| * | ||
| * So the contract lives on the client for now. A binding says where a named value comes from, and | ||
| * the SDK reads it off the snapshot. If Deepnote later grows a server-side notion of named outputs, | ||
| * this surface does not have to change — only the resolver below does. | ||
| */ | ||
|
|
||
| /** Where one named output comes from, and how to read it. */ | ||
| export interface OutputBinding<T> { | ||
| /** Called with the finished step result. Throws if the value is not there. */ | ||
| read: (result: PipelineStepResult) => T | ||
| /** Human-readable source, used in error messages. */ | ||
| describe: string | ||
| } | ||
|
|
||
| /** A block's textual output. */ | ||
| export function text(blockId: string): OutputBinding<string> { | ||
| return { read: result => outputText(result, blockId), describe: `text of block "${blockId}"` } | ||
| } | ||
|
|
||
| /** | ||
| * A block's JSON output, optionally one path into it. | ||
| * | ||
| * `path` is a dotted path with numeric indexes — `totals.eu`, `regions[0].name` — deliberately not | ||
| * a full JSONPath: a binding that needs filters or wildcards is a computation, and computations | ||
| * belong in the notebook that produced the value or in the pipeline that consumes it. | ||
| */ | ||
| export function json<T = unknown>(blockId: string, path?: string): OutputBinding<T> { | ||
| return { | ||
| read: result => pluck<T>(outputJson(result, blockId), path, `block "${blockId}"`), | ||
| describe: path ? `${path} of block "${blockId}"` : `JSON of block "${blockId}"`, | ||
| } | ||
| } | ||
|
|
||
| /** | ||
| * The run's last structured JSON value, optionally one path into it. | ||
| * | ||
| * Preferred over {@link json} when a notebook was created by Deepnote from a file, because Deepnote | ||
| * reassigns block ids on creation and this does not depend on them. | ||
| */ | ||
| export function lastJson<T = unknown>(path?: string): OutputBinding<T> { | ||
| return { | ||
| read: result => pluck<T>(lastOutputJson(result), path, 'the last JSON output'), | ||
| describe: path ? `${path} of the run's last JSON output` : "the run's last JSON output", | ||
| } | ||
| } | ||
|
|
||
| /** Every block's textual output, in notebook order. Portable across remapped block ids. */ | ||
| export function allText(): OutputBinding<string> { | ||
| return { read: allOutputText, describe: "the run's textual output" } | ||
| } | ||
|
|
||
| /** The final text of the last agent block. */ | ||
| export function agentText(): OutputBinding<string> { | ||
| return { read: lastAgentText, describe: "the last agent block's text" } | ||
| } | ||
|
|
||
| /** A binding whose value the caller derives from the whole result. */ | ||
| export function derived<T>(read: (result: PipelineStepResult) => T, describe = 'a derived value'): OutputBinding<T> { | ||
| return { read, describe } | ||
| } | ||
|
|
||
| /** Named bindings for one notebook. */ | ||
| export type OutputBindings = Record<string, OutputBinding<unknown>> | ||
|
|
||
| /** The object a set of bindings produces: each name typed by its own binding. */ | ||
| export type BoundOutputs<B extends OutputBindings> = { | ||
| [K in keyof B]: B[K] extends OutputBinding<infer T> ? T : never | ||
| } | ||
|
|
||
| /** | ||
| * Resolve every binding against a finished run. | ||
| * | ||
| * One failing binding fails the whole read, and the error names the binding rather than the block: | ||
| * a caller who declared `rowCount` should be told `rowCount` is missing, not handed a block id they | ||
| * may never have typed themselves. | ||
| */ | ||
| export function resolveBindings<B extends OutputBindings>(bindings: B, result: PipelineStepResult): BoundOutputs<B> { | ||
| const resolved: Record<string, unknown> = {} | ||
| for (const [name, binding] of Object.entries(bindings)) { | ||
| try { | ||
| resolved[name] = binding.read(result) | ||
| } catch (error) { | ||
| throw new Error( | ||
| `Output "${name}" could not be read from ${binding.describe} of run ${result.runId ?? result.id}: ${ | ||
| error instanceof Error ? error.message : String(error) | ||
| }`, | ||
| { cause: error } | ||
| ) | ||
| } | ||
| } | ||
| return resolved as BoundOutputs<B> | ||
| } | ||
|
|
||
| /** Walk a dotted path with numeric indexes into a parsed JSON value. */ | ||
| function pluck<T>(value: unknown, path: string | undefined, source: string): T { | ||
| if (path === undefined) { | ||
| return value as T | ||
| } | ||
| let current = value | ||
| for (const segment of path.split('.')) { | ||
| const match = /^([A-Za-z_$][\w$]*)((?:\[\d+\])*)$/.exec(segment) | ||
| if (!match) { | ||
| throw new Error(`"${path}" is not a dotted path with numeric indexes.`) | ||
| } | ||
| current = step(current, match[1], path, source) | ||
| for (const index of match[2].matchAll(/\[(\d+)\]/g)) { | ||
| current = step(current, Number(index[1]), path, source) | ||
| } | ||
| } | ||
| return current as T | ||
| } | ||
|
|
||
| function step(value: unknown, key: string | number, path: string, source: string): unknown { | ||
| if (value === null || typeof value !== 'object') { | ||
| throw new Error(`"${path}" does not exist in ${source}: ${JSON.stringify(value)} has no "${key}".`) | ||
| } | ||
| const next = (value as Record<string | number, unknown>)[key] | ||
| if (next === undefined) { | ||
| throw new Error(`"${path}" does not exist in ${source}.`) | ||
| } | ||
| return next | ||
| } | ||
|
|
||
| /** | ||
| * The binding constructors, namespaced. | ||
| * | ||
| * A namespace rather than bare exports because `text` and `json` are far too generic to occupy a | ||
| * package's root, and because `outputs.json("block-stats", "row_count")` reads as what it is. | ||
| */ | ||
| export const outputs = { text, json, lastJson, allText, agentText, derived } as const |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
🎯 Functional Correctness | 🟡 Minor | ⚡ Quick win
Print each run id after Deepnote accepts the run.
Lines 43-50 use
runAndWait(), so no id is logged until every run completes successfully. If the process stops during the wait, users cannot usedeepnote.getRun(id)asexamples/pipelines/sdk/README.mdstates. Callrun(), logrun.id, then callwait().Proposed fix
📝 Committable suggestion
🤖 Prompt for AI Agents