Repository navigation
feat(pipelines): durable notebook steps for Workflow SDK #499
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
Closed
Closed
Changes from 6 commits
Commits
Show all changes
7 commits
Select commit
Hold shift + click to select a range
a81bf74
feat(local-runner): durable notebook steps for Workflow SDK
jamesbhobbs 00be664
Merge branch 'feat/orchestration-app' into feat/orchestration-durable
jamesbhobbs 44424a6
Merge branch 'feat/orchestration-app' into feat/orchestration-durable
jamesbhobbs b80e701
fix(local-runner): address review on durable notebook steps
jamesbhobbs edf4d71
fix(local-runner): do not declare workflow as a peer dependency
jamesbhobbs 6ab5229
Merge branch 'feat/orchestration-app' into feat/orchestration-durable
jamesbhobbs 73a3500
Merge branch 'feat/orchestration-app' into feat/orchestration-durable
jamesbhobbs File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 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
109
packages/local-runner/src/workflows/run-notebook-step.test.ts
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 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') | ||
| }) | ||
| }) |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 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 |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
Uh oh!
There was an error while loading. Please reload this page.