Skip to content

Commit 572982d

Browse files
committed
Merge remote-tracking branch 'origin/staging' into fix/cli-run-latency
2 parents be5a21b + b2c5b76 commit 572982d

38 files changed

Lines changed: 1447 additions & 150 deletions

‎.github/workflows/test-build.yml‎

Lines changed: 16 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -253,27 +253,31 @@ jobs:
253253
VERSION_COMPARE_E2E_REPORT_PATH="$report_dir/version-compare-http-report.json" \
254254
bun run test:workflow-version-compare:e2e
255255
256-
# Self-hosted without billing: hosted billing admits runs through Redis, which
257-
# this job does not provision.
258-
- name: Verify embedded CLI workflow runs and log reads over real HTTP
256+
# A self-hosted app: hosted billing admits a run only through a Redis usage
257+
# reservation, and the SCIM suite above asserts PostgreSQL rate-limit storage,
258+
# so workflow execution gets its own app rather than adding Redis to that one.
259+
- name: Verify workflow runs, stop-after, and embedded CLI reads over real HTTP
259260
working-directory: apps/sim
260261
env:
261262
NEXT_PUBLIC_APP_URL: http://127.0.0.1:3018
262263
BETTER_AUTH_URL: http://127.0.0.1:3018
263264
NEXT_PUBLIC_FORCE_HOSTED: 'false'
264-
INTERNAL_API_SECRET: cli-http-ci-local-secret-at-least-32-characters
265+
INTERNAL_API_SECRET: workflow-http-ci-local-secret-at-least-32-characters
266+
DB_TX_TRIPWIRE: throw
265267
DISABLE_TELEMETRY: 'true'
266268
NEXT_TELEMETRY_DISABLED: '1'
269+
NEXT_PUBLIC_CHAT_DISABLED: 'true'
267270
READY_TIMEOUT_SECONDS: 300
268271
run: |
269272
report_dir="$RUNNER_TEMP/e2e"
270-
server_log="$report_dir/cli-next.log"
273+
server_log="$report_dir/workflow-next.log"
271274
mkdir -p "$report_dir"
272275
node ../../node_modules/next/dist/bin/next dev --hostname 127.0.0.1 --port 3018 > "$server_log" 2>&1 &
273276
server_pid=$!
274277
finish() {
275278
kill "$server_pid" 2>/dev/null || true
276279
wait "$server_pid" 2>/dev/null || true
280+
awk '/^ (GET|POST|PUT|PATCH|DELETE|HEAD) \/api\// { print }' "$server_log" > "$report_dir/workflow-http-status.log"
277281
}
278282
trap finish EXIT
279283
fail_startup() {
@@ -283,12 +287,16 @@ jobs:
283287
}
284288
started=$SECONDS
285289
until curl --fail --silent --max-time 10 http://127.0.0.1:3018/api/health > /dev/null; do
286-
kill -0 "$server_pid" 2>/dev/null || fail_startup 'Local CLI app exited during startup.'
290+
kill -0 "$server_pid" 2>/dev/null || fail_startup 'Local workflow app exited during startup.'
287291
[ $((SECONDS - started)) -lt "$READY_TIMEOUT_SECONDS" ] ||
288-
fail_startup "Local CLI app did not become ready within $READY_TIMEOUT_SECONDS seconds."
292+
fail_startup "Local workflow app did not become ready within $READY_TIMEOUT_SECONDS seconds."
289293
sleep 2
290294
done
291-
echo "Local CLI app ready after $((SECONDS - started))s"
295+
echo "Local workflow app ready after $((SECONDS - started))s"
296+
STOP_AFTER_E2E_BASE_URL="$NEXT_PUBLIC_APP_URL" \
297+
STOP_AFTER_E2E_DATABASE_URL="$DATABASE_URL" \
298+
STOP_AFTER_E2E_REPORT_PATH="$report_dir/stop-after-http-report.json" \
299+
bun run test:workflow-stop-after:e2e
292300
CLI_LATENCY_E2E_BASE_URL="$NEXT_PUBLIC_APP_URL" \
293301
CLI_LATENCY_E2E_DATABASE_URL="$DATABASE_URL" \
294302
CLI_LATENCY_E2E_RUNS=3 \

‎apps/docs/content/docs/cli/reference.mdx‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6704,6 +6704,7 @@ sim workflows run <workflowId> [options]
67046704
| `--mock-payload` | No | Use the selected trigger's server-derived mock payload; runs the current saved workflow state (implies --manual). |
67056705
| `--from-block <blockId>` | No | Run manually from this saved workflow block. |
67066706
| `--source-run <runId>` | No | Prior run whose persisted state supplies upstream outputs (requires --from-block). |
6707+
| `--stop-after <blockId>` | No | Stop the run after this saved block; with --from-block on the same block, re-runs only that block (implies --manual). |
67076708
| `--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. |
67086709
| `--include-thinking` | No | Show model reasoning while following (requires --follow). |
67096710
| `--include-tool-calls` | No | Show tool calls while following (requires --follow). |

‎apps/docs/content/docs/cli/workflows.mdx‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -642,6 +642,7 @@ sim workflows run <workflowId> [options]
642642
| `--mock-payload` | No | Use the selected trigger's server-derived mock payload; runs the current saved workflow state (implies --manual). |
643643
| `--from-block <blockId>` | No | Run manually from this saved workflow block. |
644644
| `--source-run <runId>` | No | Prior run whose persisted state supplies upstream outputs (requires --from-block). |
645+
| `--stop-after <blockId>` | No | Stop the run after this saved block; with --from-block on the same block, re-runs only that block (implies --manual). |
645646
| `--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. |
646647
| `--include-thinking` | No | Show model reasoning while following (requires --follow). |
647648
| `--include-tool-calls` | No | Show tool calls while following (requires --follow). |

‎apps/docs/openapi-v2-workflows.json‎

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -12615,6 +12615,11 @@
1261512615
"additionalProperties": false
1261612616
}
1261712617
]
12618+
},
12619+
"stopAfterBlockId": {
12620+
"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.",
12621+
"type": "string",
12622+
"minLength": 1
1261812623
}
1261912624
},
1262012625
"required": ["source"],
@@ -12703,6 +12708,17 @@
1270312708
"sourceRunId": "run_123"
1270412709
}
1270512710
}
12711+
},
12712+
{
12713+
"run": {
12714+
"source": "manual",
12715+
"entry": {
12716+
"type": "block",
12717+
"blockId": "block_123",
12718+
"sourceRunId": "run_123"
12719+
},
12720+
"stopAfterBlockId": "block_123"
12721+
}
1270612722
}
1270712723
]
1270812724
},

‎apps/sim/app/api/v2/workflows/[workflowId]/execute/route.ts‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -448,6 +448,7 @@ export const POST = withRouteHandler(
448448
mode: body.stream ? 'stream' : resultStream ? 'sync-result-stream' : 'sync',
449449
blockId: manualRun.entry.blockId,
450450
sourceRunId: manualRun.entry.sourceRunId,
451+
stopAfterBlockId: manualRun.stopAfterBlockId,
451452
},
452453
request: req,
453454
})
@@ -460,6 +461,7 @@ export const POST = withRouteHandler(
460461
mode: body.stream ? 'stream' : resultStream ? 'sync-result-stream' : 'sync',
461462
triggerBlockId: manualRun.entry?.blockId,
462463
useMockPayload: manualRun.entry?.useMockPayload === true,
464+
stopAfterBlockId: manualRun.stopAfterBlockId,
463465
},
464466
request: req,
465467
})

‎apps/sim/executor/execution/block-executor.ts‎

Lines changed: 4 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -424,9 +424,8 @@ export class BlockExecutor {
424424
typeof normalizedOutput._childWorkflowInstanceId === 'string'
425425
? normalizedOutput._childWorkflowInstanceId
426426
: undefined
427-
const displayOutput = filterOutputForLog(block.metadata?.id || '', normalizedOutput, {
428-
block,
429-
})
427+
// Shallow top-level copy: streaming later writes token/cost keys onto blockLog.output.
428+
const displayOutput = { ...blockLog.output }
430429
const displayInput = this.projectInputsForDisplay(inputsForLog, block, inputDisplayRegistry)
431430
blockLog.input = displayInput
432431
const displayProvenance = settledBlockRegistry?.exportCommittedProvenanceForValue({
@@ -680,7 +679,7 @@ export class BlockExecutor {
680679

681680
if (!isSentinel && blockLog) {
682681
const displayInput = this.projectInputsForDisplay(input, block, inputDisplayRegistry)
683-
const displayOutput = filterOutputForLog(block.metadata?.id || '', softOutput, { block })
682+
const displayOutput = { ...blockLog.output }
684683
const displayProvenance =
685684
ctx.resolvedSecretTraceRegistry?.exportCommittedProvenanceForValue({
686685
input: displayInput,
@@ -798,7 +797,7 @@ export class BlockExecutor {
798797
const childWorkflowInstanceId = ChildWorkflowError.isChildWorkflowError(error)
799798
? error.childWorkflowInstanceId
800799
: undefined
801-
const displayOutput = filterOutputForLog(block.metadata?.id || '', errorOutput, { block })
800+
const displayOutput = { ...blockLog.output }
802801
const displayInput = this.projectInputsForDisplay(input, block, inputDisplayRegistry)
803802
const displayProvenance = errorRegistry?.exportCommittedProvenanceForValue({
804803
input: displayInput,

‎apps/sim/executor/execution/engine.ts‎

Lines changed: 21 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,10 @@ import { subscribeToExecutionCancellation } from '@/lib/execution/cancellation'
55
import { BlockType, EDGE } from '@/executor/constants'
66
import type { DAG } from '@/executor/dag/builder'
77
import type { EdgeManager } from '@/executor/execution/edge-manager'
8-
import { serializePauseSnapshot } from '@/executor/execution/snapshot-serializer'
8+
import {
9+
buildCompletedExecutionState,
10+
serializePauseSnapshot,
11+
} from '@/executor/execution/snapshot-serializer'
912
import type { SerializableExecutionState } from '@/executor/execution/types'
1013
import type { NodeExecutionOrchestrator } from '@/executor/orchestrators/node'
1114
import type {
@@ -147,7 +150,7 @@ export class ExecutionEngine {
147150
success: true,
148151
output: this.finalOutput,
149152
logs: this.context.blockLogs,
150-
executionState: this.getSerializableExecutionState(),
153+
executionState: this.getCompletedExecutionState(),
151154
metadata: this.context.metadata,
152155
}
153156
} catch (error) {
@@ -591,6 +594,22 @@ export class ExecutionEngine {
591594
}
592595
}
593596

597+
/**
598+
* State for a run whose blocks have all settled, without the JSON round-trip.
599+
* Cancelled and failed runs keep {@link getSerializableExecutionState}: a block
600+
* still running there could mutate the logs after the run returns.
601+
*/
602+
private getCompletedExecutionState(): SerializableExecutionState | undefined {
603+
try {
604+
return buildCompletedExecutionState(this.context, this.dag, this.edgeManager)
605+
} catch (error) {
606+
this.execLogger.warn('Failed to serialize execution state', {
607+
error: toError(error).message,
608+
})
609+
return undefined
610+
}
611+
}
612+
594613
private collectPauseResponses(): NormalizedBlockOutput {
595614
const responses = Array.from(this.pausedBlocks.values()).map((pause) => pause.response)
596615

‎apps/sim/executor/execution/snapshot-serializer.test.ts‎

Lines changed: 84 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,11 @@
11
import { describe, expect, it, vi } from 'vitest'
22
import type { DAG, DAGNode } from '@/executor/dag/builder'
33
import { EdgeManager } from '@/executor/execution/edge-manager'
4-
import { serializePauseSnapshot } from '@/executor/execution/snapshot-serializer'
4+
import {
5+
buildCompletedExecutionState,
6+
isLiveExecutionState,
7+
serializePauseSnapshot,
8+
} from '@/executor/execution/snapshot-serializer'
59
import type { ExecutionContext } from '@/executor/types'
610
import { ResolvedSecretTraceRegistry } from '@/executor/utils/resolved-secret-trace-registry'
711

@@ -307,3 +311,82 @@ describe('serializePauseSnapshot', () => {
307311
expect(serialized.metadata.capabilityGovernedUserId).toBe('requesting-member')
308312
})
309313
})
314+
315+
describe('buildCompletedExecutionState', () => {
316+
/** The state a completed run used to carry: the pause snapshot's state, JSON round-tripped. */
317+
function jsonClonedState(context: ExecutionContext): unknown {
318+
return JSON.parse(serializePauseSnapshot(context, []).snapshot).state
319+
}
320+
321+
function contextWithOutput(output: unknown): ExecutionContext {
322+
return createContext({
323+
blockStates: new Map([['block-1', { output, executed: true, executionTime: 5 }]]),
324+
executedBlocks: new Set(['block-1']),
325+
blockLogs: [
326+
{
327+
blockId: 'block-1',
328+
blockName: 'Block',
329+
blockType: 'function',
330+
startedAt: '2026-01-01T00:00:00.000Z',
331+
endedAt: '2026-01-01T00:00:01.000Z',
332+
durationMs: 1000,
333+
success: true,
334+
executionOrder: 1,
335+
output,
336+
input: { note: undefined, list: [undefined, 1] },
337+
},
338+
] as ExecutionContext['blockLogs'],
339+
})
340+
}
341+
342+
const shared = { rows: [{ id: 1, at: new Date(0) }] }
343+
const cyclic: Record<string, unknown> = { a: 1 }
344+
cyclic.self = cyclic
345+
const cycleBehindToJSON: Record<string, unknown> = { toJSON: () => ({ safe: true }) }
346+
cycleBehindToJSON.self = cycleBehindToJSON
347+
348+
it.each([
349+
['a subtree shared twice', { first: shared, second: shared }],
350+
['a cycle hidden behind toJSON', { value: cycleBehindToJSON }],
351+
[
352+
'a boxed number carrying a BigInt property',
353+
{ value: Object.assign(Object(1), { big: BigInt(1) }) },
354+
],
355+
])('serializes like the JSON-cloned pause state for %s', (_name, output) => {
356+
const context = contextWithOutput(output)
357+
expect(JSON.stringify(buildCompletedExecutionState(context))).toBe(
358+
JSON.stringify(jsonClonedState(context))
359+
)
360+
})
361+
362+
it.each([
363+
['a cycle', cyclic],
364+
['a BigInt', { big: BigInt(1) }],
365+
])('throws like the pause snapshot for %s', (_name, output) => {
366+
const context = contextWithOutput(output)
367+
expect(() => jsonClonedState(context)).toThrow(TypeError)
368+
expect(() => buildCompletedExecutionState(context)).toThrow(TypeError)
369+
})
370+
371+
it('shares block outputs with the run instead of cloning them', () => {
372+
const output = { rows: [{ id: 1 }] }
373+
const context = contextWithOutput(output)
374+
const state = buildCompletedExecutionState(context)
375+
376+
expect(state.blockLogs[0].output).toBe(output)
377+
expect(state.blockStates['block-1'].output).toBe(output)
378+
379+
context.blockLogs[0].endedAt = 'later'
380+
expect(state.blockLogs[0].endedAt).toBe('2026-01-01T00:00:01.000Z')
381+
})
382+
383+
it('marks completed state as live through spreads but not through JSON', () => {
384+
const context = contextWithOutput({ rows: [{ id: 1 }] })
385+
const state = buildCompletedExecutionState(context)
386+
387+
expect(isLiveExecutionState(state)).toBe(true)
388+
expect(isLiveExecutionState({ ...state, sourceExecutionId: 'other' })).toBe(true)
389+
expect(isLiveExecutionState(jsonClonedState(context))).toBe(false)
390+
expect(JSON.stringify(state)).toBe(JSON.stringify(jsonClonedState(context)))
391+
})
392+
})

0 commit comments

Comments
 (0)