diff --git a/apps/docs/content/docs/cli/reference.mdx b/apps/docs/content/docs/cli/reference.mdx index 80cadfee795..9655f81a3ad 100644 --- a/apps/docs/content/docs/cli/reference.mdx +++ b/apps/docs/content/docs/cli/reference.mdx @@ -6704,7 +6704,7 @@ sim workflows run [options] | `--mock-payload` | No | Use the selected trigger's server-derived mock payload; runs the current saved workflow state (implies --manual). | | `--from-block ` | No | Run manually from this saved workflow block. | | `--source-run ` | No | Prior run whose persisted state supplies upstream outputs (requires --from-block). | -| `--stop-after ` | No | Stop the run after this saved block; with --from-block on the same block, re-runs only that block (implies --manual). | +| `--stop-after ` | No | Stop the run after this saved block, failing it if the run takes a path that skips the block; with --from-block on the same block, re-runs only that block (implies --manual). | | `--follow` | No | Stream the run as it happens; progress on stderr, result on stdout. The stream reports only success and output, so the result omits the run id and timings a non-streaming run returns. | | `--include-thinking` | No | Show model reasoning while following (requires --follow). | | `--include-tool-calls` | No | Show tool calls while following (requires --follow). | diff --git a/apps/docs/content/docs/cli/workflows.mdx b/apps/docs/content/docs/cli/workflows.mdx index dc239167f2f..4d03ddb8c55 100644 --- a/apps/docs/content/docs/cli/workflows.mdx +++ b/apps/docs/content/docs/cli/workflows.mdx @@ -642,7 +642,7 @@ sim workflows run [options] | `--mock-payload` | No | Use the selected trigger's server-derived mock payload; runs the current saved workflow state (implies --manual). | | `--from-block ` | No | Run manually from this saved workflow block. | | `--source-run ` | No | Prior run whose persisted state supplies upstream outputs (requires --from-block). | -| `--stop-after ` | No | Stop the run after this saved block; with --from-block on the same block, re-runs only that block (implies --manual). | +| `--stop-after ` | No | Stop the run after this saved block, failing it if the run takes a path that skips the block; with --from-block on the same block, re-runs only that block (implies --manual). | | `--follow` | No | Stream the run as it happens; progress on stderr, result on stdout. The stream reports only success and output, so the result omits the run id and timings a non-streaming run returns. | | `--include-thinking` | No | Show model reasoning while following (requires --follow). | | `--include-tool-calls` | No | Show tool calls while following (requires --follow). | diff --git a/apps/docs/openapi-v2-workflows.json b/apps/docs/openapi-v2-workflows.json index 2d34a84cd65..7625f7f1e0c 100644 --- a/apps/docs/openapi-v2-workflows.json +++ b/apps/docs/openapi-v2-workflows.json @@ -12617,7 +12617,7 @@ ] }, "stopAfterBlockId": { - "description": "Saved workflow block after which the run stops; downstream blocks do not execute. Must not be inside a loop or parallel. With a block entry naming the same block, re-runs only that block against the source run.", + "description": "Saved workflow block after which the run stops; downstream blocks do not execute. Must not be inside a loop or parallel. If a router, condition, or untaken error path routes the run away from the block, the run fails as soon as that is decided, without running the other branches. With a block entry naming the same block, re-runs only that block against the source run.", "type": "string", "minLength": 1 } diff --git a/apps/sim/executor/execution/edge-manager.ts b/apps/sim/executor/execution/edge-manager.ts index c190beafc6d..46aff1a3f34 100644 --- a/apps/sim/executor/execution/edge-manager.ts +++ b/apps/sim/executor/execution/edge-manager.ts @@ -122,6 +122,17 @@ export class EdgeManager { return node.incomingEdges.size === 0 || this.countActiveIncomingEdges(node) === 0 } + /** + * Whether a node that has not been queued can still run: it has received an activated edge, or + * an incoming edge is still undecided. False means every path into it was deactivated (a router + * or condition chose another route, or an error path was not taken), so nothing will queue it. + */ + canNodeStillRun(nodeId: string): boolean { + if (this.nodesWithActivatedEdge.has(nodeId)) return true + const node = this.dag.nodes.get(nodeId) + return node !== undefined && this.countActiveIncomingEdges(node) > 0 + } + restoreIncomingEdge(targetNodeId: string, sourceNodeId: string): void { const targetNode = this.dag.nodes.get(targetNodeId) if (!targetNode) { diff --git a/apps/sim/executor/execution/engine.test.ts b/apps/sim/executor/execution/engine.test.ts index b36e88bffcf..be1d8307f96 100644 --- a/apps/sim/executor/execution/engine.test.ts +++ b/apps/sim/executor/execution/engine.test.ts @@ -21,13 +21,17 @@ vi.mock('@/lib/execution/cancellation', () => ({ }, })) -import { EDGE } from '@/executor/constants' -import type { DAG, DAGNode } from '@/executor/dag/builder' -import type { EdgeManager } from '@/executor/execution/edge-manager' +import { BlockType, EDGE } from '@/executor/constants' +import { type DAG, DAGBuilder, type DAGNode } from '@/executor/dag/builder' +import { EdgeManager } from '@/executor/execution/edge-manager' import type { NodeExecutionOrchestrator } from '@/executor/orchestrators/node' import type { ExecutionContext, ExecutionResult } from '@/executor/types' import { ResolvedSecretTraceRegistry } from '@/executor/utils/resolved-secret-trace-registry' -import type { SerializedBlock } from '@/serializer/types' +import { + buildLoopSentinelEndId, + buildLoopSentinelStartId, +} from '@/executor/utils/subflow-node-id-codec' +import type { SerializedBlock, SerializedWorkflow } from '@/serializer/types' import { ExecutionEngine } from './engine' const executionEngineLoggerCallIndex = loggerMock.createLogger.mock.calls.findIndex( @@ -1088,4 +1092,309 @@ describe('ExecutionEngine', () => { expect(result.output).toEqual({ data: { response: true }, status: 200, headers: {} }) }) }) + /** + * A stop-after run is a promise that the run ends with the stop block. These run the real DAG + * builder and edge manager, so a router or condition deciding a path is the same decision the + * engine makes in production. + */ + describe('Stop-after block', () => { + function block(id: string, type = BlockType.FUNCTION): SerializedBlock { + return { ...createMockBlock(id), metadata: { id: type, name: id } } + } + + function buildRun( + workflow: SerializedWorkflow, + stopAfterBlockId: string, + outputs: Record = {} + ) { + const dag = new DAGBuilder().build(workflow, { triggerBlockId: 'start' }) + const executed: string[] = [] + const nodeOrchestrator = { + executeNode: vi.fn(async (_ctx: ExecutionContext, nodeId: string) => { + executed.push(nodeId) + return { nodeId, output: outputs[nodeId] ?? {}, isFinalOutput: false } + }), + handleNodeCompletion: vi.fn(), + } as unknown as NodeExecutionOrchestrator + const engine = new ExecutionEngine( + createMockContext({ + stopAfterBlockId, + decisions: { router: new Map(), condition: new Map() }, + }), + dag, + new EdgeManager(dag), + nodeOrchestrator + ) + return { engine, executed } + } + + /** start → condition; `if` → taken → takenTail, `else` → stop → after. */ + const conditionWorkflow: SerializedWorkflow = { + version: '1', + blocks: [ + block('start', BlockType.STARTER), + block('condition', BlockType.CONDITION), + block('taken'), + block('takenTail'), + block('stop'), + block('after'), + ], + connections: [ + { source: 'start', target: 'condition' }, + { source: 'condition', target: 'taken', sourceHandle: 'condition-if' }, + { source: 'condition', target: 'stop', sourceHandle: 'condition-else' }, + { source: 'taken', target: 'takenTail' }, + { source: 'stop', target: 'after' }, + ], + loops: {}, + parallels: {}, + } + + it('fails as soon as a condition routes away from the stop block', async () => { + const { engine, executed } = buildRun(conditionWorkflow, 'stop', { + condition: { selectedOption: 'if' }, + }) + + await expect(engine.run('start')).rejects.toThrow('Stop block "stop" (stop) was not reached') + expect(executed).toEqual(['start', 'condition']) + }) + + it('fails as soon as a router routes away from the stop block', async () => { + const { engine, executed } = buildRun( + { + version: '1', + blocks: [ + block('start', BlockType.STARTER), + block('router', BlockType.ROUTER_V2), + block('taken'), + block('takenTail'), + block('stop'), + ], + connections: [ + { source: 'start', target: 'router' }, + { source: 'router', target: 'taken', sourceHandle: 'router-route-a' }, + { source: 'router', target: 'stop', sourceHandle: 'router-route-b' }, + { source: 'taken', target: 'takenTail' }, + ], + loops: {}, + parallels: {}, + }, + 'stop', + { router: { selectedRoute: 'route-a' } } + ) + + await expect(engine.run('start')).rejects.toThrow('Stop block "stop" (stop) was not reached') + expect(executed).toEqual(['start', 'router']) + }) + + it('fails rather than pausing when another branch pauses after the stop block is skipped', async () => { + const { engine, executed } = buildRun( + { + version: '1', + blocks: [ + block('start', BlockType.STARTER), + block('approval'), + block('condition', BlockType.CONDITION), + block('taken'), + block('stop'), + ], + connections: [ + { source: 'start', target: 'approval' }, + { source: 'start', target: 'condition' }, + { source: 'condition', target: 'taken', sourceHandle: 'condition-if' }, + { source: 'condition', target: 'stop', sourceHandle: 'condition-else' }, + ], + loops: {}, + parallels: {}, + }, + 'stop', + { + approval: { + response: { status: 'paused' }, + _pauseMetadata: { + contextId: 'pause-1', + blockId: 'approval', + response: { status: 'paused' }, + timestamp: new Date().toISOString(), + pauseKind: 'hitl', + }, + }, + condition: { selectedOption: 'if' }, + } + ) + + await expect(engine.run('start')).rejects.toThrow('Stop block "stop" (stop) was not reached') + expect(executed).toEqual(['start', 'approval', 'condition']) + }) + + it('fails when the stop block sits on an error path the run never takes', async () => { + const { engine, executed } = buildRun( + { + version: '1', + blocks: [block('start', BlockType.STARTER), block('work'), block('next'), block('stop')], + connections: [ + { source: 'start', target: 'work' }, + { source: 'work', target: 'next', sourceHandle: 'source' }, + { source: 'work', target: 'stop', sourceHandle: 'error' }, + ], + loops: {}, + parallels: {}, + }, + 'stop' + ) + + await expect(engine.run('start')).rejects.toThrow('Stop block "stop" (stop) was not reached') + expect(executed).toEqual(['start', 'work']) + }) + + it('stops after the stop block when the condition routes to it', async () => { + const { engine, executed } = buildRun(conditionWorkflow, 'stop', { + condition: { selectedOption: 'else' }, + }) + + const result = await engine.run('start') + + expect(result.success).toBe(true) + expect(executed).toEqual(['start', 'condition', 'stop']) + }) + + it('runs a join reached through one taken and one skipped branch', async () => { + const { engine, executed } = buildRun( + { + version: '1', + blocks: [ + block('start', BlockType.STARTER), + block('condition', BlockType.CONDITION), + block('taken'), + block('skipped'), + block('stop'), + block('after'), + ], + connections: [ + { source: 'start', target: 'condition' }, + { source: 'condition', target: 'taken', sourceHandle: 'condition-if' }, + { source: 'condition', target: 'skipped', sourceHandle: 'condition-else' }, + { source: 'taken', target: 'stop' }, + { source: 'skipped', target: 'stop' }, + { source: 'stop', target: 'after' }, + ], + loops: {}, + parallels: {}, + }, + 'stop', + { condition: { selectedOption: 'if' } } + ) + + const result = await engine.run('start') + + expect(result.success).toBe(true) + expect(executed).toEqual(['start', 'condition', 'taken', 'stop']) + }) + + it('runs a loop stop block queued by a dead end inside the loop', async () => { + const sentinelStart = buildLoopSentinelStartId('loop') + const sentinelEnd = buildLoopSentinelEndId('loop') + const { engine, executed } = buildRun( + { + version: '1', + blocks: [ + block('start', BlockType.STARTER), + block('loop', BlockType.LOOP), + block('condition', BlockType.CONDITION), + block('inner'), + block('after'), + ], + connections: [ + { source: 'start', target: 'loop' }, + { source: 'loop', target: 'condition', sourceHandle: 'loop-start-source' }, + { source: 'condition', target: 'inner', sourceHandle: 'condition-if' }, + { source: 'loop', target: 'after', sourceHandle: 'loop-end-source' }, + ], + loops: { loop: { id: 'loop', nodes: ['condition', 'inner'], iterations: 1 } }, + parallels: {}, + }, + sentinelEnd, + { condition: { selectedOption: 'else' } } + ) + + const result = await engine.run('start') + + expect(result.success).toBe(true) + expect(executed).toEqual(['start', sentinelStart, 'condition', sentinelEnd]) + }) + + it('stops after a loop stop block whose loop has nothing to run', async () => { + const sentinelStart = buildLoopSentinelStartId('loop') + const { engine, executed } = buildRun( + { + version: '1', + blocks: [ + block('start', BlockType.STARTER), + block('loop', BlockType.LOOP), + block('inner'), + block('after'), + ], + connections: [ + { source: 'start', target: 'loop' }, + { source: 'loop', target: 'inner', sourceHandle: 'loop-start-source' }, + { source: 'loop', target: 'after', sourceHandle: 'loop-end-source' }, + ], + loops: { loop: { id: 'loop', nodes: ['inner'], iterations: 0 } }, + parallels: {}, + }, + buildLoopSentinelEndId('loop'), + { + [sentinelStart]: { sentinelStart: true, shouldExit: true, selectedRoute: EDGE.LOOP_EXIT }, + } + ) + + const result = await engine.run('start') + + expect(result.success).toBe(true) + expect(executed).toEqual(['start', sentinelStart]) + }) + + it('succeeds when the stop block is a Response block, and fails when one ends the run first', async () => { + const workflow: SerializedWorkflow = { + version: '1', + blocks: [ + block('start', BlockType.STARTER), + block('respond', BlockType.RESPONSE), + block('stop'), + ], + connections: [ + { source: 'start', target: 'respond' }, + { source: 'respond', target: 'stop' }, + ], + loops: {}, + parallels: {}, + } + + const atResponse = buildRun(workflow, 'respond') + expect((await atResponse.engine.run('start')).success).toBe(true) + + const pastResponse = buildRun(workflow, 'stop') + await expect(pastResponse.engine.run('start')).rejects.toThrow( + 'Stop block "stop" (stop) was not reached: a Response block ended the run first' + ) + expect(pastResponse.executed).toEqual(['start', 'respond']) + }) + + it('stops after the entry block when the stop block is the entry', async () => { + const { engine, executed } = buildRun(conditionWorkflow, 'start') + + const result = await engine.run('start') + + expect(result.success).toBe(true) + expect(executed).toEqual(['start']) + }) + + it('fails when the stop block is not in the graph the run executes', async () => { + const { engine } = buildRun(conditionWorkflow, 'missing', { + condition: { selectedOption: 'if' }, + }) + + await expect(engine.run('start')).rejects.toThrow('Stop block missing was not reached') + }) + }) }) diff --git a/apps/sim/executor/execution/engine.ts b/apps/sim/executor/execution/engine.ts index 86f696694cd..a2c0aebcc83 100644 --- a/apps/sim/executor/execution/engine.ts +++ b/apps/sim/executor/execution/engine.ts @@ -3,7 +3,7 @@ import { toError } from '@sim/utils/errors' import { combineExecutionAbortSignals } from '@/lib/core/execution-limits' import { subscribeToExecutionCancellation } from '@/lib/execution/cancellation' import { BlockType, EDGE } from '@/executor/constants' -import type { DAG } from '@/executor/dag/builder' +import type { DAG, DAGNode } from '@/executor/dag/builder' import type { EdgeManager } from '@/executor/execution/edge-manager' import { buildCompletedExecutionState, @@ -35,6 +35,9 @@ export class ExecutionEngine { private cancelledFlag = false private errorFlag = false private stoppedEarlyFlag = false + private stopBlockQueued = false + private stopBlockReached = false + private stopBlockUnreachable = false private executionError: Error | null = null private abortPromise!: Promise private abortResolve!: () => void @@ -125,6 +128,11 @@ export class ExecutionEngine { throw this.executionError } + /** A pause keeps a run whose stop block can still run; one proven unreachable fails. */ + if (!this.cancelledFlag && (this.stopBlockUnreachable || this.pausedBlocks.size === 0)) { + this.assertStopBlockReached() + } + if (this.pausedBlocks.size > 0) { return this.buildPausedResult(startTime) } @@ -228,6 +236,9 @@ export class ExecutionEngine { if (!this.readyQueue.includes(nodeId)) { this.readyQueue.push(nodeId) } + if (nodeId === this.context.stopAfterBlockId) { + this.stopBlockQueued = true + } } private addMultipleToQueue(nodeIds: string[]): void { @@ -483,6 +494,9 @@ export class ExecutionEngine { this.setFinalOutput(nodeId, output) this.responseOutputLocked = true } + if (this.context.stopAfterBlockId === nodeId) { + this.stopBlockReached = true + } this.stoppedEarlyFlag = true return } @@ -491,21 +505,71 @@ export class ExecutionEngine { this.setFinalOutput(nodeId, output) } - if (this.context.stopAfterBlockId === nodeId) { - // For loop/parallel sentinels, only stop if the subflow has fully exited (all iterations done) - // shouldContinue: true means more iterations, shouldExit: true means loop is done - const shouldContinue = - output.shouldContinue === true || output.selectedRoute === EDGE.PARALLEL_CONTINUE - if (!shouldContinue) { - this.execLogger.info('Stopping execution after target block', { nodeId }) - this.stoppedEarlyFlag = true - return - } + if (this.completesStopBlock(node, output)) { + this.execLogger.info('Stopping execution after target block', { nodeId }) + this.stopBlockReached = true + this.stoppedEarlyFlag = true + return } const readyNodes = this.edgeManager.processOutgoingEdges(node, output, false) this.addMultipleToQueue(readyNodes) + this.stopIfStopBlockCannotRun() + } + + /** + * Whether this completion finishes the stop block. A loop or parallel stop resolves to its end + * sentinel, which finishes only once no iteration remains; a subflow with nothing to run exits + * from its start sentinel, and its end sentinel never runs. + */ + private completesStopBlock(node: DAGNode, output: NormalizedBlockOutput): boolean { + const stopBlockId = this.context.stopAfterBlockId + if (!stopBlockId) return false + if (node.id === stopBlockId) { + return output.shouldContinue !== true && output.selectedRoute !== EDGE.PARALLEL_CONTINUE + } + const stopNode = this.dag.nodes.get(stopBlockId) + return ( + (output.selectedRoute === EDGE.LOOP_EXIT || output.selectedRoute === EDGE.PARALLEL_EXIT) && + node.metadata.sentinelType === 'start' && + stopNode?.metadata.sentinelType === 'end' && + node.metadata.subflowId === stopNode.metadata.subflowId + ) + } + + /** + * Ends the run once its stop block can no longer execute, rather than running every other branch + * to completion first; {@link assertStopBlockReached} then fails it. + */ + private stopIfStopBlockCannotRun(): void { + const stopBlockId = this.context.stopAfterBlockId + if (!stopBlockId || this.stopBlockQueued || this.edgeManager.canNodeStillRun(stopBlockId)) { + return + } + this.execLogger.info('Stopping execution: the stop block can no longer run', { stopBlockId }) + this.stopBlockUnreachable = true + this.stoppedEarlyFlag = true + } + + /** + * A stop-after run succeeds only by completing its stop block. A run whose routing skipped it, + * or that a Response block ended first, fails instead of passing for a run that stopped there. + */ + private assertStopBlockReached(): void { + const stopBlockId = this.context.stopAfterBlockId + if (!stopBlockId || this.stopBlockReached) return + const node = this.dag.nodes.get(stopBlockId) + const label = node?.metadata.isSentinel + ? (node.metadata.subflowId ?? stopBlockId) + : node?.block.metadata?.name + ? `"${node.block.metadata.name}" (${stopBlockId})` + : stopBlockId + const reason = + this.responseOutputLocked && !this.stopBlockUnreachable + ? 'a Response block ended the run first' + : 'no path this run took leads to it' + throw new Error(`Stop block ${label} was not reached: ${reason}`) } private setFinalOutput(nodeId: string, output: NormalizedBlockOutput): void { diff --git a/apps/sim/lib/api/contracts/v2/workflows.ts b/apps/sim/lib/api/contracts/v2/workflows.ts index 917d935556f..2ab8cf6d46b 100644 --- a/apps/sim/lib/api/contracts/v2/workflows.ts +++ b/apps/sim/lib/api/contracts/v2/workflows.ts @@ -1302,7 +1302,7 @@ export const v2WorkflowRunSelectionSchema = z.discriminatedUnion('source', [ .min(1, 'run.stopAfterBlockId cannot be empty') .optional() .describe( - 'Saved workflow block after which the run stops; downstream blocks do not execute. Must not be inside a loop or parallel. With a block entry naming the same block, re-runs only that block against the source run.' + 'Saved workflow block after which the run stops; downstream blocks do not execute. Must not be inside a loop or parallel. If a router, condition, or untaken error path routes the run away from the block, the run fails as soon as that is decided, without running the other branches. With a block entry naming the same block, re-runs only that block against the source run.' ), }) .strict(), diff --git a/apps/sim/scripts/test-workflow-stop-after-e2e.ts b/apps/sim/scripts/test-workflow-stop-after-e2e.ts index 691583fa980..c27ab483826 100644 --- a/apps/sim/scripts/test-workflow-stop-after-e2e.ts +++ b/apps/sim/scripts/test-workflow-stop-after-e2e.ts @@ -29,6 +29,10 @@ import { readResponseTextWithLimit } from '@/lib/core/utils/stream-limits' * in for an expensive upstream block, so a single-block re-run of Check that * finishes well under Slow's delay proves Slow was not re-executed, and an * absent After output proves the run stopped where it was told to. + * + * A second fixture branches: `Start → Gate (condition, always if)`, with `if → + * Taken → Tail (slow wait)` and `else → Skipped`. A stop on Skipped must fail the + * run once Gate routes away from it, well before Tail's delay would elapse. */ const logger = createLogger('WorkflowStopAfterE2E') const execFileAsync = promisify(execFile) @@ -76,6 +80,14 @@ interface PipelineFixture { const pipeline = fixtureIds() const otherPipeline = fixtureIds() +const branch = { + workflowId: generateId(), + start: generateId(), + gate: generateId(), + taken: generateId(), + tail: generateId(), + skipped: generateId(), +} function fixtureIds(): PipelineFixture { return { @@ -116,11 +128,6 @@ function record(value: unknown): Record { async function seedPipeline(tx: postgres.TransactionSql, fixture: PipelineFixture) { await tx`insert into workflow (id, user_id, workspace_id, name, last_synced, created_at, updated_at) values (${fixture.workflowId}, ${ownerId}, ${workspaceId}, ${`Stop-after fixture ${fixture.workflowId}`}, now(), now(), now())` - const wait = (seconds: number) => ({ - timeValue: { id: 'timeValue', type: 'short-input', value: String(seconds) }, - timeUnit: { id: 'timeUnit', type: 'dropdown', value: 'seconds' }, - async: { id: 'async', type: 'switch', value: false }, - }) const blocks = [ { id: fixture.start, @@ -128,9 +135,9 @@ async function seedPipeline(tx: postgres.TransactionSql, fixture: PipelineFixtur name: 'Start', subBlocks: { inputFormat: { id: 'inputFormat', type: 'input-format', value: [] } }, }, - { id: fixture.slow, type: 'wait', name: 'Slow', subBlocks: wait(SLOW_SECONDS) }, - { id: fixture.check, type: 'wait', name: 'Check', subBlocks: wait(0.2) }, - { id: fixture.after, type: 'wait', name: 'After', subBlocks: wait(0.2) }, + { id: fixture.slow, type: 'wait', name: 'Slow', subBlocks: waitSubBlocks(SLOW_SECONDS) }, + { id: fixture.check, type: 'wait', name: 'Check', subBlocks: waitSubBlocks(0.2) }, + { id: fixture.after, type: 'wait', name: 'After', subBlocks: waitSubBlocks(0.2) }, ] for (const [index, block] of blocks.entries()) { await tx`insert into workflow_blocks (id, workflow_id, type, name, position_x, position_y, sub_blocks) @@ -146,6 +153,59 @@ async function seedPipeline(tx: postgres.TransactionSql, fixture: PipelineFixtur } } +function waitSubBlocks(seconds: number) { + return { + timeValue: { id: 'timeValue', type: 'short-input', value: String(seconds) }, + timeUnit: { id: 'timeUnit', type: 'dropdown', value: 'seconds' }, + async: { id: 'async', type: 'switch', value: false }, + } +} + +async function seedBranch(tx: postgres.TransactionSql) { + await tx`insert into workflow (id, user_id, workspace_id, name, last_synced, created_at, updated_at) + values (${branch.workflowId}, ${ownerId}, ${workspaceId}, ${`Stop-after branch fixture ${branch.workflowId}`}, now(), now(), now())` + const conditions = [ + { id: `${branch.gate}-if`, title: 'if', value: 'true' }, + { id: `${branch.gate}-else`, title: 'else', value: '' }, + ] + const blocks = [ + { + id: branch.start, + type: 'start_trigger', + name: 'Start', + subBlocks: { inputFormat: { id: 'inputFormat', type: 'input-format', value: [] } }, + }, + { + id: branch.gate, + type: 'condition', + name: 'Gate', + subBlocks: { + conditions: { + id: 'conditions', + type: 'condition-input', + value: JSON.stringify(conditions), + }, + }, + }, + { id: branch.taken, type: 'wait', name: 'Taken', subBlocks: waitSubBlocks(0.2) }, + { id: branch.tail, type: 'wait', name: 'Tail', subBlocks: waitSubBlocks(SLOW_SECONDS) }, + { id: branch.skipped, type: 'wait', name: 'Skipped', subBlocks: waitSubBlocks(0.2) }, + ] + for (const [index, block] of blocks.entries()) { + await tx`insert into workflow_blocks (id, workflow_id, type, name, position_x, position_y, sub_blocks) + values (${block.id}, ${branch.workflowId}, ${block.type}, ${block.name}, ${index * 300}, 0, ${JSON.stringify(block.subBlocks)}::text::jsonb)` + } + for (const [source, target, sourceHandle] of [ + [branch.start, branch.gate, 'source'], + [branch.gate, branch.taken, `condition-${branch.gate}-if`], + [branch.gate, branch.skipped, `condition-${branch.gate}-else`], + [branch.taken, branch.tail, 'source'], + ]) { + await tx`insert into workflow_edges (id, workflow_id, source_block_id, target_block_id, source_handle, target_handle) + values (${generateId()}, ${branch.workflowId}, ${source}, ${target}, ${sourceHandle}, 'target')` + } +} + async function seed() { directory = await mkdtemp(resolve(tmpdir(), 'sim-stop-after-')) await sql.begin(async (tx) => { @@ -161,6 +221,7 @@ async function seed() { values (${generateId()}, ${ownerId}, 'Stop-after fixture', ${personalKey}, ${sha256Hex(personalKey)}, 'personal')` await seedPipeline(tx, pipeline) await seedPipeline(tx, otherPipeline) + await seedBranch(tx) }) } @@ -240,10 +301,22 @@ async function runCli(args: string[]): Promise { return v2ExecuteWorkflowDataSchema.parse(JSON.parse(stdout)) } +/** A CLI run the command itself must fail: exits non-zero and prints the failed run. */ +async function runCliExpectingFailure(args: string[]): Promise { + try { + await runCli(args) + } catch (error) { + assert(isRecordLike(error) && typeof error.stdout === 'string', getErrorMessage(error)) + assert.notEqual(error.code, 0, 'a failed run must exit non-zero') + return v2ExecuteWorkflowDataSchema.parse(JSON.parse(error.stdout)) + } + assert.fail('the CLI exited 0 for a run that must fail') +} + const selectAll = ['Slow.status', 'Check.status', 'After.status'] try { - await check('seed disposable workspace, personal key and two pipelines', seed) + await check('seed disposable workspace, personal key and fixtures', seed) let sourceRunId = '' await check('a full manual run executes every block and persists its state', async () => { @@ -294,6 +367,52 @@ try { assert.deepEqual(until.blockOutputs, { 'Slow.status': 'completed' }) }) + const branchOutputs = ['Taken.status', 'Tail.status', 'Skipped.status'] + + await check( + 'a stop block on the branch the condition takes still stops the run there', + async () => { + const until = await run(branch.workflowId, { + run: { source: 'manual', stopAfterBlockId: branch.taken }, + selectedOutputs: branchOutputs, + }) + assert.deepEqual(until.blockOutputs, { 'Taken.status': 'completed' }) + } + ) + + await check( + 'a stop block the condition routes away from fails the run before the other branch finishes', + async () => { + const skipped = v2ExecuteWorkflowDataSchema.parse( + record( + await execute(branch.workflowId, { + run: { source: 'manual', stopAfterBlockId: branch.skipped }, + selectedOutputs: branchOutputs, + }) + ).data + ) + assert.equal(skipped.status, 'failed') + assert.match(skipped.error?.message ?? '', /Stop block "Skipped" \(.+\) was not reached/) + assert.equal(skipped.blockOutputs?.['Skipped.status'], undefined) + assert.equal(skipped.blockOutputs?.['Tail.status'], undefined) + assert( + (skipped.durationMs ?? Number.POSITIVE_INFINITY) < SLOW_MS, + `the run must end once Gate decides, not after Tail's ${SLOW_MS} ms, took ${skipped.durationMs} ms` + ) + } + ) + + await check('the CLI exits non-zero when the stop block is not reached', async () => { + const skipped = await runCliExpectingFailure([ + branch.workflowId, + '--stop-after', + branch.skipped, + ...branchOutputs.flatMap((selector) => ['--select-output', selector]), + ]) + assert.equal(skipped.status, 'failed') + assert.match(skipped.error?.message ?? '', /was not reached/) + }) + await check('a block entry without stopAfterBlockId still runs downstream blocks', async () => { const fromCheck = await run(pipeline.workflowId, { run: { @@ -375,7 +494,7 @@ try { await check('remove disposable fixtures', async () => { // A response returns before its run finishes persisting logs and large-value // references; a cascade delete racing those writes can be chosen as a deadlock victim. - const workflowIds = [pipeline.workflowId, otherPipeline.workflowId] + const workflowIds = [pipeline.workflowId, otherPipeline.workflowId, branch.workflowId] for (let attempt = 0; attempt < 120; attempt++) { const [{ open }] = await sql`select count(*)::int as open from workflow_execution_logs where workflow_id in ${sql(workflowIds)} and ended_at is null` diff --git a/packages/sim-cli/src/commands/protocol/workflow-run-follow.ts b/packages/sim-cli/src/commands/protocol/workflow-run-follow.ts index 00f8ca36807..307b47ad01a 100644 --- a/packages/sim-cli/src/commands/protocol/workflow-run-follow.ts +++ b/packages/sim-cli/src/commands/protocol/workflow-run-follow.ts @@ -495,7 +495,7 @@ export function attachWorkflowRunFollow(workflows: Command): void { ) .option( '--stop-after ', - 'Stop the run after this saved block; with --from-block on the same block, re-runs only that block (implies --manual)' + 'Stop the run after this saved block, failing it if the run takes a path that skips the block; with --from-block on the same block, re-runs only that block (implies --manual)' ) .option( '--follow',