Skip to content
Closed
Show file tree
Hide file tree
Changes from 6 commits
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
49 changes: 49 additions & 0 deletions packages/local-runner/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -381,6 +381,55 @@ boundary worth keeping.

See [`examples/local-runner/sales-pipeline.deepnote`](../../examples/local-runner/sales-pipeline.deepnote).

### Make a pipeline durable

`orchestrate` holds its state in one process and is gone if that process is. That is the right trade
for a script or an interactive page, and the wrong one for anything scheduled or long-lived.

Rather than growing a checkpoint/resume layer — which is how orchestration libraries turn into bad
workflow engines — durability is delegated. `@deepnote/local-runner/workflows` exposes one notebook
run as a step you compose inside a [Workflow SDK](https://www.npmjs.com/package/workflow) function:

```ts
import { lastOutputJson } from "@deepnote/local-runner";
import { runNotebookStep } from "@deepnote/local-runner/workflows";

export async function salesReview() {
"use workflow";

const regions = await Promise.all(
REGIONS.map((region) =>
runNotebookStep({ id: region.name, notebookId: region.notebookId }),
),
);
const failing = regions.filter((r) => lastOutputJson(r).qualityScore < 0.95);
Comment thread
coderabbitai[bot] marked this conversation as resolved.
Outdated
return runNotebookStep({
id: "arbiter",
notebookId: ARBITER,
inputs: { failing: failing.length },
});
}
```

Replay, retries, timers, and observability are that engine's job.

Nothing here imports `workflow` — `'use step'` is a directive its compiler reads, and without that
compiler the directive is inert and `runNotebookStep` is an ordinary async function. So this package
declares no dependency on it, not even a peer one: install
[`workflow`](https://www.npmjs.com/package/workflow) (>= 4) alongside it if you want durability, and
nothing is imposed on consumers who do not.

Two deliberate choices:

- **The token is read from the environment inside the step**, not passed as an argument, so the
credential stays out of the workflow's arguments and therefore out of its event log.
- **`maxRetries` is 0.** A notebook may write files, mutate databases, or spend model budget;
repeating that implicitly is not a safe default. A consumer who has made a notebook idempotent can
wrap it in their own step with whatever policy they want.

This is a server-side concern by definition — a durable engine needs a process that outlives a page —
which is why it is a separate entry point from the rest of the package.

### Read a snapshot — no Python, no kernel

A snapshot is a `.deepnote` file with the outputs stored inline, so reading one is parsing, not
Expand Down
7 changes: 6 additions & 1 deletion packages/local-runner/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,12 @@
"require": "./dist/index.cjs"
},
"./snapshot-reader": "./dist/snapshot-reader.iife.js",
"./orchestrator": "./dist/orchestrator.iife.js"
"./orchestrator": "./dist/orchestrator.iife.js",
"./workflows": {
"types": "./dist/workflows.d.ts",
"import": "./dist/workflows.js",
"require": "./dist/workflows.cjs"
}
},
"main": "./dist/index.cjs",
"module": "./dist/index.js",
Expand Down
8 changes: 8 additions & 0 deletions packages/local-runner/src/workflows/index.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
/**
* Durable Deepnote steps for Workflow SDK.
*
* Separate from the main entry point because it is a server-side concern: a durable engine needs a
* process that outlives a page, and this module reads the API token from the environment.
*/
export type { WorkflowCloudPollOptions, WorkflowNotebookStep } from './run-notebook-step'
export { runNotebookStep } from './run-notebook-step'
109 changes: 109 additions & 0 deletions packages/local-runner/src/workflows/run-notebook-step.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,109 @@
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'

const cloudMock = vi.hoisted(() => ({
triggerNotebookRun: vi.fn(),
pollRunUntilComplete: vi.fn(),
waitForRunSnapshot: vi.fn(),
}))

vi.mock('@deepnote/cloud', async () => {
const actual = await vi.importActual<typeof import('@deepnote/cloud')>('@deepnote/cloud')
return { ...actual, ...cloudMock }
})

import { runNotebookStep } from './run-notebook-step'

const SNAPSHOT = `metadata:
createdAt: '2026-01-01T00:00:00.000Z'
project:
id: p1
name: Durable
notebooks:
- id: nb-1
name: Main
blocks:
- blockGroup: g1
content: emit
id: b1
metadata: {}
sortingKey: a0
type: code
executionCount: 1
outputs:
- output_type: execute_result
data:
application/json:
revenue: 12
metadata: {}
version: '1.0.0'
`

beforeEach(() => {
vi.clearAllMocks()
process.env.DEEPNOTE_TOKEN = 'workflow-token'
cloudMock.triggerNotebookRun.mockResolvedValue({ runId: 'run-1', status: 'running', raw: {} })
cloudMock.pollRunUntilComplete.mockResolvedValue({ runId: 'run-1', status: 'success', raw: {} })
cloudMock.waitForRunSnapshot.mockResolvedValue({
run: { runId: 'run-1', status: 'success', raw: {} },
content: SNAPSHOT,
})
})

afterEach(() => {
delete process.env.DEEPNOTE_TOKEN
})

describe('runNotebookStep', () => {
it('runs one notebook and returns a serializable result', async () => {
const result = await runNotebookStep({ id: 'analyze', notebookId: 'nb-regional', inputs: { region: 'Europe' } })

expect(result.id).toBe('analyze')
expect(result.success).toBe(true)
expect(result.runId).toBe('run-1')
// The result crosses a durable step boundary, so it must survive a round trip.
expect(JSON.parse(JSON.stringify(result)).outputs[0].blockId).toBe('b1')
})

it('reads the token from the environment, keeping it out of the step arguments', async () => {
await runNotebookStep({ id: 'a', notebookId: 'nb-a' })

expect(cloudMock.triggerNotebookRun).toHaveBeenCalledWith(
'https://api.deepnote.com',
'workflow-token',
{ notebookId: 'nb-a', inputs: {} },
expect.anything()
)
})

it('honours a custom API origin', async () => {
await runNotebookStep({ id: 'a', notebookId: 'nb-a', baseUrl: 'https://api.example.test' })
expect(cloudMock.triggerNotebookRun).toHaveBeenCalledWith(
'https://api.example.test',
'workflow-token',
expect.anything(),
expect.anything()
)
})

it('says what is missing when there is no token', async () => {
delete process.env.DEEPNOTE_TOKEN
await expect(runNotebookStep({ id: 'a', notebookId: 'nb-a' })).rejects.toThrow('requires DEEPNOTE_TOKEN')
})

it('never retries by default, because a notebook run has side effects', () => {
expect(runNotebookStep.maxRetries).toBe(0)
})

it('returns a failed run when it is allowed to fail, rather than throwing', async () => {
cloudMock.pollRunUntilComplete.mockResolvedValue({
runId: 'run-1',
status: 'error',
error: { message: 'the warehouse is down' },
raw: {},
})

const result = await runNotebookStep({ id: 'a', notebookId: 'nb-a', allowFailure: true })
expect(result.success).toBe(false)
expect(result.error).toContain('warehouse')
})
})
71 changes: 71 additions & 0 deletions packages/local-runner/src/workflows/run-notebook-step.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,71 @@
import type { PollOptions } from '@deepnote/cloud'
import { createCloudStepExecutor } from '../cloud-executor'
import type { OrchestrationStepResult } from '../orchestrate'
import { runOrchestration } from '../orchestrate'

/**
* A Deepnote notebook run as a durable step.
*
* `orchestrate` is deliberately not durable: it holds its state in one process and is gone if that
* process is. That is the right trade for a script or an interactive page, and the wrong one for
* anything scheduled or long-lived.
*
* Rather than growing a checkpoint/resume layer of its own — which is how orchestration libraries
* become bad workflow engines — this delegates. Compose these steps inside a
* [Workflow SDK](https://www.npmjs.com/package/workflow) function and durability, replay, and
* observability are that engine's job. `workflow` is an optional peer dependency: without its
* compiler the `'use step'` directive is inert and this is an ordinary async function.
*/

/**
* Poll options that can safely cross a step boundary.
*
* Callback and injectable-clock fields are deliberately absent because functions are not
* serializable, and a durable engine provides its own timeline and observability anyway.
*/
export type WorkflowCloudPollOptions = Pick<
PollOptions,
'intervalMs' | 'timeoutMs' | 'requestTimeoutMs' | 'maxTransientRetries' | 'snapshotDelivery'
>

/** Serializable description of one Deepnote notebook run. */
export interface WorkflowNotebookStep {
id: string
/** The Deepnote notebook to run. It must already exist. */
notebookId: string
inputs?: Record<string, unknown>
/** Return a failed run instead of throwing, so the workflow can decide what to do about it. */
allowFailure?: boolean
/** API origin. Defaults to Deepnote Cloud. */
baseUrl?: string
poll?: WorkflowCloudPollOptions
}

/**
* Run one Deepnote notebook as a durable step.
*
* The token is read from the environment inside the step rather than taken as an argument, which
* keeps the credential out of the workflow's arguments and therefore out of its event log.
*/
export async function runNotebookStep(step: WorkflowNotebookStep): Promise<OrchestrationStepResult> {
'use step'

// Matches the CLI's variable, so a workflow host configured for Deepnote already has it.
const token = process.env.DEEPNOTE_TOKEN
if (!token) {
throw new Error('runNotebookStep requires DEEPNOTE_TOKEN in the environment.')
}

const { id, notebookId, inputs, allowFailure } = step
const result = await runOrchestration(
({ run }) => run({ id, notebookId, inputs, allowFailure }),
{},
createCloudStepExecutor({ token, baseUrl: step.baseUrl, poll: step.poll })
)
return result.value
}

// A notebook may write files, mutate databases, or spend model budget. Do not repeat those side
// effects implicitly. A consumer who has made a notebook idempotent can wrap it in their own step
// with whatever retry policy they want.
runNotebookStep.maxRetries = 0
3 changes: 2 additions & 1 deletion packages/local-runner/tsdown.config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,11 +4,12 @@ export default defineConfig([
{
entry: {
index: 'src/index.ts',
workflows: 'src/workflows/index.ts',
},
format: ['esm', 'cjs'],
fixedExtension: false,
dts: true,
external: ['@deepnote/blocks', '@deepnote/cloud', '@deepnote/convert', '@deepnote/runtime-core'],
external: ['@deepnote/blocks', '@deepnote/cloud', '@deepnote/convert', '@deepnote/runtime-core', 'workflow'],
},
{
// The browser build ships as one self-contained file that a static page can <script> in, so
Expand Down
Loading