diff --git a/.github/workflows/test-build.yml b/.github/workflows/test-build.yml index e6ef55b6ae2..98abe3bf0ab 100644 --- a/.github/workflows/test-build.yml +++ b/.github/workflows/test-build.yml @@ -253,6 +253,51 @@ jobs: VERSION_COMPARE_E2E_REPORT_PATH="$report_dir/version-compare-http-report.json" \ bun run test:workflow-version-compare:e2e + # A self-hosted app: hosted billing admits a run only through a Redis usage + # reservation, and the SCIM suite above asserts PostgreSQL rate-limit storage, + # so workflow execution gets its own app rather than adding Redis to that one. + - name: Verify single-block workflow runs over real HTTP + working-directory: apps/sim + env: + NEXT_PUBLIC_APP_URL: http://127.0.0.1:3018 + BETTER_AUTH_URL: http://127.0.0.1:3018 + NEXT_PUBLIC_FORCE_HOSTED: 'false' + INTERNAL_API_SECRET: stop-after-http-ci-local-secret-at-least-32-characters + DB_TX_TRIPWIRE: throw + DISABLE_TELEMETRY: 'true' + NEXT_TELEMETRY_DISABLED: '1' + NEXT_PUBLIC_CHAT_DISABLED: 'true' + READY_TIMEOUT_SECONDS: 300 + run: | + report_dir="$RUNNER_TEMP/e2e" + server_log="$report_dir/stop-after-next.log" + mkdir -p "$report_dir" + node ../../node_modules/next/dist/bin/next dev --hostname 127.0.0.1 --port 3018 > "$server_log" 2>&1 & + server_pid=$! + finish() { + kill "$server_pid" 2>/dev/null || true + wait "$server_pid" 2>/dev/null || true + awk '/^ (GET|POST|PUT|PATCH|DELETE|HEAD) \/api\// { print }' "$server_log" > "$report_dir/stop-after-http-status.log" + } + trap finish EXIT + fail_startup() { + echo "::error::$1" + tail -n 200 "$server_log" + exit 1 + } + started=$SECONDS + until curl --fail --silent --max-time 10 http://127.0.0.1:3018/api/health > /dev/null; do + kill -0 "$server_pid" 2>/dev/null || fail_startup 'Local workflow app exited during startup.' + [ $((SECONDS - started)) -lt "$READY_TIMEOUT_SECONDS" ] || + fail_startup "Local workflow app did not become ready within $READY_TIMEOUT_SECONDS seconds." + sleep 2 + done + echo "Local workflow app ready after $((SECONDS - started))s" + STOP_AFTER_E2E_BASE_URL="$NEXT_PUBLIC_APP_URL" \ + STOP_AFTER_E2E_DATABASE_URL="$DATABASE_URL" \ + STOP_AFTER_E2E_REPORT_PATH="$report_dir/stop-after-http-report.json" \ + bun run test:workflow-stop-after:e2e + - name: Upload end-to-end reports and server logs if: failure() uses: actions/upload-artifact@ea165f8d65b6e75b540449e92b4886f43607fa02 # v4 diff --git a/apps/docs/content/docs/cli/reference.mdx b/apps/docs/content/docs/cli/reference.mdx index f7e2a05d058..80cadfee795 100644 --- a/apps/docs/content/docs/cli/reference.mdx +++ b/apps/docs/content/docs/cli/reference.mdx @@ -6704,6 +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). | | `--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 359f2be643a..dc239167f2f 100644 --- a/apps/docs/content/docs/cli/workflows.mdx +++ b/apps/docs/content/docs/cli/workflows.mdx @@ -642,6 +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). | | `--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 4d843361eb6..2d34a84cd65 100644 --- a/apps/docs/openapi-v2-workflows.json +++ b/apps/docs/openapi-v2-workflows.json @@ -12615,6 +12615,11 @@ "additionalProperties": false } ] + }, + "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.", + "type": "string", + "minLength": 1 } }, "required": ["source"], @@ -12703,6 +12708,17 @@ "sourceRunId": "run_123" } } + }, + { + "run": { + "source": "manual", + "entry": { + "type": "block", + "blockId": "block_123", + "sourceRunId": "run_123" + }, + "stopAfterBlockId": "block_123" + } } ] }, diff --git a/apps/sim/app/api/v2/workflows/[workflowId]/execute/route.ts b/apps/sim/app/api/v2/workflows/[workflowId]/execute/route.ts index 0bb2bc272a2..b1a5a8e9c44 100644 --- a/apps/sim/app/api/v2/workflows/[workflowId]/execute/route.ts +++ b/apps/sim/app/api/v2/workflows/[workflowId]/execute/route.ts @@ -448,6 +448,7 @@ export const POST = withRouteHandler( mode: body.stream ? 'stream' : resultStream ? 'sync-result-stream' : 'sync', blockId: manualRun.entry.blockId, sourceRunId: manualRun.entry.sourceRunId, + stopAfterBlockId: manualRun.stopAfterBlockId, }, request: req, }) @@ -460,6 +461,7 @@ export const POST = withRouteHandler( mode: body.stream ? 'stream' : resultStream ? 'sync-result-stream' : 'sync', triggerBlockId: manualRun.entry?.blockId, useMockPayload: manualRun.entry?.useMockPayload === true, + stopAfterBlockId: manualRun.stopAfterBlockId, }, request: req, }) diff --git a/apps/sim/lib/api/contracts/v2/workflows.ts b/apps/sim/lib/api/contracts/v2/workflows.ts index 7b3c4812b99..917d935556f 100644 --- a/apps/sim/lib/api/contracts/v2/workflows.ts +++ b/apps/sim/lib/api/contracts/v2/workflows.ts @@ -1297,6 +1297,13 @@ export const v2WorkflowRunSelectionSchema = z.discriminatedUnion('source', [ .describe( 'Manual entry mode. Omit to enter through the workflow trigger; a block entry requires an exact source run.' ), + stopAfterBlockId: z + .string() + .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.' + ), }) .strict(), ]) @@ -1412,6 +1419,13 @@ export const v2ExecuteWorkflowBodySchema = z entry: { type: 'block', blockId: 'block_123', sourceRunId: 'run_123' }, }, }, + { + run: { + source: 'manual', + entry: { type: 'block', blockId: 'block_123', sourceRunId: 'run_123' }, + stopAfterBlockId: 'block_123', + }, + }, ], }) export type V2ExecuteWorkflowBody = z.input diff --git a/apps/sim/lib/workflows/application/execute-manual-workflow.test.ts b/apps/sim/lib/workflows/application/execute-manual-workflow.test.ts index ce403f1dabd..200241c8656 100644 --- a/apps/sim/lib/workflows/application/execute-manual-workflow.test.ts +++ b/apps/sim/lib/workflows/application/execute-manual-workflow.test.ts @@ -300,6 +300,124 @@ describe('manual workflow execution application operations', () => { expect(mocks.loadSourceState).not.toHaveBeenCalled() }) + it('rejects a stop block missing from the saved workflow before anything runs', async () => { + await expect( + executeManualWorkflowOperation.execute({ + principal, + input: { ...baseInput, useMockPayload: false, stopAfterBlockId: 'missing' }, + }) + ).rejects.toMatchObject({ code: 'validation', message: expect.stringContaining('not a block') }) + await expect( + executeManualWorkflowFromBlockOperation.execute({ + principal, + input: { + ...baseInput, + blockId: 'agent-1', + sourceRunId: 'source-run-1', + stopAfterBlockId: 'missing', + }, + }) + ).rejects.toMatchObject({ code: 'validation', message: expect.stringContaining('not a block') }) + expect(mocks.loadSourceState).not.toHaveBeenCalled() + expect(mocks.executeService).not.toHaveBeenCalled() + }) + + it('rejects a stop block nested in a loop or parallel, which the engine cannot stop on', async () => { + mockLoadManualState.mockResolvedValue({ + blocks: { + 'trigger-1': {}, + 'loop-1': { type: 'loop' }, + 'agent-1': { data: { parentId: 'loop-1' } }, + }, + edges: [], + }) + + await expect( + executeManualWorkflowOperation.execute({ + principal, + input: { ...baseInput, useMockPayload: false, stopAfterBlockId: 'agent-1' }, + }) + ).rejects.toMatchObject({ code: 'validation', message: expect.stringContaining('inside loop') }) + expect(mocks.executeService).not.toHaveBeenCalled() + }) + + it('rejects a stop block the run cannot reach from its entry, which would run everything', async () => { + mockLoadManualState.mockResolvedValue({ + blocks: { 'trigger-1': {}, 'agent-1': {}, 'agent-2': {}, unconnected: {} }, + edges: [ + { source: 'trigger-1', target: 'agent-1' }, + { source: 'agent-1', target: 'agent-2' }, + ], + }) + + await expect( + executeManualWorkflowFromBlockOperation.execute({ + principal, + input: { + ...baseInput, + blockId: 'agent-2', + sourceRunId: 'source-run-1', + stopAfterBlockId: 'agent-1', + }, + }) + ).rejects.toMatchObject({ + code: 'validation', + message: expect.stringContaining('not reachable'), + }) + await expect( + executeManualWorkflowOperation.execute({ + principal, + input: { ...baseInput, useMockPayload: false, stopAfterBlockId: 'unconnected' }, + }) + ).rejects.toMatchObject({ + code: 'validation', + message: expect.stringContaining('not reachable'), + }) + expect(mocks.loadSourceState).not.toHaveBeenCalled() + expect(mocks.executeService).not.toHaveBeenCalled() + }) + + it('rejects a disabled stop block, or one reached only through a disabled block', async () => { + mockLoadManualState.mockResolvedValue({ + blocks: { + 'trigger-1': {}, + 'agent-1': { enabled: false }, + 'agent-2': {}, + }, + edges: [ + { source: 'trigger-1', target: 'agent-1' }, + { source: 'agent-1', target: 'agent-2' }, + ], + }) + + await expect( + executeManualWorkflowOperation.execute({ + principal, + input: { ...baseInput, useMockPayload: false, stopAfterBlockId: 'agent-1' }, + }) + ).rejects.toMatchObject({ code: 'validation', message: expect.stringContaining('is disabled') }) + await expect( + executeManualWorkflowOperation.execute({ + principal, + input: { ...baseInput, useMockPayload: false, stopAfterBlockId: 'agent-2' }, + }) + ).rejects.toMatchObject({ + code: 'validation', + message: expect.stringContaining('not reachable'), + }) + expect(mocks.executeService).not.toHaveBeenCalled() + }) + + it('rejects a stop block named by an inherited object key', async () => { + await expect( + executeManualWorkflowOperation.execute({ + principal, + input: { ...baseInput, useMockPayload: false, stopAfterBlockId: 'toString' }, + }) + ).rejects.toMatchObject({ code: 'validation', message: expect.stringContaining('not a block') }) + expect(mocks.executeService).not.toHaveBeenCalled() + }) + it('rejects a source run without persisted state for this workflow', async () => { mocks.loadSourceState.mockResolvedValueOnce(null) diff --git a/apps/sim/lib/workflows/application/execute-manual-workflow.ts b/apps/sim/lib/workflows/application/execute-manual-workflow.ts index 840e89d5108..4d7049e7161 100644 --- a/apps/sim/lib/workflows/application/execute-manual-workflow.ts +++ b/apps/sim/lib/workflows/application/execute-manual-workflow.ts @@ -20,6 +20,7 @@ interface ManualExecutionInput extends Omit { input?: unknown mode: 'sync' | 'stream' | 'sync-result-stream' + stopAfterBlockId?: string } export interface ExecuteManualWorkflowInput extends ManualExecutionInput { @@ -47,6 +48,69 @@ async function loadManualState(workflowId: string) { return state } +type ManualWorkflowState = Awaited> + +/** Blocks a run entering at `entryBlockId` can reach; the executor skips disabled blocks. */ +function reachableFrom(state: ManualWorkflowState, entryBlockId: string): Set { + const targetsBySource = new Map() + for (const edge of state.edges) { + const targets = targetsBySource.get(edge.source) ?? [] + targets.push(edge.target) + targetsBySource.set(edge.source, targets) + } + const reached = new Set([entryBlockId]) + const queue = [entryBlockId] + for (let next = queue.pop(); next !== undefined; next = queue.pop()) { + for (const target of targetsBySource.get(next) ?? []) { + if (reached.has(target) || state.blocks[target]?.enabled === false) continue + reached.add(target) + queue.push(target) + } + } + return reached +} + +/** + * The engine stops only when it completes a node whose id equals the target, so + * a target the run cannot reach would silently run everything after the entry: + * an unknown or disabled block, a block upstream of the entry or only behind a + * disabled one, or a block inside a loop or parallel (which would stop after its + * first iteration or never). All are refused, matching the editor, which offers + * "Run until block" only outside subflows. + */ +function assertStopAfterBlock( + state: ManualWorkflowState, + blockId: string | undefined, + entryBlockId: string +): void { + if (blockId === undefined) return + const block = Object.hasOwn(state.blocks, blockId) ? state.blocks[blockId] : undefined + if (!block) { + throw new OrchestrationError( + 'validation', + `run.stopAfterBlockId "${blockId}" is not a block in the current saved workflow.` + ) + } + if (block.enabled === false) { + throw new OrchestrationError( + 'validation', + `run.stopAfterBlockId "${blockId}" is disabled, so the run never executes it.` + ) + } + if (block.data?.parentId) { + throw new OrchestrationError( + 'validation', + `run.stopAfterBlockId "${blockId}" is inside loop or parallel "${block.data.parentId}"; stop after that container instead.` + ) + } + if (!reachableFrom(state, entryBlockId).has(blockId)) { + throw new OrchestrationError( + 'validation', + `run.stopAfterBlockId "${blockId}" is not reachable from entry block "${entryBlockId}"; stop after the entry block or one downstream of it.` + ) + } +} + function listTriggers(options: ReturnType): string { return options.map((option) => `${option.triggerBlockId} (${option.blockName})`).join(', ') } @@ -68,6 +132,7 @@ function executionServiceInput(params: { includeFileBase64: params.input.includeFileBase64, base64MaxBytes: params.input.base64MaxBytes, selectedOutputs: params.input.selectedOutputs, + stopAfterBlockId: params.input.stopAfterBlockId, rateLimitCounter: 'sync' as const, abortSignal: params.input.abortSignal, mode: params.input.mode, @@ -115,6 +180,7 @@ export const executeManualWorkflowOperation = defineAuthorizedWorkflowUseCase({ ) } + assertStopAfterBlock(state, input.stopAfterBlockId, selected.triggerBlockId) const executionInput = input.useMockPayload ? selected.mockPayload : input.input const validation = validateTriggerInput(selected, executionInput) if (!validation.ok) { @@ -141,6 +207,7 @@ export const executeManualWorkflowFromBlockOperation = defineAuthorizedWorkflowU `run.entry.blockId "${input.blockId}" is not a block in the current saved workflow.` ) } + assertStopAfterBlock(state, input.stopAfterBlockId, input.blockId) const sourceSnapshot = await getExecutionStateForWorkflow(input.sourceRunId, context.workflowId) if (!sourceSnapshot) { diff --git a/apps/sim/lib/workflows/executor/execute-service.ts b/apps/sim/lib/workflows/executor/execute-service.ts index d76c07f6595..dec05aac944 100644 --- a/apps/sim/lib/workflows/executor/execute-service.ts +++ b/apps/sim/lib/workflows/executor/execute-service.ts @@ -120,6 +120,8 @@ export interface ExecuteWorkflowServiceParams { /** Mocked upstream outputs (block name/id → output object) overlaid on the snapshot. */ variableInputs?: Record } + /** Saved block after which execution stops, validated by the application use case. */ + stopAfterBlockId?: string } export interface ExecuteWorkflowServiceFailure { @@ -272,6 +274,7 @@ export async function executeWorkflowService( useDraftState = false, triggerBlockId, runFromBlock, + stopAfterBlockId, } = params let reqLogger = logger.withMetadata({ requestId, workflowId, userId }) @@ -289,6 +292,9 @@ export async function executeWorkflowService( if (runFromBlock && !useDraftState) { throw new Error('Run-from-block requires manual execution state') } + if (stopAfterBlockId && !useDraftState) { + throw new Error('Stop-after-block requires manual execution state') + } if (callChain) { const chainError = validateCallChain(callChain) @@ -589,6 +595,7 @@ export async function executeWorkflowService( triggerBlockId, useDraftState, runFromBlock, + stopAfterBlockId, onStream, onBlockComplete: (blockId, data) => onBlockComplete(blockId, data.output, data.outputBlockId), @@ -689,6 +696,7 @@ export async function executeWorkflowService( base64MaxBytes, abortSignal: timeoutController.signal, runFromBlock, + stopAfterBlockId, }) await handlePostExecutionPauseState({ result, workflowId, executionId, loggingSession }) diff --git a/apps/sim/lib/workflows/executor/execution-core.ts b/apps/sim/lib/workflows/executor/execution-core.ts index 0f994a481f6..a5b10390c5d 100644 --- a/apps/sim/lib/workflows/executor/execution-core.ts +++ b/apps/sim/lib/workflows/executor/execution-core.ts @@ -907,6 +907,14 @@ async function executeWorkflowCoreImpl( resolvedStopAfterBlockId = buildLoopSentinelEndId(stopAfterBlockId) } else if (serializedWorkflow.parallels?.[stopAfterBlockId]) { resolvedStopAfterBlockId = buildParallelSentinelEndId(stopAfterBlockId) + } else if ( + !serializedWorkflow.blocks.some((block) => block.id === stopAfterBlockId && block.enabled) + ) { + // The engine stops on an exact node id and skips disabled blocks, so an absent or + // disabled target would run everything. + throw new Error( + `Stop block ${stopAfterBlockId} is not an enabled block in the workflow being executed` + ) } } diff --git a/apps/sim/package.json b/apps/sim/package.json index f537910f724..cbb1ef66c78 100644 --- a/apps/sim/package.json +++ b/apps/sim/package.json @@ -23,6 +23,7 @@ "test": "vitest run", "test:scim:e2e": "bun run scripts/test-scim-e2e.ts", "test:workflow-version-compare:e2e": "bun --no-env-file scripts/test-workflow-version-compare-e2e.ts", + "test:workflow-stop-after:e2e": "bun --no-env-file scripts/test-workflow-stop-after-e2e.ts", "test:watch": "vitest", "test:coverage": "vitest run --coverage", "email:dev": "email dev --dir components/emails", diff --git a/apps/sim/scripts/test-workflow-stop-after-e2e.ts b/apps/sim/scripts/test-workflow-stop-after-e2e.ts new file mode 100644 index 00000000000..691583fa980 --- /dev/null +++ b/apps/sim/scripts/test-workflow-stop-after-e2e.ts @@ -0,0 +1,410 @@ +import assert from 'node:assert/strict' +import { execFile } from 'node:child_process' +import { mkdtemp, rm, writeFile } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { resolve } from 'node:path' +import { fileURLToPath } from 'node:url' +import { promisify } from 'node:util' +import { assertDisposableTestDatabaseUrl } from '@sim/db/testing/test-infrastructure' +import { createLogger } from '@sim/logger' +import { sha256Hex } from '@sim/security/hash' +import { getErrorMessage } from '@sim/utils/errors' +import { sleep } from '@sim/utils/helpers' +import { generateId } from '@sim/utils/id' +import { isRecordLike } from '@sim/utils/object' +import { truncate } from '@sim/utils/string' +import postgres from 'postgres' +import { + type V2ExecuteWorkflowBody, + type V2ExecuteWorkflowData, + v2ExecuteWorkflowDataSchema, +} from '@/lib/api/contracts/v2/workflows' +import { readResponseTextWithLimit } from '@/lib/core/utils/stream-limits' + +/** + * Exercises `run.stopAfterBlockId` against a running local Next app, through the + * v2 execute route and the published CLI, on a disposable database. + * + * The fixture is `Start → Slow (wait) → Check (wait) → After (wait)`. Slow stands + * 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. + */ +const logger = createLogger('WorkflowStopAfterE2E') +const execFileAsync = promisify(execFile) +const MAX_RESPONSE_BYTES = 2 * 1024 * 1024 +/** The first execute request cold-compiles the route under `next dev`. */ +const REQUEST_TIMEOUT_MS = 300_000 +const SLOW_SECONDS = 4 +const SLOW_MS = SLOW_SECONDS * 1000 +const startedAt = new Date().toISOString() + +function requiredEnvironment(name: string): string { + const value = process.env[name] + assert(value, `${name} must be explicitly provided`) + return value +} + +const baseUrl = new URL(requiredEnvironment('STOP_AFTER_E2E_BASE_URL')) +const databaseUrl = assertDisposableTestDatabaseUrl( + requiredEnvironment('STOP_AFTER_E2E_DATABASE_URL') +) +const reportPath = requiredEnvironment('STOP_AFTER_E2E_REPORT_PATH') +assert(new Set(['localhost', '127.0.0.1', '[::1]']).has(baseUrl.hostname), 'Use a loopback app') +assert.equal(baseUrl.protocol, 'http:', 'Use a local HTTP app') +assert(!baseUrl.username && !baseUrl.password, 'App URL cannot contain credentials') +assert.equal(baseUrl.pathname, '/', 'App URL must be an origin') +assert(!baseUrl.search && !baseUrl.hash, 'App URL cannot contain a query or fragment') + +const sql = postgres(databaseUrl.toString(), { max: 2 }) +const ownerId = generateId() +const workspaceId = generateId() +const personalKey = `sk-sim-fixture-${generateId()}` +const cliPath = fileURLToPath(new URL('../../../packages/sim-cli/src/index.ts', import.meta.url)) +const checks: { name: string; status: 'passed' | 'failed'; durationMs: number; error?: string }[] = + [] +const requests: { method: string; path: string; status: number; durationMs: number }[] = [] +let directory: string | undefined + +interface PipelineFixture { + workflowId: string + start: string + slow: string + check: string + after: string +} + +const pipeline = fixtureIds() +const otherPipeline = fixtureIds() + +function fixtureIds(): PipelineFixture { + return { + workflowId: generateId(), + start: generateId(), + slow: generateId(), + check: generateId(), + after: generateId(), + } +} + +function failureMessage(error: unknown): string { + return truncate(getErrorMessage(error).replaceAll(personalKey, '[redacted]'), 2000) +} + +async function check(name: string, run: () => Promise) { + const started = performance.now() + try { + await run() + checks.push({ name, status: 'passed', durationMs: Math.round(performance.now() - started) }) + logger.info(`PASS ${name}`) + } catch (error) { + checks.push({ + name, + status: 'failed', + durationMs: Math.round(performance.now() - started), + error: failureMessage(error), + }) + throw error + } +} + +function record(value: unknown): Record { + assert(isRecordLike(value), 'Expected a JSON object') + return value +} + +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, + type: 'start_trigger', + 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) }, + ] + 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}, ${fixture.workflowId}, ${block.type}, ${block.name}, ${index * 300}, 0, ${JSON.stringify(block.subBlocks)}::text::jsonb)` + } + for (const [source, target] of [ + [fixture.start, fixture.slow], + [fixture.slow, fixture.check], + [fixture.check, fixture.after], + ]) { + await tx`insert into workflow_edges (id, workflow_id, source_block_id, target_block_id, source_handle, target_handle) + values (${generateId()}, ${fixture.workflowId}, ${source}, ${target}, 'source', 'target')` + } +} + +async function seed() { + directory = await mkdtemp(resolve(tmpdir(), 'sim-stop-after-')) + await sql.begin(async (tx) => { + const email = `${ownerId}@stop-after.test` + await tx`insert into "user" (id, name, email, normalized_email, email_verified, created_at, updated_at) + values (${ownerId}, 'Stop-after fixture', ${email}, ${email}, true, now(), now())` + await tx`insert into user_stats (id, user_id) values (${generateId()}, ${ownerId})` + await tx`insert into workspace (id, name, owner_id, billed_account_user_id) + values (${workspaceId}, 'Stop-after fixture', ${ownerId}, ${ownerId})` + await tx`insert into permissions (id, user_id, entity_type, entity_id, permission_type) + values (${generateId()}, ${ownerId}, 'workspace', ${workspaceId}, 'admin')` + await tx`insert into api_key (id, user_id, name, key, key_hash, type) + values (${generateId()}, ${ownerId}, 'Stop-after fixture', ${personalKey}, ${sha256Hex(personalKey)}, 'personal')` + await seedPipeline(tx, pipeline) + await seedPipeline(tx, otherPipeline) + }) +} + +async function execute( + workflowId: string, + body: V2ExecuteWorkflowBody, + expectedStatus = 200 +): Promise> { + const url = new URL(`/api/v2/workflows/${workflowId}/execute`, baseUrl) + const started = performance.now() + // boundary-raw-fetch: protocol E2E exercises a separately running local app over real HTTP + const response = await fetch(url, { + method: 'POST', + redirect: 'error', + signal: AbortSignal.timeout(REQUEST_TIMEOUT_MS), + headers: { + accept: 'application/json', + 'content-type': 'application/json', + 'x-api-key': personalKey, + 'x-forwarded-for': '127.0.0.1', + }, + body: JSON.stringify(body), + }) + requests.push({ + method: 'POST', + path: url.pathname, + status: response.status, + durationMs: Math.round(performance.now() - started), + }) + const text = await readResponseTextWithLimit(response, { + maxBytes: MAX_RESPONSE_BYTES, + label: 'Stop-after E2E response', + }) + assert.equal(response.status, expectedStatus, `${url.pathname}: ${truncate(text, 500)}`) + return record(JSON.parse(text)) +} + +async function run( + workflowId: string, + body: V2ExecuteWorkflowBody +): Promise { + const data = v2ExecuteWorkflowDataSchema.parse(record(await execute(workflowId, body)).data) + assert.equal(data.status, 'completed', `run ${data.runId} did not complete`) + return data +} + +async function expectBadRequest(workflowId: string, body: V2ExecuteWorkflowBody, code: string) { + const status = code === 'NOT_FOUND' ? 404 : 400 + const error = record((await execute(workflowId, body, status)).error) + assert.equal(error.code, code) +} + +async function runCli(args: string[]): Promise { + assert(directory, 'CLI fixture directory must exist') + const { stdout } = await execFileAsync( + 'bun', + [ + '--no-env-file', + cliPath, + '--endpoint', + baseUrl.origin, + '--workspace', + workspaceId, + '--output', + 'json', + 'workflows', + 'run', + ...args, + ], + { + cwd: directory, + env: { ...process.env, SIM_CONFIG_DIR: directory, SIM_API_KEY: personalKey, NO_COLOR: '1' }, + timeout: REQUEST_TIMEOUT_MS, + maxBuffer: MAX_RESPONSE_BYTES, + } + ) + return v2ExecuteWorkflowDataSchema.parse(JSON.parse(stdout)) +} + +const selectAll = ['Slow.status', 'Check.status', 'After.status'] + +try { + await check('seed disposable workspace, personal key and two pipelines', seed) + + let sourceRunId = '' + await check('a full manual run executes every block and persists its state', async () => { + const full = await run(pipeline.workflowId, { + run: { source: 'manual' }, + selectedOutputs: selectAll, + }) + assert.deepEqual(full.blockOutputs, { + 'Slow.status': 'completed', + 'Check.status': 'completed', + 'After.status': 'completed', + }) + assert( + (full.durationMs ?? 0) >= SLOW_MS, + `a full run must include Slow's ${SLOW_MS} ms, took ${full.durationMs} ms` + ) + sourceRunId = full.runId + }) + + await check( + 'the CLI re-runs only Check against the source run, without re-running Slow', + async () => { + const single = await runCli([ + pipeline.workflowId, + '--from-block', + pipeline.check, + '--source-run', + sourceRunId, + '--stop-after', + pipeline.check, + ...selectAll.flatMap((selector) => ['--select-output', selector]), + ]) + assert.equal(single.status, 'completed') + assert.notEqual(single.runId, sourceRunId) + assert.deepEqual(single.blockOutputs, { 'Check.status': 'completed' }) + assert( + (single.durationMs ?? Number.POSITIVE_INFINITY) < SLOW_MS, + `a single-block run must not wait on Slow, took ${single.durationMs} ms` + ) + } + ) + + await check('a trigger-entry run stops after the named block', async () => { + const until = await run(pipeline.workflowId, { + run: { source: 'manual', stopAfterBlockId: pipeline.slow }, + selectedOutputs: selectAll, + }) + assert.deepEqual(until.blockOutputs, { 'Slow.status': 'completed' }) + }) + + await check('a block entry without stopAfterBlockId still runs downstream blocks', async () => { + const fromCheck = await run(pipeline.workflowId, { + run: { + source: 'manual', + entry: { type: 'block', blockId: pipeline.check, sourceRunId }, + }, + selectedOutputs: selectAll, + }) + assert.deepEqual(fromCheck.blockOutputs, { + 'Check.status': 'completed', + 'After.status': 'completed', + }) + }) + + await check('invalid stop-after selections are refused before anything runs', async () => { + const before = + await sql`select count(*)::int as count from workflow_execution_logs where workflow_id = ${pipeline.workflowId}` + await expectBadRequest( + pipeline.workflowId, + { run: { source: 'manual', stopAfterBlockId: otherPipeline.check } }, + 'BAD_REQUEST' + ) + await expectBadRequest( + pipeline.workflowId, + { run: { source: 'manual', stopAfterBlockId: '' } }, + 'BAD_REQUEST' + ) + await expectBadRequest( + pipeline.workflowId, + { + run: { + source: 'manual', + entry: { type: 'block', blockId: pipeline.check, sourceRunId }, + stopAfterBlockId: pipeline.slow, + }, + }, + 'BAD_REQUEST' + ) + await expectBadRequest( + pipeline.workflowId, + { run: { source: 'manual', stopAfterBlockId: pipeline.check }, async: true }, + 'BAD_REQUEST' + ) + await expectBadRequest( + pipeline.workflowId, + { + run: { + source: 'manual', + entry: { type: 'block', blockId: pipeline.check, sourceRunId: 'not-a-run' }, + stopAfterBlockId: pipeline.check, + }, + }, + 'NOT_FOUND' + ) + const after = + await sql`select count(*)::int as count from workflow_execution_logs where workflow_id = ${pipeline.workflowId}` + assert.equal(after[0].count, before[0].count, 'a refused request must not start a run') + }) + + await check('a source run from another workflow is not accepted', async () => { + const foreign = await run(otherPipeline.workflowId, { run: { source: 'manual' } }) + await expectBadRequest( + pipeline.workflowId, + { + run: { + source: 'manual', + entry: { type: 'block', blockId: pipeline.check, sourceRunId: foreign.runId }, + stopAfterBlockId: pipeline.check, + }, + }, + 'NOT_FOUND' + ) + }) +} catch (error) { + logger.error(failureMessage(error)) + process.exitCode = 1 +} finally { + 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] + 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` + if (open === 0) break + await sleep(500) + } + for (let attempt = 1; ; attempt++) { + try { + await sql.begin(async (tx) => { + await tx`delete from workspace where id = ${workspaceId}` + await tx`delete from "user" where id = ${ownerId}` + }) + break + } catch (error) { + const deadlocked = isRecordLike(error) && error.code === '40P01' + if (!deadlocked || attempt === 5) throw error + await sleep(1000) + } + } + if (directory) await rm(directory, { recursive: true, force: true }) + }) + } catch (error) { + logger.error(failureMessage(error)) + process.exitCode = 1 + } finally { + await sql.end() + await writeFile( + reportPath, + `${JSON.stringify({ startedAt, finishedAt: new Date().toISOString(), status: process.exitCode ? 'failed' : 'passed', checks, requests }, null, 2)}\n` + ) + } +} diff --git a/packages/sim-cli/src/commands/protocol/workflow-run-follow.test.ts b/packages/sim-cli/src/commands/protocol/workflow-run-follow.test.ts index 9033de264cf..d0d54c207f2 100644 --- a/packages/sim-cli/src/commands/protocol/workflow-run-follow.test.ts +++ b/packages/sim-cli/src/commands/protocol/workflow-run-follow.test.ts @@ -203,6 +203,61 @@ describe('sim workflows run --follow', () => { }) }) + it('sends --stop-after with a block entry so one block re-runs against the source run', async () => { + vi.spyOn(console, 'log').mockImplementation(() => {}) + + await run( + WORKFLOW_ID, + '--from-block', + 'agent-1', + '--source-run', + 'run-1', + '--stop-after', + 'agent-1' + ) + + expect(requestRaw.mock.calls[0][1].body).toEqual({ + run: { + source: 'manual', + entry: { type: 'block', blockId: 'agent-1', sourceRunId: 'run-1' }, + stopAfterBlockId: 'agent-1', + }, + }) + }) + + it('lets --stop-after alone imply a manual run through the trigger', async () => { + vi.spyOn(console, 'log').mockImplementation(() => {}) + + await run(WORKFLOW_ID, '--stop-after', 'agent-1') + await run(WORKFLOW_ID, '--trigger', 'start', '--stop-after', 'agent-1') + + expect(requestRaw.mock.calls[0][1].body).toEqual({ + run: { source: 'manual', stopAfterBlockId: 'agent-1' }, + }) + expect(requestRaw.mock.calls[1][1].body).toEqual({ + run: { + source: 'manual', + entry: { type: 'trigger', blockId: 'start' }, + stopAfterBlockId: 'agent-1', + }, + }) + }) + + it('refuses an empty --stop-after rather than running the whole draft', async () => { + await expect(run(WORKFLOW_ID, '--stop-after', '')).rejects.toThrow( + '--stop-after requires a block ID' + ) + expect(requestRaw).not.toHaveBeenCalled() + }) + + it('refuses --stop-after with --async before sending anything', async () => { + await expect(run(WORKFLOW_ID, '--stop-after', 'agent-1', '--async')).rejects.toThrow( + 'Manual execution does not support --async' + ) + expect(request).not.toHaveBeenCalled() + expect(requestRaw).not.toHaveBeenCalled() + }) + it('prints then fails for a failed NDJSON run just like the JSON path', async () => { requestRaw.mockResolvedValue( ndjsonResponse({ 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 da5c9001848..00f8ca36807 100644 --- a/packages/sim-cli/src/commands/protocol/workflow-run-follow.ts +++ b/packages/sim-cli/src/commands/protocol/workflow-run-follow.ts @@ -43,16 +43,13 @@ export interface FollowOptions { stderr: CommentaryWriter } -type WorkflowRunSelection = - | { source: 'manual' } - | { - source: 'manual' - entry: { type: 'trigger'; blockId?: string; useMockPayload?: boolean } - } - | { - source: 'manual' - entry: { type: 'block'; blockId: string; sourceRunId: string } - } +type WorkflowRunSelection = { + source: 'manual' + entry?: + | { type: 'trigger'; blockId?: string; useMockPayload?: boolean } + | { type: 'block'; blockId: string; sourceRunId: string } + stopAfterBlockId?: string +} /** Projects friendly CLI flags into the API's strict nested run selector. */ export function resolveWorkflowRunSelection( @@ -60,8 +57,10 @@ export function resolveWorkflowRunSelection( ): WorkflowRunSelection | undefined { const trigger = typeof flags.trigger === 'string' ? flags.trigger : undefined const useMockPayload = flags.mockPayload === true - /** A trigger entry only exists on the draft, so these flags imply `--manual`. */ - const manual = flags.manual === true || trigger !== undefined || useMockPayload + const stopAfter = typeof flags.stopAfter === 'string' ? flags.stopAfter : undefined + /** A trigger entry and a stop block only exist on the draft, so these flags imply `--manual`. */ + const manual = + flags.manual === true || trigger !== undefined || useMockPayload || stopAfter !== undefined const fromBlock = typeof flags.fromBlock === 'string' ? flags.fromBlock : undefined const sourceRun = typeof flags.sourceRun === 'string' ? flags.sourceRun : undefined @@ -80,15 +79,20 @@ export function resolveWorkflowRunSelection( if (useMockPayload && flags.input !== undefined) { throw new SimApiError('--mock-payload cannot be combined with --input', 0) } + if (stopAfter !== undefined && stopAfter.trim() === '') { + throw new SimApiError('--stop-after requires a block ID', 0) + } + const stop = stopAfter !== undefined ? { stopAfterBlockId: stopAfter } : {} if (fromBlock && sourceRun) { return { source: 'manual', entry: { type: 'block', blockId: fromBlock, sourceRunId: sourceRun }, + ...stop, } } if (!manual) return undefined - if (!trigger && !useMockPayload) return { source: 'manual' } + if (!trigger && !useMockPayload) return { source: 'manual', ...stop } return { source: 'manual', entry: { @@ -96,6 +100,7 @@ export function resolveWorkflowRunSelection( ...(trigger ? { blockId: trigger } : {}), ...(useMockPayload ? { useMockPayload: true } : {}), }, + ...stop, } } @@ -488,6 +493,10 @@ export function attachWorkflowRunFollow(workflows: Command): void { '--source-run ', 'Prior run whose persisted state supplies upstream outputs (requires --from-block)' ) + .option( + '--stop-after ', + 'Stop the run after this saved block; with --from-block on the same block, re-runs only that block (implies --manual)' + ) .option( '--follow', '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' diff --git a/packages/sim-cli/src/generated/v2-api.ts b/packages/sim-cli/src/generated/v2-api.ts index c6a15f56e57..aad83e0a553 100644 --- a/packages/sim-cli/src/generated/v2-api.ts +++ b/packages/sim-cli/src/generated/v2-api.ts @@ -4893,6 +4893,7 @@ export type ExecuteWorkflowBody = { blockId: string sourceRunId: string } + stopAfterBlockId?: string } async?: boolean executionTimeoutSeconds?: number diff --git a/scripts/check-unused-exports.baseline.json b/scripts/check-unused-exports.baseline.json index f238aa01c30..013b0444c84 100644 --- a/scripts/check-unused-exports.baseline.json +++ b/scripts/check-unused-exports.baseline.json @@ -3808,8 +3808,6 @@ "apps/sim/lib/api/contracts/v2/workflows.ts#V2DeployedWebhook", "apps/sim/lib/api/contracts/v2/workflows.ts#V2DownloadRunFileParams", "apps/sim/lib/api/contracts/v2/workflows.ts#V2DuplicateWorkflowBody", - "apps/sim/lib/api/contracts/v2/workflows.ts#V2ExecuteWorkflowBody", - "apps/sim/lib/api/contracts/v2/workflows.ts#V2ExecuteWorkflowData", "apps/sim/lib/api/contracts/v2/workflows.ts#V2ExecuteWorkflowHeaders", "apps/sim/lib/api/contracts/v2/workflows.ts#V2ExecuteWorkflowQueued", "apps/sim/lib/api/contracts/v2/workflows.ts#V2ExecutionError",