diff --git a/.github/workflows/test-build.yml b/.github/workflows/test-build.yml index e61fb6e9f7d..0fc7d4b9027 100644 --- a/.github/workflows/test-build.yml +++ b/.github/workflows/test-build.yml @@ -269,6 +269,49 @@ jobs: VERSION_COMPARE_E2E_REPORT_PATH="$report_dir/version-compare-http-report.json" \ bun run test:workflow-version-compare:e2e + # Self-hosted without billing: hosted billing admits runs through Redis, which + # this job does not provision. + - name: Verify embedded CLI workflow runs and log reads 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: cli-http-ci-local-secret-at-least-32-characters + DISABLE_TELEMETRY: 'true' + NEXT_TELEMETRY_DISABLED: '1' + READY_TIMEOUT_SECONDS: 300 + run: | + report_dir="$RUNNER_TEMP/e2e" + server_log="$report_dir/cli-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 + } + 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 CLI app exited during startup.' + [ $((SECONDS - started)) -lt "$READY_TIMEOUT_SECONDS" ] || + fail_startup "Local CLI app did not become ready within $READY_TIMEOUT_SECONDS seconds." + sleep 2 + done + echo "Local CLI app ready after $((SECONDS - started))s" + CLI_LATENCY_E2E_BASE_URL="$NEXT_PUBLIC_APP_URL" \ + CLI_LATENCY_E2E_DATABASE_URL="$DATABASE_URL" \ + CLI_LATENCY_E2E_RUNS=3 \ + CLI_LATENCY_E2E_WARMUP=1 \ + CLI_LATENCY_E2E_REPORT_PATH="$report_dir/cli-run-latency-report.json" \ + bun run test:cli-run-latency: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. diff --git a/apps/docs/content/docs/cli/logs.mdx b/apps/docs/content/docs/cli/logs.mdx index 22403eec7a9..159e4a4909c 100644 --- a/apps/docs/content/docs/cli/logs.mdx +++ b/apps/docs/content/docs/cli/logs.mdx @@ -31,7 +31,7 @@ sim logs get [options] | Option | Required | Description | | --- | --- | --- | -| `--include-workflow-state` | No | Include the saved workflow snapshot (default: true). Set false to omit block configuration from a log read. Other run fields are unchanged. | +| `--include-workflow-state` | No | Include the saved workflow snapshot (default: true; Sim’s in-app agent leaves it out unless this flag is passed). Set false to omit block configuration from a log read. Other run fields are unchanged. | | `--no-include-workflow-state` | No | Send --include-workflow-state as false. | | `--trace` | No | Show expanded trace spans with inputs, outputs, errors, timing, and cost. | diff --git a/apps/docs/content/docs/cli/reference.mdx b/apps/docs/content/docs/cli/reference.mdx index 9655f81a3ad..d56a61a68cf 100644 --- a/apps/docs/content/docs/cli/reference.mdx +++ b/apps/docs/content/docs/cli/reference.mdx @@ -2751,7 +2751,7 @@ sim logs get [options] | Option | Required | Description | | --- | --- | --- | -| `--include-workflow-state` | No | Include the saved workflow snapshot (default: true). Set false to omit block configuration from a log read. Other run fields are unchanged. | +| `--include-workflow-state` | No | Include the saved workflow snapshot (default: true; Sim’s in-app agent leaves it out unless this flag is passed). Set false to omit block configuration from a log read. Other run fields are unchanged. | | `--no-include-workflow-state` | No | Send --include-workflow-state as false. | | `--trace` | No | Show expanded trace spans with inputs, outputs, errors, timing, and cost. | @@ -6695,7 +6695,7 @@ sim workflows run [options] | `--async` | No | Queue the run and return immediately. | | `--execution-timeout-seconds ` | No | Maximum duration of an asynchronous run, in seconds, capped by the plan's execution timeout. Requires `async: true`; otherwise returns `400`. | | `--select-output ` | No | Return blockName.path values (e.g. agent_1.content), or childWorkflowId.blockName.path for a child workflow (applies to every invocation) — in blockOutputs on a sync run, or from the streamed result with --follow; missing paths are omitted. Not available with --async (space-separated, or @path / @- with one value per line; @@value for a literal leading @). | -| `--include-file-base64` | No | Inline eligible output files as base64 content. Rejected when `async` is true. | +| `--include-file-base64` | No | Inline eligible output files as base64 content (default: true; Sim’s in-app agent gets file references only unless this flag is passed). Rejected when `async` is true. | | `--no-include-file-base64` | No | Send --include-file-base64 as false. | | `--base64-max-bytes ` | No | Maximum total bytes of file content to inline as base64, lowering but never raising the server limit of 16 MiB. Rejected when `async` is true. | | `--run-id ` | No | One-shot identifier for this run; NOT an idempotency key — reusing a claimed value fails with RUN_ID_CONFLICT instead of replaying the first result, and a fresh value starts another run. | diff --git a/apps/docs/content/docs/cli/workflows.mdx b/apps/docs/content/docs/cli/workflows.mdx index 4d03ddb8c55..019d9691255 100644 --- a/apps/docs/content/docs/cli/workflows.mdx +++ b/apps/docs/content/docs/cli/workflows.mdx @@ -633,7 +633,7 @@ sim workflows run [options] | `--async` | No | Queue the run and return immediately. | | `--execution-timeout-seconds ` | No | Maximum duration of an asynchronous run, in seconds, capped by the plan's execution timeout. Requires `async: true`; otherwise returns `400`. | | `--select-output ` | No | Return blockName.path values (e.g. agent_1.content), or childWorkflowId.blockName.path for a child workflow (applies to every invocation) — in blockOutputs on a sync run, or from the streamed result with --follow; missing paths are omitted. Not available with --async (space-separated, or @path / @- with one value per line; @@value for a literal leading @). | -| `--include-file-base64` | No | Inline eligible output files as base64 content. Rejected when `async` is true. | +| `--include-file-base64` | No | Inline eligible output files as base64 content (default: true; Sim’s in-app agent gets file references only unless this flag is passed). Rejected when `async` is true. | | `--no-include-file-base64` | No | Send --include-file-base64 as false. | | `--base64-max-bytes ` | No | Maximum total bytes of file content to inline as base64, lowering but never raising the server limit of 16 MiB. Rejected when `async` is true. | | `--run-id ` | No | One-shot identifier for this run; NOT an idempotency key — reusing a claimed value fails with RUN_ID_CONFLICT instead of replaying the first result, and a fresh value starts another run. | diff --git a/apps/sim/lib/workflows/application/execute-manual-workflow.ts b/apps/sim/lib/workflows/application/execute-manual-workflow.ts index 4d7049e7161..d949bea2f5b 100644 --- a/apps/sim/lib/workflows/application/execute-manual-workflow.ts +++ b/apps/sim/lib/workflows/application/execute-manual-workflow.ts @@ -10,7 +10,10 @@ import { executeWorkflowService, } from '@/lib/workflows/executor/execute-service' import { getExecutionStateForWorkflow } from '@/lib/workflows/executor/execution-state' -import { loadWorkflowFromNormalizedTables } from '@/lib/workflows/persistence/utils' +import { + loadWorkflowFromNormalizedTables, + type NormalizedWorkflowData, +} from '@/lib/workflows/persistence/utils' import { resolveTriggerRunOptions, validateTriggerInput, @@ -119,6 +122,7 @@ function executionServiceInput(params: { principal: WorkflowExecutionPrincipal context: Awaited> input: ManualExecutionInput + draftState: NormalizedWorkflowData }) { return { workflowId: params.context.workflowId, @@ -141,6 +145,7 @@ function executionServiceInput(params: { includeToolCalls: params.input.includeToolCalls, triggerType: 'manual' as const, useDraftState: true, + draftState: params.draftState, } } @@ -189,7 +194,7 @@ export const executeManualWorkflowOperation = defineAuthorizedWorkflowUseCase({ } return executeWorkflowService({ - ...executionServiceInput({ principal, context, input }), + ...executionServiceInput({ principal, context, input, draftState: state }), input: executionInput, triggerBlockId: selected.triggerBlockId, }) @@ -218,7 +223,7 @@ export const executeManualWorkflowFromBlockOperation = defineAuthorizedWorkflowU } return executeWorkflowService({ - ...executionServiceInput({ principal, context, input }), + ...executionServiceInput({ principal, context, input, draftState: state }), input: input.input, runFromBlock: { startBlockId: input.blockId, diff --git a/apps/sim/lib/workflows/executor/execute-service.ts b/apps/sim/lib/workflows/executor/execute-service.ts index dec05aac944..7cb9707f635 100644 --- a/apps/sim/lib/workflows/executor/execute-service.ts +++ b/apps/sim/lib/workflows/executor/execute-service.ts @@ -30,6 +30,7 @@ import { loadDeployedWorkflowState, loadWorkflowDeploymentVersionState, loadWorkflowFromNormalizedTables, + type NormalizedWorkflowData, } from '@/lib/workflows/persistence/utils' import { shouldEmitAgentStreamEvents } from '@/lib/workflows/streaming/agent-stream-protocol' import { resolveOutputSelectors } from '@/lib/workflows/streaming/resolve-output-selectors' @@ -110,6 +111,12 @@ export interface ExecuteWorkflowServiceParams { includeToolCalls?: boolean /** Execute the current saved state manually instead of the active deployment. */ useDraftState?: boolean + /** + * The saved state the application use case loaded to choose and validate the entry + * point. Requires `useDraftState`; the run executes this same snapshot instead of + * reading the draft tables again. + */ + draftState?: NormalizedWorkflowData /** Explicit trigger entry point selected and validated by the application use case. */ triggerBlockId?: string /** Trusted prior-run snapshot resolved by the application use case. */ @@ -272,6 +279,7 @@ export async function executeWorkflowService( includeThinking = false, includeToolCalls = false, useDraftState = false, + draftState, triggerBlockId, runFromBlock, stopAfterBlockId, @@ -295,6 +303,9 @@ export async function executeWorkflowService( if (stopAfterBlockId && !useDraftState) { throw new Error('Stop-after-block requires manual execution state') } + if (draftState && !useDraftState) { + throw new Error('A preloaded draft state requires manual execution state') + } if (callChain) { const chainError = validateCallChain(callChain) @@ -467,7 +478,7 @@ export async function executeWorkflowService( let workflowBlocks: Record = {} try { const workflowData = useDraftState - ? await loadWorkflowFromNormalizedTables(workflowId) + ? (draftState ?? (await loadWorkflowFromNormalizedTables(workflowId))) : deploymentVersionId ? await loadWorkflowDeploymentVersionState(workflowId, deploymentVersionId, workspaceId) : await loadDeployedWorkflowState(workflowId, workspaceId) @@ -594,6 +605,7 @@ export async function executeWorkflowService( workflowTriggerType: triggerType, triggerBlockId, useDraftState, + draftState, runFromBlock, stopAfterBlockId, onStream, @@ -697,6 +709,7 @@ export async function executeWorkflowService( abortSignal: timeoutController.signal, runFromBlock, stopAfterBlockId, + draftState, }) await handlePostExecutionPauseState({ result, workflowId, executionId, loggingSession }) diff --git a/apps/sim/lib/workflows/executor/execute-workflow.ts b/apps/sim/lib/workflows/executor/execute-workflow.ts index fc6332679ac..01c14e76eb8 100644 --- a/apps/sim/lib/workflows/executor/execute-workflow.ts +++ b/apps/sim/lib/workflows/executor/execute-workflow.ts @@ -11,6 +11,7 @@ import { LoggingSession } from '@/lib/logs/execution/logging-session' import { captureServerEvent } from '@/lib/posthog/server' import { executeWorkflowCore } from '@/lib/workflows/executor/execution-core' import { handlePostExecutionPauseState } from '@/lib/workflows/executor/pause-persistence' +import type { NormalizedWorkflowData } from '@/lib/workflows/persistence/utils' import { ExecutionSnapshot } from '@/executor/execution/snapshot' import type { BlockCompletionCallbackData, @@ -54,6 +55,8 @@ export interface ExecuteWorkflowOptions { abortSignal?: AbortSignal /** Use the live/draft workflow state instead of the deployed state. Used by copilot. */ useDraftState?: boolean + /** Draft state the caller already loaded, reused instead of reading the draft tables again. */ + draftState?: NormalizedWorkflowData /** Immutable workflow state selected by a trusted server-side trigger boundary. */ workflowStateOverride?: NonNullable /** Stop execution after this block completes. Used for "run until block" feature. */ @@ -226,6 +229,7 @@ export async function executeWorkflow( trustedInitialResolvedSecretTraceProvenance: streamConfig?.trustedInitialResolvedSecretTraceProvenance, runFromBlock: streamConfig?.runFromBlock, + draftState: streamConfig?.draftState, })) const blockTypes = [ diff --git a/apps/sim/lib/workflows/executor/execution-core.ts b/apps/sim/lib/workflows/executor/execution-core.ts index a5b10390c5d..d7c05f420f9 100644 --- a/apps/sim/lib/workflows/executor/execution-core.ts +++ b/apps/sim/lib/workflows/executor/execution-core.ts @@ -43,6 +43,7 @@ import { loadDeployedWorkflowState, loadWorkflowDeploymentVersionState, loadWorkflowFromNormalizedTables, + type NormalizedWorkflowData, } from '@/lib/workflows/persistence/utils' import { TriggerUtils } from '@/lib/workflows/triggers/triggers' import { updateWorkflowRunCounts } from '@/lib/workflows/utils' @@ -143,6 +144,12 @@ export interface ExecuteWorkflowCoreOptions { trustedInitialResolvedSecretTraceProvenance?: ResolvedSecretTraceProvenanceV1 /** Immutable deployment admitted by the durable parent log for a resumed execution. */ resumeDeploymentVersionId?: string + /** + * Draft state the caller already loaded for this draft-state run. Reused instead of + * reading the draft tables again, so the graph that executes is the snapshot the + * caller validated its trigger and output selectors against. + */ + draftState?: NormalizedWorkflowData /** * Environment the caller already resolved for this run, reused instead of loading * and decrypting it again. Used only when it was resolved for exactly the @@ -627,6 +634,7 @@ async function executeWorkflowCoreImpl( stopAfterBlockId, runFromBlock, resumeDeploymentVersionId, + draftState, } = options loggingSession.setExecutionDeadlineAt(getExecutionDeadlineAt(abortSignal)) const { metadata, input, workflowVariables, selectedOutputs } = snapshot @@ -721,7 +729,7 @@ async function executeWorkflowCoreImpl( } if (useDraftState) { - const draftData = await loadWorkflowFromNormalizedTables(workflowId) + const draftData = draftState ?? (await loadWorkflowFromNormalizedTables(workflowId)) if (!draftData) { throw new Error('Workflow not found or not yet saved') diff --git a/apps/sim/package.json b/apps/sim/package.json index dffa24627c4..feee817b989 100644 --- a/apps/sim/package.json +++ b/apps/sim/package.json @@ -25,6 +25,7 @@ "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:desktop-inbox:e2e": "bun --no-env-file scripts/test-desktop-inbox-e2e.ts", + "test:cli-run-latency:e2e": "bun --no-env-file scripts/test-cli-run-latency-e2e.ts", "test:watch": "vitest", "test:coverage": "vitest run --coverage", "email:dev": "email dev --dir components/emails", diff --git a/apps/sim/scripts/test-cli-run-latency-e2e.ts b/apps/sim/scripts/test-cli-run-latency-e2e.ts new file mode 100644 index 00000000000..0c26f73fe11 --- /dev/null +++ b/apps/sim/scripts/test-cli-run-latency-e2e.ts @@ -0,0 +1,476 @@ +import assert from 'node:assert/strict' +import { mkdir, mkdtemp, rm, writeFile } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { dirname, resolve } from 'node:path' +import { fileURLToPath, pathToFileURL } from 'node:url' +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 { EmbeddedCliIdentity, EmbeddedCliResult } from 'sim/embed' + +/** + * Times the two CLI commands the chat agent issues most, `workflows run --manual` + * and `logs get --trace`, against a separately running local app over real HTTP. + * + * Each command runs through the embedded CLI (the same command tree, request + * builder and JSON renderer the agent's `sim_cli` tool uses) and its stdout is + * written to a file the way an `outputFile` sink lands it. Every iteration is + * split into segments so a before/after comparison shows where time moved: + * CLI work before the request, server time before execution started, execution, + * server time after execution ended, CLI rendering, and the sink write. + * + * Fixtures are seeded straight into a disposable database and removed at the end. + * With this checkout's CLI, the run also asserts the embedded results stay lean: + * file references without inline bytes, and a log without its workflow snapshot. + * `CLI_LATENCY_E2E_CLI_MODULE` points the run at another checkout's `embed.ts`, + * so two CLI builds can be compared against the same server (measured, not asserted), and + * `CLI_LATENCY_E2E_SINK_DIR` keeps every command's stdout for inspection. + */ +const logger = createLogger('CliRunLatencyE2E') +const REQUEST_TIMEOUT_MS = 120_000 +const LOG_SETTLE_TIMEOUT_MS = 30_000 +const startedAt = new Date().toISOString() + +function requiredEnvironment(name: string): string { + const value = process.env[name] + assert(value, `${name} must be explicitly provided`) + return value +} + +function positiveInteger(name: string, fallback: number): number { + const raw = process.env[name] + if (raw === undefined) return fallback + const value = Number(raw) + assert(Number.isInteger(value) && value > 0, `${name} must be a positive integer`) + return value +} + +const baseUrl = new URL(requiredEnvironment('CLI_LATENCY_E2E_BASE_URL')) +const databaseUrl = assertDisposableTestDatabaseUrl( + requiredEnvironment('CLI_LATENCY_E2E_DATABASE_URL') +) +const reportPath = requiredEnvironment('CLI_LATENCY_E2E_REPORT_PATH') +const label = process.env.CLI_LATENCY_E2E_LABEL ?? 'local' +const measuredRuns = positiveInteger('CLI_LATENCY_E2E_RUNS', 30) +const warmupRuns = positiveInteger('CLI_LATENCY_E2E_WARMUP', 3) +const comparedCliModule = process.env.CLI_LATENCY_E2E_CLI_MODULE +const cliModulePath = + comparedCliModule ?? + fileURLToPath(new URL('../../../packages/sim-cli/src/embed.ts', import.meta.url)) +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.equal(baseUrl.pathname, '/', 'App URL must be an origin') + +const sql = postgres(databaseUrl.toString(), { max: 2 }) +const userId = generateId() +const workspaceId = generateId() +const workflowId = generateId() +const apiKey = `sk-sim-fixture-${generateId()}` + +type Segment = + | 'cliBeforeRequestMs' + | 'serverBeforeExecutionMs' + | 'executionMs' + | 'serverAfterExecutionMs' + | 'httpMs' + | 'cliRenderMs' + | 'sinkWriteMs' + | 'totalMs' + +interface CommandTiming { + command: 'workflows run' | 'logs get' + iteration: number + segments: Partial> + stdoutBytes: number + /** What the command's stdout carried, so a size change is attributable. */ + carried: { inlineFileBytes?: boolean; workflowState?: boolean } +} + +const checks: { name: string; status: 'passed' | 'failed'; durationMs: number; error?: string }[] = + [] +const timings: CommandTiming[] = [] + +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: truncate(getErrorMessage(error).replaceAll(apiKey, '[redacted]'), 2000), + }) + throw error + } +} + +function record(value: unknown): Record { + assert(isRecordLike(value), 'Expected a JSON object') + return value +} + +const FUNCTION_CODE = `const rows = Array.from({ length: 2000 }, (_, i) => ({ i, label: 'row-' + i, value: i * 7 })) +return { rows, total: rows.reduce((sum, row) => sum + row.value, 0) }` + +const SUMMARY_CODE = `const rows = .rows +return { count: rows.length, csv: rows.map((row) => row.i + ',' + row.label + ',' + row.value).join('\\n') }` + +const blocks = [ + { + id: generateId(), + type: 'start_trigger', + name: 'Start', + x: 0, + subBlocks: { inputFormat: { id: 'inputFormat', type: 'input-format', value: [] } }, + }, + { + id: generateId(), + type: 'function', + name: 'Build', + x: 300, + subBlocks: { + language: { id: 'language', type: 'dropdown', value: 'javascript' }, + code: { id: 'code', type: 'code', value: FUNCTION_CODE }, + }, + }, + { + id: generateId(), + type: 'function', + name: 'Summarize', + x: 600, + subBlocks: { + language: { id: 'language', type: 'dropdown', value: 'javascript' }, + code: { id: 'code', type: 'code', value: SUMMARY_CODE }, + }, + }, + { + id: generateId(), + type: 'file_v5', + name: 'Report', + x: 900, + subBlocks: { + operation: { id: 'operation', type: 'dropdown', value: 'file_write' }, + fileName: { id: 'fileName', type: 'short-input', value: 'latency-report.csv' }, + content: { id: 'content', type: 'long-input', value: '' }, + overwrite: { id: 'overwrite', type: 'switch', value: true }, + }, + }, + { + id: generateId(), + type: 'file_v5', + name: 'Read Report', + x: 1200, + advancedMode: true, + subBlocks: { + operation: { id: 'operation', type: 'dropdown', value: 'file_read' }, + readFileId: { id: 'readFileId', type: 'short-input', value: '' }, + }, + }, +] + +async function seed() { + await sql.begin(async (tx) => { + const email = `${userId}@cli-latency.test` + await tx`insert into "user" (id, name, email, normalized_email, email_verified, created_at, updated_at) + values (${userId}, 'CLI latency fixture', ${email}, ${email}, true, now(), now())` + await tx`insert into user_stats (id, user_id) values (${generateId()}, ${userId})` + await tx`insert into workspace (id, name, owner_id, billed_account_user_id) + values (${workspaceId}, 'CLI latency fixture', ${userId}, ${userId})` + await tx`insert into permissions (id, user_id, entity_type, entity_id, permission_type) + values (${generateId()}, ${userId}, 'workspace', ${workspaceId}, 'admin')` + await tx`insert into api_key (id, user_id, name, key, key_hash, type) + values (${generateId()}, ${userId}, 'CLI latency fixture', ${apiKey}, ${sha256Hex(apiKey)}, 'personal')` + await tx`insert into workflow (id, user_id, workspace_id, name, last_synced, created_at, updated_at) + values (${workflowId}, ${userId}, ${workspaceId}, 'CLI latency fixture', now(), now(), now())` + for (const block of blocks) { + await tx`insert into workflow_blocks (id, workflow_id, type, name, position_x, position_y, advanced_mode, sub_blocks) + values (${block.id}, ${workflowId}, ${block.type}, ${block.name}, ${block.x}, 0, ${'advancedMode' in block}, ${JSON.stringify(block.subBlocks)}::text::jsonb)` + } + for (let index = 1; index < blocks.length; index++) { + await tx`insert into workflow_edges (id, workflow_id, source_block_id, target_block_id, source_handle, target_handle) + values (${generateId()}, ${workflowId}, ${blocks[index - 1].id}, ${blocks[index].id}, 'source', 'target')` + } + }) +} + +async function cleanup() { + // Snapshots outlive their workflow (`workflow_id` is set null), so they go first, + // after the logs that reference them. + await sql`delete from workflow_execution_logs where workflow_id = ${workflowId}` + await sql`delete from workflow_execution_snapshots where workflow_id = ${workflowId}` + await sql`delete from workspace where id = ${workspaceId}` + await sql`delete from "user" where id = ${userId}` +} + +interface TransportMarks { + firstRequestAt?: number + firstRequestEpochMs?: number + lastBodyAt?: number + lastBodyEpochMs?: number +} + +/** Real HTTP to the loopback app, recording when the first request left and the last body byte arrived. */ +function timingTransport(marks: TransportMarks): typeof fetch { + const transport = async (input: string | URL | Request, init?: RequestInit) => { + const url = new URL(input instanceof Request ? input.url : input) + assert.equal(url.origin, baseUrl.origin, 'Requests must stay on the loopback app') + if (marks.firstRequestAt === undefined) { + marks.firstRequestAt = performance.now() + marks.firstRequestEpochMs = Date.now() + } + const signal = init?.signal ?? (input instanceof Request ? input.signal : undefined) + // boundary-raw-fetch: protocol E2E exercises a separately running local app over real HTTP + const response = await fetch(input, { + ...init, + redirect: 'error', + signal: signal + ? AbortSignal.any([signal, AbortSignal.timeout(REQUEST_TIMEOUT_MS)]) + : AbortSignal.timeout(REQUEST_TIMEOUT_MS), + }) + const markBodyEnd = () => { + marks.lastBodyAt = performance.now() + marks.lastBodyEpochMs = Date.now() + } + if (!response.body) { + markBodyEnd() + return response + } + return new Response( + // The NDJSON reader returns on its `final` frame without reading to EOF, so + // the last chunk to arrive marks the end of the response as the CLI sees it. + response.body.pipeThrough( + new TransformStream({ + transform(chunk, controller) { + markBodyEnd() + controller.enqueue(chunk) + }, + flush: markBodyEnd, + }) + ), + { status: response.status, statusText: response.statusText, headers: response.headers } + ) + } + return transport as typeof fetch +} + +type RunEmbeddedCli = (argv: string[], identity: EmbeddedCliIdentity) => Promise + +async function timedCommand( + runCli: RunEmbeddedCli, + argv: string[], + sinkPath: string +): Promise<{ + result: EmbeddedCliResult + marks: TransportMarks + segments: Partial> +}> { + const marks: TransportMarks = {} + const started = performance.now() + const result = await runCli(argv, { + endpoint: baseUrl.origin, + apiKey, + workspaceId, + transport: timingTransport(marks), + }) + const rendered = performance.now() + assert.equal( + result.exitCode, + 0, + `${argv.slice(0, 2).join(' ')} failed: ${truncate(result.stderr, 500)}` + ) + await writeFile(sinkPath, result.stdout) + const finished = performance.now() + assert( + marks.firstRequestAt !== undefined && marks.lastBodyAt !== undefined, + 'No request was made' + ) + return { + result, + marks, + segments: { + cliBeforeRequestMs: marks.firstRequestAt - started, + httpMs: marks.lastBodyAt - marks.firstRequestAt, + cliRenderMs: rendered - marks.lastBodyAt, + sinkWriteMs: finished - rendered, + totalMs: finished - started, + }, + } +} + +/** The log row is finalized after the response, so a read waits for its terminal state first. */ +async function waitForSettledLog(runId: string) { + const deadline = Date.now() + LOG_SETTLE_TIMEOUT_MS + while (Date.now() < deadline) { + // The columns are UTC wall-clock without a zone; epoch extraction reads them as UTC. + const [row] = await sql<{ status: string; startedAtMs: string; endedAtMs: string | null }[]>` + select status, + extract(epoch from started_at) * 1000 as "startedAtMs", + extract(epoch from ended_at) * 1000 as "endedAtMs" + from workflow_execution_logs where execution_id = ${runId}` + if (row?.endedAtMs && row.status !== 'running') + return { startedAtMs: Number(row.startedAtMs), endedAtMs: Number(row.endedAtMs) } + await sleep(50) + } + throw new Error(`Run ${runId} log did not settle within ${LOG_SETTLE_TIMEOUT_MS} ms`) +} + +async function iterate( + runCli: RunEmbeddedCli, + directory: string, + iteration: number, + keep: boolean +) { + const run = await timedCommand( + runCli, + ['workflows', 'run', workflowId, '--manual', '--select-output', 'Summarize.result.count'], + resolve(directory, `run-${iteration}.json`) + ) + const payload = record(JSON.parse(run.result.stdout)) + assert.equal(payload.status, 'completed', 'The fixture run must complete') + const runId = payload.runId + assert(typeof runId === 'string', 'The run result names its run') + assert.deepEqual(payload.blockOutputs, { 'Summarize.result.count': 2000 }) + const file = record(record(payload.output).file) + assert( + typeof file.id === 'string' && file.name === 'latency-report.csv', + 'The run names its file' + ) + const inlineFileBytes = typeof file.base64 === 'string' + // Another CLI build is measured as it behaves; this checkout's CLI must stay lean. + if (!comparedCliModule) assert(!inlineFileBytes, 'An embedded run returns file references only') + + const logRow = await waitForSettledLog(runId) + const logs = await timedCommand( + runCli, + ['logs', 'get', runId, '--trace'], + resolve(directory, `log-${iteration}.json`) + ) + const log = record(JSON.parse(logs.result.stdout)) + assert.equal(log.runId, runId) + assert(Array.isArray(log.traceSpans) && log.traceSpans.length > 0, 'The log carries trace spans') + const workflowState = isRecordLike(log.workflowState) + if (!comparedCliModule) assert(!workflowState, 'An embedded log read omits the workflow snapshot') + + if (!keep) return + const { firstRequestEpochMs, lastBodyEpochMs } = run.marks + assert(firstRequestEpochMs !== undefined && lastBodyEpochMs !== undefined) + timings.push({ + command: 'workflows run', + iteration, + stdoutBytes: Buffer.byteLength(run.result.stdout), + carried: { inlineFileBytes }, + segments: { + ...run.segments, + serverBeforeExecutionMs: logRow.startedAtMs - firstRequestEpochMs, + executionMs: logRow.endedAtMs - logRow.startedAtMs, + serverAfterExecutionMs: lastBodyEpochMs - logRow.endedAtMs, + }, + }) + timings.push({ + command: 'logs get', + iteration, + stdoutBytes: Buffer.byteLength(logs.result.stdout), + carried: { workflowState }, + segments: logs.segments, + }) +} + +function percentile(sorted: number[], p: number): number { + const rank = (sorted.length - 1) * p + const low = Math.floor(rank) + const high = Math.ceil(rank) + return sorted[low] + (sorted[high] - sorted[low]) * (rank - low) +} + +function summarize() { + const summary: Record> = {} + for (const command of ['workflows run', 'logs get'] as const) { + const rows = timings.filter((timing) => timing.command === command) + const segments = new Set(rows.flatMap((row) => Object.keys(row.segments))) + summary[command] = {} + for (const segment of segments) { + const values = rows + .map((row) => row.segments[segment as Segment]) + .filter((value): value is number => value !== undefined) + .sort((a, b) => a - b) + summary[command][segment] = { + n: values.length, + p50: Math.round(percentile(values, 0.5) * 10) / 10, + p90: Math.round(percentile(values, 0.9) * 10) / 10, + } + } + const bytes = rows.map((row) => row.stdoutBytes).sort((a, b) => a - b) + summary[command].stdoutBytes = { + n: bytes.length, + p50: percentile(bytes, 0.5), + p90: percentile(bytes, 0.9), + } + } + return summary +} + +async function main() { + const keptSinkDirectory = process.env.CLI_LATENCY_E2E_SINK_DIR + const directory = keptSinkDirectory ?? (await mkdtemp(resolve(tmpdir(), 'sim-cli-latency-'))) + await mkdir(directory, { recursive: true }) + let failure: unknown + try { + const cli = (await import(pathToFileURL(cliModulePath).href)) as { + runEmbeddedCli: RunEmbeddedCli + } + await check('seed fixtures', seed) + await check(`warm up (${warmupRuns} iterations, not recorded)`, async () => { + for (let index = 0; index < warmupRuns; index++) + await iterate(cli.runEmbeddedCli, directory, -1 - index, false) + }) + await check(`measure (${measuredRuns} iterations)`, async () => { + for (let index = 0; index < measuredRuns; index++) + await iterate(cli.runEmbeddedCli, directory, index, true) + }) + } catch (error) { + failure = error + logger.error('CLI latency E2E failed', { + error: getErrorMessage(error).replaceAll(apiKey, '[redacted]'), + }) + } finally { + await cleanup().catch((error) => + logger.warn('Fixture cleanup failed', { error: getErrorMessage(error) }) + ) + await sql.end() + if (!keptSinkDirectory) await rm(directory, { recursive: true, force: true }) + } + + await mkdir(dirname(reportPath), { recursive: true }) + await writeFile( + reportPath, + `${JSON.stringify( + { + suite: 'cli-run-latency', + label, + startedAt, + finishedAt: new Date().toISOString(), + baseUrl: baseUrl.origin, + cliModule: cliModulePath, + warmupRuns, + measuredRuns, + checks, + summary: summarize(), + timings, + }, + null, + 2 + )}\n` + ) + if (failure) process.exit(1) +} + +await main() 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 307b47ad01a..1a94ceb8cd2 100644 --- a/packages/sim-cli/src/commands/protocol/workflow-run-follow.ts +++ b/packages/sim-cli/src/commands/protocol/workflow-run-follow.ts @@ -4,6 +4,7 @@ import { writeStderr } from '#sim-cli/output/io' import { styles } from '#sim-cli/output/presentation' import { clientFrom } from '../../context' import { CLI_CONTRACT } from '../../contract/commands' +import { embedStore } from '../../embed-context' import { V2_OPERATIONS } from '../../generated/v2-api' import { SimApiError } from '../../http/client' import { readNdjson } from '../../http/ndjson' @@ -151,6 +152,20 @@ async function readWorkflowResult(response: Response): Promise | undefined +): Record | undefined { + if (!embedStore.getStore() || body?.includeFileBase64 !== undefined) return body + return { ...body, includeFileBase64: false } +} + /** Runs synchronously while keeping idle-limited HTTP paths active. */ async function runWithResultStream(workflowId: string, command: Command): Promise { const flags = command.optsWithGlobals() as Record @@ -163,7 +178,7 @@ async function runWithResultStream(workflowId: string, command: Command): Promis const response = await client.requestRaw(request.path, { method: operation.method, query: request.query, - body: request.body, + body: withEmbeddedFileReferences(request.body), headers: { ...request.headers, accept: WORKFLOW_RESULT_STREAM_CONTENT_TYPE }, }) const payload = await readWorkflowResult(response) @@ -362,7 +377,7 @@ async function followRun(workflowId: string, command: Command): Promise { method: 'POST', query: request.query, body: { - ...(request.body ?? {}), + ...withEmbeddedFileReferences(request.body), stream: true, ...(includeThinking ? { includeThinking: true } : {}), ...(includeToolCalls ? { includeToolCalls: true } : {}), diff --git a/packages/sim-cli/src/contract/commands.ts b/packages/sim-cli/src/contract/commands.ts index 5525f471aec..62dc0258cdd 100644 --- a/packages/sim-cli/src/contract/commands.ts +++ b/packages/sim-cli/src/contract/commands.ts @@ -419,6 +419,16 @@ export const CLI_CONTRACT: CliContract = { getLog: { describe: 'Show run diagnostics', expandedTrace: true, + flags: { + // The snapshot repeats the block configuration `workflows get` already + // serves, and is the bulk of a log after its trace. An agent diagnosing a + // run reads the trace, so the agent's invocations leave it out unless asked. + includeWorkflowState: { + embeddedRequestDefault: 'false', + describe: + 'Include the saved workflow snapshot (default: true; Sim’s in-app agent leaves it out unless this flag is passed). Set false to omit block configuration from a log read. Other run fields are unchanged.', + }, + }, fields: [ { header: 'run', path: 'runId' }, { header: 'workflow', path: 'workflow.name' }, @@ -1926,6 +1936,12 @@ export const CLI_CONTRACT: CliContract = { stream: { omit: true }, includeThinking: { omit: true }, includeToolCalls: { omit: true }, + // Its embedded default lives with the synchronous request in + // `workflow-run-follow.ts`, because an `--async` run rejects the field. + includeFileBase64: { + describe: + 'Inline eligible output files as base64 content (default: true; Sim’s in-app agent gets file references only unless this flag is passed). Rejected when `async` is true.', + }, // Exposed under its domain name: every other flag in the CLI is one, and // `--x-run-id` would be the only place the raw HTTP header spelling // surfaced. The describe denies idempotency outright because the name diff --git a/packages/sim-cli/src/contract/types.ts b/packages/sim-cli/src/contract/types.ts index d62d546cfb5..ed8f8aa951d 100644 --- a/packages/sim-cli/src/contract/types.ts +++ b/packages/sim-cli/src/contract/types.ts @@ -95,6 +95,16 @@ export interface FlagSpec { * including a deliberate `--details basic`. */ requestDefault?: string + /** + * Value an embedded invocation sends when the caller passes nothing; wins over + * `requestDefault` there and is ignored by the installed CLI. + * + * The embedding host returns stdout to a model or writes it to the model's + * workbench, which addresses files and run state by id rather than reading + * bulk payloads inline. A field whose server default inlines such a payload + * names the leaner value here. Whatever the caller types still wins. + */ + embeddedRequestDefault?: string /** Accepted values when the generated descriptor cannot recover an enum. */ choices?: readonly string[] /** diff --git a/packages/sim-cli/src/embed.test.ts b/packages/sim-cli/src/embed.test.ts index 1fa113f3a63..9f6f72bbe80 100644 --- a/packages/sim-cli/src/embed.test.ts +++ b/packages/sim-cli/src/embed.test.ts @@ -143,3 +143,49 @@ describe('embedded artifact destinations', () => { expect(noWriter.stderr).toContain('no machine to write to') }) }) + +describe('embedded request defaults', () => { + function capture(body: unknown) { + const requests: { url: URL; body: Record | undefined }[] = [] + const transport = async (input: string | URL | Request, init?: RequestInit) => { + requests.push({ + url: new URL(String(input)), + body: typeof init?.body === 'string' ? JSON.parse(init.body) : undefined, + }) + return jsonResponse({ data: body }) + } + return { requests, identity: { ...IDENTITY, transport } } + } + + it('asks a synchronous run for file references unless the caller wants inline bytes', async () => { + const { requests, identity } = capture({ runId: 'run-1', status: 'completed', output: null }) + expect((await runEmbeddedCli(['workflows', 'run', 'wf', '--manual'], identity)).exitCode).toBe( + 0 + ) + expect( + (await runEmbeddedCli(['workflows', 'run', 'wf', '--include-file-base64'], identity)).exitCode + ).toBe(0) + expect(requests.map((request) => request.body?.includeFileBase64)).toEqual([false, true]) + }) + + it('asks a followed run for file references too', async () => { + const { requests, identity } = capture({}) + await runEmbeddedCli(['workflows', 'run', 'wf', '--follow'], identity) + expect(requests[0]?.body).toMatchObject({ stream: true, includeFileBase64: false }) + }) + + it('never sends the field on an async run, which the server rejects', async () => { + const { requests, identity } = capture({ runId: 'run-1', statusUrl: 'https://x.test/s' }) + expect((await runEmbeddedCli(['workflows', 'run', 'wf', '--async'], identity)).exitCode).toBe(0) + expect(requests[0]?.body).not.toHaveProperty('includeFileBase64') + }) + + it('reads a log without its workflow snapshot unless the caller asks for it', async () => { + const { requests, identity } = capture({ runId: 'run-1', traceSpans: [] }) + await runEmbeddedCli(['logs', 'get', 'run-1'], identity) + await runEmbeddedCli(['logs', 'get', 'run-1', '--include-workflow-state'], identity) + expect(requests.map((request) => request.url.searchParams.get('includeWorkflowState'))).toEqual( + ['false', 'true'] + ) + }) +}) diff --git a/packages/sim-cli/src/runtime/request.ts b/packages/sim-cli/src/runtime/request.ts index 30f1a086c79..4d07ea27b8c 100644 --- a/packages/sim-cli/src/runtime/request.ts +++ b/packages/sim-cli/src/runtime/request.ts @@ -643,7 +643,10 @@ export async function buildRequest( // A contract default only applies to what the caller left unsaid, so // typing the flag — including typing the server's own default back — still // decides. It is validated like any other value, enum choices included. - const raw = provided ?? flag.requestDefault + const raw = + provided ?? + (embedStore.getStore() ? flag.embeddedRequestDefault : undefined) ?? + flag.requestDefault /** * A blank filter is a mistake, and every v2 JSON route says so