Skip to content
Draft
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
refactor(pipelines): trim the public surface and drop dead result fields
Stop exporting the engine's internals (finishResult, toRunInputs, the
standalone output helpers) from the package index; callers reach the
helpers through pipelineOutputs / context.outputs. Remove the fields and
event variants nothing ever set or emitted: viewUrl and summary on step
results and graph nodes, the block_output and agent_event events, the
'branch' control kind and the node-definition version. registerNode
already validates ids, so the second validateNodeId call goes too.

RunBlockOutput moves into extract-outputs.ts; block-output.ts only
existed to mirror runtime-core's AgentStreamEvent and ExecutionSummary
for the removed fields, so it and the local-runner test asserting that
mirror are deleted.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
  • Loading branch information
jamesbhobbs and claude committed Sep 2, 2026
commit b0f5afa95ad6026bc2e40991774ce913e1e5acd8
32 changes: 0 additions & 32 deletions packages/local-runner/src/pipeline-types.test.ts

This file was deleted.

40 changes: 0 additions & 40 deletions packages/pipelines/src/block-output.ts

This file was deleted.

3 changes: 1 addition & 2 deletions packages/pipelines/src/cloud-executor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,8 +6,7 @@ import {
triggerNotebookRun,
waitForRunSnapshot,
} from '@deepnote/cloud'
import type { RunBlockOutput } from './block-output'
import { extractOutputs } from './extract-outputs'
import { extractOutputs, type RunBlockOutput } from './extract-outputs'
import { finishResult, type PipelineStepExecutor } from './pipeline'
import type { SnapshotView } from './snapshot-view'
import { parseSnapshot } from './snapshot-view'
Expand Down
14 changes: 13 additions & 1 deletion packages/pipelines/src/extract-outputs.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,19 @@
import type { RunBlockOutput } from './block-output'
import type { IOutput } from '@jupyterlab/nbformat'
import type { SnapshotView } from './snapshot-view'
import { parseSnapshot } from './snapshot-view'

/**
* What one block of a run produced.
*
* Both sides of a run want this shape and neither should reach for the other: a pipeline reads
* these off a cloud snapshot, and `@deepnote/local-runner` returns the same shape from a local kernel.
*/
export interface RunBlockOutput {
blockId: string
outputs: IOutput[]
executionCount: number | null
}

/**
* Read the per-block outputs out of a cloud snapshot's YAML, in document order.
*
Expand Down
17 changes: 3 additions & 14 deletions packages/pipelines/src/index.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
export type { AgentStreamEvent, ExecutionSummary, RunBlockOutput } from './block-output'
export type { CloudExecutorOptions } from './cloud-executor'
export { createCloudStepExecutor, DEFAULT_CLOUD_API_URL, toRunInputs } from './cloud-executor'
export { createCloudStepExecutor, DEFAULT_CLOUD_API_URL } from './cloud-executor'
export type { RunBlockOutput } from './extract-outputs'
export { extractOutputs } from './extract-outputs'
export type { InputBlockInfo } from './input-info'
export { inputInfoFor } from './input-info'
Expand All @@ -24,17 +24,6 @@ export type {
PipelineStepExecutor,
PipelineStepResult,
} from './pipeline'
export {
allOutputText,
finishResult,
lastAgentText,
lastOutputJson,
outputJson,
outputText,
PipelineStepError,
pipelineOutputs,
runPipeline,
runPipelineWithExecutor,
} from './pipeline'
export { PipelineStepError, pipelineOutputs, runPipeline, runPipelineWithExecutor } from './pipeline'
export type { SnapshotBlock, SnapshotInput, SnapshotNotebook, SnapshotView } from './snapshot-view'
export { parseSnapshot, toSnapshotView } from './snapshot-view'
38 changes: 9 additions & 29 deletions packages/pipelines/src/pipeline.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
import type { IOutput } from '@jupyterlab/nbformat'
import type { AgentStreamEvent, ExecutionSummary, RunBlockOutput } from './block-output'
import { type CloudExecutorOptions, createCloudStepExecutor } from './cloud-executor'
import type { RunBlockOutput } from './extract-outputs'
import type { SnapshotBlock, SnapshotView } from './snapshot-view'

/**
Expand Down Expand Up @@ -35,8 +35,6 @@ interface PipelineNodeDefinition {
dependsOn?: PipelineDependencyInput[]
/** Marks the node that a result viewer should select first. Only one node may be concluding. */
concluding?: boolean
/** Bump this when node behavior changes but its notebook and inputs do not. */
version?: string
/** Small, JSON-serializable annotations for generated graph renderers. */
metadata?: Record<string, string | number | boolean | null>
}
Expand All @@ -63,14 +61,14 @@ export interface PipelineStep extends PipelineNodeDefinition {
allowFailure?: boolean
}

export type PipelineControlKind = 'control' | 'gate' | 'join' | 'branch'
export type PipelineControlKind = 'control' | 'gate' | 'join'
export type PipelineGraphNodeKind = 'notebook' | PipelineControlKind

/** A local JavaScript decision or transformation that should appear in the execution graph. */
export interface PipelineControlNode extends PipelineNodeDefinition {
/** Unique within one pipeline, shared with notebook step IDs. */
id: string
/** More specific kinds let renderers distinguish validation, joining, and branching. */
/** More specific kinds let renderers distinguish a validation gate from a join. */
kind?: PipelineControlKind
}

Expand All @@ -87,8 +85,6 @@ export interface PipelineStepResult {
/** Parsed, renderer-friendly view of `snapshotYaml`, when present. */
snapshot: SnapshotView | null
runId?: string
viewUrl?: string
summary?: ExecutionSummary
error?: string
startedAt: string
finishedAt: string
Expand All @@ -110,7 +106,6 @@ export interface PipelineGraphNode {
finishedAt?: string
durationMs?: number
runId?: string
viewUrl?: string
error?: string
}

Expand All @@ -131,8 +126,6 @@ export interface PipelineGraph {
export type PipelineEvent =
| { type: 'step_started'; stepId: string; startedAt: string }
| { type: 'step_status'; stepId: string; status: string }
| { type: 'block_output'; stepId: string; blockId: string; output: IOutput }
| { type: 'agent_event'; stepId: string; event: AgentStreamEvent }
| { type: 'step_completed'; stepId: string; result: PipelineStepResult }
| { type: 'step_failed'; stepId: string; error: string; result?: PipelineStepResult }
| { type: 'control_started'; node: PipelineGraphNode }
Expand Down Expand Up @@ -218,9 +211,8 @@ export const pipelineOutputs: PipelineOutputHelpers = {
/**
* Run a pipeline in Deepnote Cloud.
*
* This is the imperative interface: the callback's own control flow is the pipeline. For a pipeline
* that should live in a file rather than in code, see `runPipelineFile`, which compiles a
* `.deepnote` definition into exactly this callback.
* This is the imperative interface: the callback's own control flow is the pipeline. Anything that
* can produce such a callback — a script, a page, a compiled definition — can drive it.
*/
export async function runPipeline<T>(
workflow: (context: PipelineContext) => T | Promise<T>,
Expand Down Expand Up @@ -290,22 +282,20 @@ export async function runPipelineWithExecutor<T>(
node: PipelineGraphNode,
status: Exclude<PipelineGraphNodeStatus, 'running'>,
startedMs: number,
details: { target?: string; runId?: string; viewUrl?: string; error?: string } = {}
details: { target?: string; runId?: string; error?: string } = {}
): void => {
node.status = status
node.finishedAt = new Date().toISOString()
node.durationMs = Date.now() - startedMs
node.target = details.target
node.runId = details.runId
node.viewUrl = details.viewUrl
node.error = details.error
}

const runNode = async (step: PipelineStep): Promise<PipelineStepResult> => {
const id = step.id.trim()
const startedMs = Date.now()
const startedAt = new Date(startedMs).toISOString()
validateNodeId(id, usedIds)
const node = registerNode(id, 'notebook', step, startedAt)
resultOrder.set(id, resultOrder.size)
emit({ type: 'step_started', stepId: id, startedAt })
Expand All @@ -316,24 +306,15 @@ export async function runPipelineWithExecutor<T>(
results.push(result)
if (!result.success) {
const error = result.error ?? `the notebook finished with status "${result.status}"`
finishNode(node, 'failed', startedMs, {
target: result.target,
runId: result.runId,
viewUrl: result.viewUrl,
error,
})
finishNode(node, 'failed', startedMs, { target: result.target, runId: result.runId, error })
emit({ type: 'step_failed', stepId: id, error, result })
if (!step.allowFailure) {
throw new PipelineStepError(id, error, { result })
}
return result
}

finishNode(node, 'success', startedMs, {
target: result.target,
runId: result.runId,
viewUrl: result.viewUrl,
})
finishNode(node, 'success', startedMs, { target: result.target, runId: result.runId })
emit({ type: 'step_completed', stepId: id, result })
return result
} catch (error) {
Expand All @@ -352,7 +333,6 @@ export async function runPipelineWithExecutor<T>(
const kind = definition.kind ?? 'control'
const startedMs = Date.now()
const startedAt = new Date(startedMs).toISOString()
validateNodeId(id, usedIds)
const node = registerNode(id, kind, definition, startedAt)
emit({ type: 'control_started', node: { ...node } })

Expand Down Expand Up @@ -413,7 +393,7 @@ function normalizeDependencies(dependencies: PipelineDependencyInput[] | undefin
return normalized
}

/** Stamp an executor result with its timing. Exported for executors outside this module. */
/** Stamp an executor result with its timing, for executors outside this module. */
export function finishResult(
result: Omit<PipelineStepResult, 'startedAt' | 'finishedAt' | 'durationMs'>,
startedMs: number,
Expand Down