From 27ff10c6351951df6d5459ebd85885fecb76826c Mon Sep 17 00:00:00 2001 From: Miya Date: Tue, 15 Sep 2026 19:28:02 +0200 Subject: [PATCH 1/6] fix: run standalone authored bodies with native Node promise hooks Session-Id: 01a09c40-ce3b-7f11-a7df-b6b7ccab6fd9 --- .github/workflows/cloud-runtime-artifact.yml | 6 +- .github/workflows/publish.yml | 6 +- docs/SURFACE.md | 23 +++ packages/sdk/src/authored-flow-executor.ts | 10 ++ packages/sdk/src/authored-node-entry.ts | 63 ++++++++ packages/sdk/src/authored-node-runner.ts | 139 ++++++++++++++++++ packages/sdk/src/authored-root.ts | 37 ++--- .../sdk/src/authored-runtime-capability.ts | 31 ++++ packages/sdk/src/authored-source-authority.ts | 33 +++++ packages/sdk/src/cli/direct-run.ts | 1 + packages/sdk/src/cli/run.ts | 4 + .../sdk/tests/authored-node-runtime.test.ts | 134 +++++++++++++++++ scripts/build-standalone-cli.mjs | 30 ++++ 13 files changed, 486 insertions(+), 31 deletions(-) create mode 100644 packages/sdk/src/authored-node-entry.ts create mode 100644 packages/sdk/src/authored-node-runner.ts create mode 100644 packages/sdk/src/authored-runtime-capability.ts create mode 100644 packages/sdk/src/authored-source-authority.ts create mode 100644 packages/sdk/tests/authored-node-runtime.test.ts create mode 100644 scripts/build-standalone-cli.mjs diff --git a/.github/workflows/cloud-runtime-artifact.yml b/.github/workflows/cloud-runtime-artifact.yml index f4a43bb3b..1110896db 100644 --- a/.github/workflows/cloud-runtime-artifact.yml +++ b/.github/workflows/cloud-runtime-artifact.yml @@ -22,6 +22,7 @@ on: - "kernel/**" - "packages/sdk/**" - "scripts/cloud-artifact.mjs" + - "scripts/build-standalone-cli.mjs" - "scripts/cloud-artifact.test.mjs" - "testdata/**" @@ -159,10 +160,7 @@ jobs: - name: Build standalone flows CLI run: | mkdir -p dist/cloud-artifact-input - bun build packages/sdk/src/cli-executable.ts \ - --compile \ - --target=bun-linux-x64 \ - --outfile=dist/cloud-artifact-input/flows + node scripts/build-standalone-cli.mjs bun-linux-x64 dist/cloud-artifact-input/flows - name: Assemble artifact and smoke verifier path env: diff --git a/.github/workflows/publish.yml b/.github/workflows/publish.yml index 5d3b6a654..b517bae74 100644 --- a/.github/workflows/publish.yml +++ b/.github/workflows/publish.yml @@ -145,8 +145,7 @@ jobs: run: | mkdir -p packages/runtime-linux-x64/bin cp kernel/target/release/relayflowd packages/runtime-linux-x64/bin/relayflowd - bun build packages/sdk/src/cli-executable.ts --compile --target=bun-linux-x64 \ - --outfile=packages/runtime-linux-x64/bin/flows + node scripts/build-standalone-cli.mjs bun-linux-x64 packages/runtime-linux-x64/bin/flows chmod +x packages/runtime-linux-x64/bin/relayflowd packages/runtime-linux-x64/bin/flows packages/runtime-linux-x64/bin/relayflowd --help packages/runtime-linux-x64/bin/flows check --json testdata/hello-deterministic.flow.yaml @@ -229,8 +228,7 @@ jobs: run: | mkdir -p packages/runtime-darwin-arm64/bin cp kernel/target/release/relayflowd packages/runtime-darwin-arm64/bin/relayflowd - bun build packages/sdk/src/cli-executable.ts --compile --target=bun-darwin-arm64 \ - --outfile=packages/runtime-darwin-arm64/bin/flows + node scripts/build-standalone-cli.mjs bun-darwin-arm64 packages/runtime-darwin-arm64/bin/flows chmod +x packages/runtime-darwin-arm64/bin/relayflowd packages/runtime-darwin-arm64/bin/flows packages/runtime-darwin-arm64/bin/relayflowd --help packages/runtime-darwin-arm64/bin/flows check --json testdata/hello-deterministic.flow.yaml diff --git a/docs/SURFACE.md b/docs/SURFACE.md index 68aa8539e..c5a5e0ee5 100644 --- a/docs/SURFACE.md +++ b/docs/SURFACE.md @@ -802,3 +802,26 @@ Any bypass of that journaling makes the run non-replayable and must set This is the required contract for channel rollout. Detecting bypasses and carrying that marker through the completion protocol remain implementation work; the initial post proof does not claim to enforce uninstrumented agent I/O. + + +### Standalone authored runtime + +The standalone CLI embeds the Node authored runner in the same hashed executable. +Authored `.flow.ts` bodies require Node **22.14 or newer** on `PATH`, or an +absolute `FLOWS_AUTHORED_NODE` executable path. No runtime is downloaded. The +runner probes Node version and native promise-hook capability before root +admission/body effects; a missing, old, or incompatible runtime is refused. +Direct SDK authored execution likewise requires working native promise hooks. +Bun 1.4.0 exposes no-op `async_hooks`, so it cannot prove that a native `await` +consumed a step. Treating every `.then` call as an await would incorrectly accept +ignored operations; that verification remains unchanged inside Node. + +Only the authored body runs in the Node child. The standalone CLI retains the +root worker lease, agent workers, declarative execution, daemon lookup, and +resume protocol. The child inherits the working directory/environment and uses +the same journal socket and child admission identities. It verifies the root's +pinned source graph and Surface package before executing the body. Root aborts +and signals stop the child; loss of the parent pipe stops it immediately. +Successful root output records the Node version, executable SHA256, and embedded +payload SHA256 as `executionRuntime`. The payload remains part of the existing +artifact hash, and completed child effects remain journal results on resume. diff --git a/packages/sdk/src/authored-flow-executor.ts b/packages/sdk/src/authored-flow-executor.ts index 289cce5be..bd50ffdf7 100644 --- a/packages/sdk/src/authored-flow-executor.ts +++ b/packages/sdk/src/authored-flow-executor.ts @@ -1,3 +1,4 @@ +import { assertAuthoredPromiseHooks } from './authored-runtime-capability.js'; import { pluginHelpers } from './plugin-loader.js'; import { runPluginEffect } from './authored-plugin-effect.js'; import { randomUUID } from 'node:crypto'; @@ -71,7 +72,15 @@ export interface AuthoredFlowJournalStep { readonly completionReason: ProtocolCompletionReason; } +export interface AuthoredExecutionRuntime { + readonly kind: 'node'; + readonly version: string; + readonly executableSha256: string; + readonly payloadSha256: string; +} + export interface AuthoredFlowExecutionResult { + readonly executionRuntime?: AuthoredExecutionRuntime; readonly rootRunId?: string; readonly name: string; readonly completionReason: FlowCompletionReason; @@ -133,6 +142,7 @@ export async function executeAuthoredFlow( input?: Input, options: ExecuteAuthoredFlowOptions = {}, ): Promise { + assertAuthoredPromiseHooks(); const getDefinition = options.getDefinition ?? getAuthoredFlowDefinition; const localAgentStream = options.localAgentStream; const onProgress = options.onProgress; diff --git a/packages/sdk/src/authored-node-entry.ts b/packages/sdk/src/authored-node-entry.ts new file mode 100644 index 000000000..dce9e2a8b --- /dev/null +++ b/packages/sdk/src/authored-node-entry.ts @@ -0,0 +1,63 @@ +import { createHash } from 'node:crypto'; +import { readFileSync, writeSync } from 'node:fs'; +import { JournalClient } from './journal-client.js'; +import { executeAuthoredFlow } from './authored-flow-executor.js'; +import { loadPinnedAuthoredSource } from './authored-source-authority.js'; +import { assertAuthoredNodeVersion } from './authored-runtime-capability.js'; +import { AuthoredFlowExecutionError } from './authored-flow-error.js'; +import type { AuthoredRootMetadata } from './authored-root.js'; + +const send = (message: unknown): void => { writeSync(3, JSON.stringify(message) + '\n'); }; +const hash = (path: string): string => createHash('sha256').update(readFileSync(path)).digest('hex'); +const controller = new AbortController(); +let finished = false; +const abort = (): void => { + controller.abort(); + setTimeout(() => process.exit(1), 2000).unref(); +}; +process.on('SIGINT', abort); process.on('SIGTERM', abort); +// The parent owns the root lease. Do not continue authored effects after it dies. +process.stdin.on('end', () => { if (!finished) process.exit(1); }); +process.stdin.on('error', () => process.exit(1)); +let client: JournalClient | undefined; +try { + assertAuthoredNodeVersion(); + send({ type: 'ready', runtime: { kind: 'node', version: process.versions.node, + executableSha256: hash(process.execPath), payloadSha256: hash(process.argv[1]!) } }); + const request = await new Promise<{ metadata: AuthoredRootMetadata; socketPath: string; + rootRunId: string; dataDir: string; localAgentStream?: string }>((resolve, reject) => { + let buffer = ''; + process.stdin.setEncoding('utf8'); + const onData = (chunk: string): void => { + buffer += chunk; + if (Buffer.byteLength(buffer) > 2 * 1024 * 1024) { reject(new Error('authored runtime request exceeded limit')); return; } + const end = buffer.indexOf('\n'); + if (end < 0) return; + process.stdin.off('data', onData); + try { resolve(JSON.parse(buffer.slice(0, end))); } catch { reject(new Error('invalid authored runtime request')); } + }; + process.stdin.on('data', onData); + }); + controller.signal.throwIfAborted(); + const loaded = await loadPinnedAuthoredSource(request.metadata); + if (request.localAgentStream !== request.metadata.localAgentStream) throw new Error('authored root local agent surface mismatch'); + client = new JournalClient(request.socketPath); + await client.connect(); await client.hello('flows-authored-node'); + const result = await executeAuthoredFlow(loaded.handle, client, + request.metadata.inputPresent ? request.metadata.input : undefined, { + getDefinition: loaded.getDefinition, dataDir: request.dataDir, + flowPath: request.metadata.flowPath, rootRunId: request.rootRunId, + localAgentStream: request.localAgentStream, signal: controller.signal, + onProgress: event => send({ type: 'progress', event }), + onWait: event => send({ type: 'wait', event }), + }); + send({ type: 'result', result }); +} catch (error) { + send({ type: 'error', message: error instanceof Error ? error.message : 'authored body failed', + ...(error instanceof AuthoredFlowExecutionError ? { code: error.code, + completionReason: error.completionReason, runId: error.runId } : {}) }); + process.exitCode = 1; +} finally { + finished = true; client?.close(); process.stdin.destroy(); + process.off('SIGINT', abort); process.off('SIGTERM', abort); +} diff --git a/packages/sdk/src/authored-node-runner.ts b/packages/sdk/src/authored-node-runner.ts new file mode 100644 index 000000000..0fc542b46 --- /dev/null +++ b/packages/sdk/src/authored-node-runner.ts @@ -0,0 +1,139 @@ +import { spawn, spawnSync } from 'node:child_process'; +import { createHash } from 'node:crypto'; +import { mkdtemp, readFile, rm, writeFile } from 'node:fs/promises'; +import { readFileSync, realpathSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import { isAbsolute, join } from 'node:path'; +import type { Readable } from 'node:stream'; +import type { AuthoredRootMetadata } from './authored-root.js'; +import type { AuthoredExecutionRuntime, AuthoredFlowExecutionResult, ExecuteAuthoredFlowOptions } from './authored-flow-executor.js'; +import { AuthoredFlowExecutionError, type AuthoredFlowExecutionErrorCode } from './authored-flow-error.js'; +import { assertAuthoredPromiseHooks } from './authored-runtime-capability.js'; + +let embeddedSource: string | undefined; +/** Installed by the standalone build; never fetched or resolved from a workspace. */ +export function installAuthoredNodeSource(source: string): void { + if (embeddedSource !== undefined || source.length === 0) throw new Error('invalid authored runtime payload'); + embeddedSource = source; +} +const hash = (value: Uint8Array | string): string => createHash('sha256').update(value).digest('hex'); +const refusal = (): AuthoredFlowExecutionError => new AuthoredFlowExecutionError( + 'unsupported_promise_lifecycle', + 'authored execution requires the embedded Node runner and Node >=22.14 on PATH (or an absolute FLOWS_AUTHORED_NODE); no workflow body was executed', +); +interface NodeAuthority { path: string; version: string; executableSha256: string } +function nodeAuthority(): NodeAuthority { + if (embeddedSource === undefined) throw refusal(); + const override = process.env['FLOWS_AUTHORED_NODE']; + if (override !== undefined && !isAbsolute(override)) throw refusal(); + const probe = spawnSync(override ?? 'node', ['--input-type=module', '--eval', ` +import {createHook} from 'node:async_hooks'; +const [major,minor]=process.versions.node.split('.').map(Number); +let init=false,resolve=false; +const h=createHook({init(_id,type){if(type==='PROMISE')init=true},promiseResolve(){resolve=true}}).enable(); +new Promise(r=>r());h.disable(); +if(process.versions.bun || major<22 || (major===22 && minor<14) || !init || !resolve) process.exit(2); +process.stdout.write(JSON.stringify({path:process.execPath,version:process.versions.node})); +`], { encoding: 'utf8', timeout: 10_000, maxBuffer: 4096, env: process.env }); + if (probe.error || probe.status !== 0) throw refusal(); + try { + const value = JSON.parse(probe.stdout) as { path: string; version: string }; + if (!isAbsolute(value.path) || !/^\d+\.\d+\.\d+$/.test(value.version)) throw refusal(); + const [major, minor] = value.version.split('.').map(Number); + if (major! < 22 || (major === 22 && minor! < 14)) throw refusal(); + const path = realpathSync(value.path); + return { path, version: value.version, executableSha256: hash(readFileSync(path)) }; + } catch { throw refusal(); } +} + +/** Called before root admission/resume; declarative commands never reach it. */ +export function assertAuthoredRuntimeAvailable(): void { + if (process.versions['bun'] === undefined) assertAuthoredPromiseHooks(); + else nodeAuthority(); +} + +/** Bun retains the root lease and local workers. Only the authored body uses Node. */ +export async function runAuthoredInNode( + metadata: AuthoredRootMetadata, + socketPath: string, + rootRunId: string, + options: ExecuteAuthoredFlowOptions, +): Promise { + const authority = nodeAuthority(); + const source = embeddedSource!; + const runtime: AuthoredExecutionRuntime = { + kind: 'node', version: authority.version, + executableSha256: authority.executableSha256, payloadSha256: hash(source), + }; + options.signal?.throwIfAborted(); + const directory = await mkdtemp(join(tmpdir(), 'flows-authored-node-')); + const entry = join(directory, 'runner.mjs'); + try { + await writeFile(entry, source, { mode: 0o400, flag: 'wx' }); + if (hash(await readFile(entry)) !== runtime.payloadSha256) throw refusal(); + const child = spawn(authority.path, ['--experimental-transform-types', entry], { + cwd: process.cwd(), env: process.env, + stdio: ['pipe', 'inherit', 'inherit', 'pipe'], + }); + return await new Promise((resolve, reject) => { + let buffer = '', ready = false, result: AuthoredFlowExecutionResult | undefined; + let failure: Error | undefined, killTimer: ReturnType | undefined; + const stop = (error: Error): void => { + failure ??= error; + child.kill('SIGTERM'); + killTimer ??= setTimeout(() => child.kill('SIGKILL'), 2500); + killTimer.unref(); + }; + const abort = (): void => stop(new Error('authored root execution was aborted')); + options.signal?.addEventListener('abort', abort, { once: true }); + if (options.signal?.aborted) abort(); + const onSigint = (): void => { child.kill('SIGINT'); }; + const onSigterm = (): void => { child.kill('SIGTERM'); }; + process.on('SIGINT', onSigint); process.on('SIGTERM', onSigterm); + const startupTimer = setTimeout(() => stop(new Error('authored runtime readiness timed out')), 10_000); + startupTimer.unref(); + const pipe = child.stdio[3] as Readable; + pipe.setEncoding('utf8'); + pipe.on('data', (chunk: string) => { + buffer += chunk; + if (Buffer.byteLength(buffer) > 16 * 1024 * 1024) { stop(new Error('authored runtime response exceeded limit')); return; } + for (;;) { + const end = buffer.indexOf('\n'); if (end < 0) break; + const line = buffer.slice(0, end); buffer = buffer.slice(end + 1); + try { + const message = JSON.parse(line); + if (message.type === 'ready' && !ready && !result) { + if (JSON.stringify(message.runtime) !== JSON.stringify(runtime)) throw refusal(); + ready = true; + clearTimeout(startupTimer); + // Keep stdin open: EOF tells the child its lease-owning parent died. + child.stdin!.write(JSON.stringify({ metadata, socketPath, rootRunId, + dataDir: options.dataDir, localAgentStream: options.localAgentStream }) + '\n'); + } else if (!ready || result) throw new Error('unexpected authored runtime message'); + else if (message.type === 'progress') options.onProgress?.(message.event); + else if (message.type === 'wait') options.onWait?.(message.event); + else if (message.type === 'result') result = { ...message.result, executionRuntime: runtime }; + else if (message.type === 'error') { + failure = typeof message.code === 'string' + ? new AuthoredFlowExecutionError(message.code as AuthoredFlowExecutionErrorCode, + message.message, message.completionReason, message.runId) + : new Error(message.message); + } else throw new Error('unknown authored runtime message'); + } catch (error) { stop(error instanceof Error ? error : new Error('invalid authored runtime message')); } + } + }); + child.stdin!.on('error', () => stop(new Error('authored runtime request pipe closed'))); + pipe.on('error', () => stop(new Error('authored runtime result pipe failed'))); + child.once('error', () => { failure = refusal(); }); + child.once('close', (code) => { + clearTimeout(startupTimer); + if (killTimer !== undefined) clearTimeout(killTimer); + options.signal?.removeEventListener('abort', abort); + process.off('SIGINT', onSigint); process.off('SIGTERM', onSigterm); + if (failure) reject(failure); + else if (code !== 0 || !ready || result === undefined || buffer.length > 0) reject(new Error('authored runtime exited without a complete result')); + else resolve(result); + }); + }); + } finally { await rm(directory, { recursive: true, force: true }); } +} diff --git a/packages/sdk/src/authored-root.ts b/packages/sdk/src/authored-root.ts index 1d2c4eb4c..5c7fb59a0 100644 --- a/packages/sdk/src/authored-root.ts +++ b/packages/sdk/src/authored-root.ts @@ -1,3 +1,5 @@ +import { loadPinnedAuthoredSource } from './authored-source-authority.js'; +import { assertAuthoredRuntimeAvailable, runAuthoredInNode } from './authored-node-runner.js'; import { createHash, randomUUID } from 'node:crypto'; import { readFile } from 'node:fs/promises'; import { canonicalize } from './canonical.js'; @@ -49,6 +51,7 @@ export async function executeDurableAuthoredFlow( input: unknown, options: DurableAuthoredOptions, ): Promise { + assertAuthoredRuntimeAvailable(); const source = await readFile(loaded.sourcePath); const sources = await Promise.all(loaded.graph.map(async node => Object.freeze({ path: node.path, @@ -108,28 +111,8 @@ export async function resumeDurableAuthoredFlow( ): Promise<(AuthoredFlowExecutionResult & { readonly rootRunId: string }) | undefined> { const metadata = await readAuthoredRootMetadata(journal, rootRunId); if (metadata === undefined) return undefined; - const source = await readFile(metadata.flowPath); - if (sha256(source) !== metadata.sourceSha256) { - throw new Error('authored root source authority mismatch'); - } - for (const pinned of metadata.sources) { - if (sha256(await readFile(pinned.path)) !== pinned.sourceSha256) { - throw new Error(`authored root source authority mismatch for "${pinned.path}"`); - } - } - const loaded = await loadAuthoredFlow(metadata.flowPath); - if (canonicalize(loaded.surfaceAuthority) !== canonicalize(metadata.surface) - || loaded.getDefinition(loaded.handle).name !== metadata.flowName) { - throw new Error('authored root Surface module authority mismatch'); - } - const loadedSources = loaded.graph.map(node => ({ - path: node.path, - sourceSha256: metadata.sources.find(source => source.path === node.path)?.sourceSha256 ?? '', - surface: node.surfaceAuthority, - })); - if (canonicalize(loadedSources) !== canonicalize(metadata.sources)) { - throw new Error('authored root declared source graph authority mismatch'); - } + assertAuthoredRuntimeAvailable(); + const loaded = await loadPinnedAuthoredSource(metadata); if (metadata.localAgentStream !== options.localAgentStream) { throw new Error('authored root local agent surface mismatch'); } @@ -188,6 +171,13 @@ async function driveRoot( try { const result = await withWorkerLease(peer, dispatch, async rootSignal => { const callerSignal = options.lifecycle?.signal; + const signal = callerSignal === undefined ? rootSignal : AbortSignal.any([callerSignal, rootSignal]); + if (process.versions['bun'] !== undefined) { + return runAuthoredInNode(metadata, journal.socketPath, dispatch.run_id, { + dataDir: options.dataDir, localAgentStream: options.localAgentStream, + ...options.lifecycle, signal, + }); + } return await executeAuthoredFlow( loaded.handle, journal, @@ -209,7 +199,8 @@ async function driveRoot( dispatch.run_id, dispatch.step_id, dispatch.attempt, dispatch.idempotency_key, 'success', { output: { name: result.name, completionReason: result.completionReason, - journalSteps: result.journalSteps }, + journalSteps: result.journalSteps, + ...(result.executionRuntime === undefined ? {} : { executionRuntime: result.executionRuntime }) }, started_pins: dispatch.pins, end_pins: dispatch.pins, }, ); diff --git a/packages/sdk/src/authored-runtime-capability.ts b/packages/sdk/src/authored-runtime-capability.ts new file mode 100644 index 000000000..7297196a2 --- /dev/null +++ b/packages/sdk/src/authored-runtime-capability.ts @@ -0,0 +1,31 @@ +import { createHook } from 'node:async_hooks'; +import { AuthoredFlowExecutionError } from './authored-flow-error.js'; + +/** A runtime with no-op promise hooks cannot distinguish await from ignored then. */ +export function assertAuthoredPromiseHooks(): void { + let initialized = false; + let resolved = false; + const hook = createHook({ + init(_id, type) { if (type === 'PROMISE') initialized = true; }, + promiseResolve() { resolved = true; }, + }); + try { + hook.enable(); + void new Promise((resolve) => resolve()); + } finally { hook.disable(); } + if (!initialized || !resolved) throw new AuthoredFlowExecutionError( + 'unsupported_promise_lifecycle', + 'authored execution requires working Node promise hooks; use the standalone Node-authored runner with Node >=22.14', + ); +} + +export function assertAuthoredNodeVersion(version = process.versions.node): void { + const [major, minor] = version.split('.').map(Number); + if (!Number.isInteger(major) || !Number.isInteger(minor) + || major! < 22 || (major === 22 && minor! < 14) + || process.versions['bun'] !== undefined) { + throw new AuthoredFlowExecutionError('unsupported_promise_lifecycle', + 'the authored runner requires Node >=22.14; no workflow body was executed'); + } + assertAuthoredPromiseHooks(); +} diff --git a/packages/sdk/src/authored-source-authority.ts b/packages/sdk/src/authored-source-authority.ts new file mode 100644 index 000000000..b0cf7f35f --- /dev/null +++ b/packages/sdk/src/authored-source-authority.ts @@ -0,0 +1,33 @@ +import { createHash } from 'node:crypto'; +import { readFile } from 'node:fs/promises'; +import { canonicalize } from './canonical.js'; +import { loadAuthoredFlow } from './authored-flow-loader.js'; +import type { AuthoredRootMetadata } from './authored-root.js'; +const sha256 = (bytes: Uint8Array): string => createHash('sha256').update(bytes).digest('hex'); + +/** Re-open the exact source graph and Surface instance pinned by the durable root. */ +export async function loadPinnedAuthoredSource(metadata: AuthoredRootMetadata) { + const source = await readFile(metadata.flowPath); + if (sha256(source) !== metadata.sourceSha256) { + throw new Error('authored root source authority mismatch'); + } + for (const pinned of metadata.sources) { + if (sha256(await readFile(pinned.path)) !== pinned.sourceSha256) { + throw new Error(`authored root source authority mismatch for "${pinned.path}"`); + } + } + const loaded = await loadAuthoredFlow(metadata.flowPath); + if (canonicalize(loaded.surfaceAuthority) !== canonicalize(metadata.surface) + || loaded.getDefinition(loaded.handle).name !== metadata.flowName) { + throw new Error('authored root Surface module authority mismatch'); + } + const loadedSources = loaded.graph.map(node => ({ + path: node.path, + sourceSha256: metadata.sources.find(source => source.path === node.path)?.sourceSha256 ?? '', + surface: node.surfaceAuthority, + })); + if (canonicalize(loadedSources) !== canonicalize(metadata.sources)) { + throw new Error('authored root declared source graph authority mismatch'); + } + return loaded; +} diff --git a/packages/sdk/src/cli/direct-run.ts b/packages/sdk/src/cli/direct-run.ts index 28a9c6ba7..ba09e72ff 100644 --- a/packages/sdk/src/cli/direct-run.ts +++ b/packages/sdk/src/cli/direct-run.ts @@ -154,6 +154,7 @@ export async function runDirectFlow( || error.code === 'helper_slack.mount_required' || error.code === 'budget_syntax_invalid' || error.code === 'budget_missing_price' + || error.code === 'unsupported_promise_lifecycle' || error.code === 'unsupported_header' || error.code === 'agent_cli_unresolved' || error.code === 'llm_cli_unresolved' diff --git a/packages/sdk/src/cli/run.ts b/packages/sdk/src/cli/run.ts index 98e4d9b43..7536234b8 100644 --- a/packages/sdk/src/cli/run.ts +++ b/packages/sdk/src/cli/run.ts @@ -213,6 +213,10 @@ export async function resumeFlow( return { exitCode: 2, report: { ...base, runId, socketPath, diagnostics: [{ severity: 'refusal', kind: 'human_influenced_run', message: error.message.replace(/^human_influenced_run: /, '') }] } }; } + if (error instanceof AuthoredFlowExecutionError && error.code === 'unsupported_promise_lifecycle') { + return { exitCode: 2, report: { ...base, runId, socketPath, + diagnostics: [{ severity: 'refusal', kind: 'invalid_spec', message: error.message }] } }; + } if (error instanceof AuthoredFlowExecutionError && (error.code === 'helper_slack.credential_missing' || error.code === 'helper_slack.mount_required' || error.code === 'helper_provider.mount_required' || error.code === 'helper_provider.unsupported')) { diff --git a/packages/sdk/tests/authored-node-runtime.test.ts b/packages/sdk/tests/authored-node-runtime.test.ts new file mode 100644 index 000000000..bccb5d814 --- /dev/null +++ b/packages/sdk/tests/authored-node-runtime.test.ts @@ -0,0 +1,134 @@ +import { spawn, spawnSync } from 'node:child_process'; +import { chmodSync, cpSync, existsSync, mkdtempSync, readFileSync, rmSync, writeFileSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import { join, resolve } from 'node:path'; +import { afterAll, beforeAll, describe, expect, it } from 'vitest'; +import { JournalClient } from '../src/journal-client.js'; +import { socketFor } from '../src/cli/run.js'; +import { assertAuthoredNodeVersion } from '../src/authored-runtime-capability.js'; + +const root = resolve('../..'), sdk = resolve('.'); +const fixtures: string[] = []; +let stage: string, cli: string; +const daemon = process.env['RELAYFLOWD_BIN'] ?? resolve('../../kernel/target/debug/relayflowd'); +const bun = process.env['FLOWS_BUILD_BUN'] ?? 'bun'; +const wrapperHelper = resolve('../../testdata/preflight/wrapper-session.mjs'); + +beforeAll(() => { + expect(spawnSync(bun, ['--version'], { encoding: 'utf8' }).stdout.trim()).toBe('1.4.0'); + expect(existsSync(daemon), 'build the current kernel or set RELAYFLOWD_BIN').toBe(true); + stage = mkdtempSync(join(tmpdir(), 'authored-standalone-build-')); + cli = join(stage, 'flows'); + const built = spawnSync(process.execPath, [join(root, 'scripts/build-standalone-cli.mjs'), + process.platform === 'darwin' ? 'bun-darwin-arm64' : 'bun-linux-x64', cli], { + cwd: root, encoding: 'utf8', timeout: 120_000, env: { ...process.env, FLOWS_BUILD_BUN: bun }, + }); + expect(built.status, built.stderr + built.stdout).toBe(0); +}, 130_000); + +afterAll(() => { + for (const directory of fixtures) { + const connection = join(directory, 'data/connection.json'); + if (existsSync(connection)) { + const { pid } = JSON.parse(readFileSync(connection, 'utf8')); + if (typeof pid === 'number') { try { process.kill(pid, 'SIGTERM'); } catch (error) { + if ((error as NodeJS.ErrnoException).code !== 'ESRCH') throw error; + } } + } + rmSync(directory, { recursive: true, force: true }); + } + if (stage) rmSync(stage, { recursive: true, force: true }); +}); + +function fixture(body: string) { + const directory = mkdtempSync(join(tmpdir(), 'authored-node-runtime-')); fixtures.push(directory); + const surface = join(directory, 'node_modules/@relayflows/surface'); + cpSync(join(sdk, 'node_modules/@relayflows/surface'), surface, { recursive: true }); + const manifest = JSON.parse(readFileSync(join(surface, 'package.json'), 'utf8')); + // Same verified execution envelope Cloud uses around the unchanged Surface payload. + writeFileSync(join(surface, 'package.json'), JSON.stringify({ name: manifest.name, + version: manifest.version, private: true, type: 'module', types: './index.d.ts', + exports: { '.': './index.js', './runtime': './runtime.js', './triggers': './triggers/index.js', './triggers/*': './triggers/*.js' } })); + writeFileSync(join(surface, 'index.js'), "export * from './dist/index.js';\n"); + writeFileSync(join(surface, 'runtime.js'), "export { getFlowDefinition } from './dist/flow.js';\n"); + writeFileSync(join(surface, 'index.d.ts'), "export * from './dist/index.js';\n"); + writeFileSync(join(directory, 'package.json'), '{"type":"module"}'); + const wrapper = join(directory, 'agent.mjs'); + writeFileSync(wrapper, `#!/usr/bin/env node +import { receiveWrapperRequest } from ${JSON.stringify(wrapperHelper)}; +import { appendFileSync } from 'node:fs'; +if(process.argv[2]==='auth') process.exit(0); +const request=await receiveWrapperRequest(); +if(request){appendFileSync('agent-effects','once\\n');await new Promise(r=>setTimeout(r,100));console.log('agent-ok');} +`); + chmodSync(wrapper, 0o755); + writeFileSync(join(directory, 'flows.json'), JSON.stringify({ cli: wrapper })); + writeFileSync(join(directory, 'case.flow.ts'), `import {flow} from '@relayflows/surface'; +import {appendFileSync,existsSync,writeFileSync} from 'node:fs'; +export default flow('runtime-case',async f=>{${body}}); +`); + const env = { ...process.env, FLOWS_AUTHORED_NODE: process.execPath, + RELAYFLOWD_BIN: daemon, NODE_PATH: join(directory, 'node_modules') }; + const flags = ['--local-agent', '--data-dir', join(directory, 'data'), '--json', '--no-observer-link']; + const invoke = (args: string[], overrides: NodeJS.ProcessEnv = {}) => spawnSync(cli, [...args, ...flags], { + cwd: directory, env: { ...env, ...overrides }, encoding: 'utf8', timeout: 60_000, + }); + return { directory, env, flags, invoke, run: () => invoke(['run', 'case.flow.ts', '--input', '{}']) }; +} + +const sequential = `await f.agent('worker',{task:'local fixture'}); +await f.run("printf one >> run-effects"); +await f.run("printf two >> run-effects"); +await f.run("printf three >> run-effects");`; + +async function entries(directory: string, runId: string) { + const client = new JournalClient(socketFor(join(directory, 'data'))); + await client.connect(); await client.hello('authored-node-test'); + try { return (await client.journalRead(runId, 1)).entries as Array<{entry_type:string;step_id?:string;payload:Record}>; } + finally { client.close(); } +} + +describe('Bun 1.4.0 standalone → native Node authored lifecycle', () => { + it('awaits agent plus three run steps and resumes without repeating effects', async () => { + const f = fixture(sequential + `f.done('success');`); + const first = f.run(); expect(first.status, first.stderr + first.stdout).toBe(0); + const report = JSON.parse(first.stdout); expect(report).toMatchObject({ok:true,completionReason:'success',completedSteps:5}); + const resumed = f.invoke(['resume', report.runId]); expect(resumed.status, resumed.stderr + resumed.stdout).toBe(0); + expect(JSON.parse(resumed.stdout).runId).toBe(report.runId); + expect(readFileSync(join(f.directory,'agent-effects'),'utf8')).toBe('once\n'); + expect(readFileSync(join(f.directory,'run-effects'),'utf8')).toBe('onetwothree'); + const rootEntries = await entries(f.directory, report.runId); + expect(rootEntries.filter(e=>e.entry_type==='step.attempt.started')).toHaveLength(1); + const output = rootEntries.find(e=>e.entry_type==='step.completed')!.payload['output']; + expect(output.executionRuntime).toMatchObject({kind:'node',version:process.versions.node}); + expect(output.executionRuntime.executableSha256).toMatch(/^[a-f0-9]{64}$/); + expect(output.executionRuntime.payloadSha256).toMatch(/^[a-f0-9]{64}$/); + expect(output.journalSteps.map((s:{id:string})=>s.id)).toEqual(['agent-1','run-2','run-3','run-4','complete-5']); + }, 90_000); + + it.each([ + ['unawaited', `f.run("printf ignored >> forbidden-effects");f.done('success');`], + ['manual then', `f.run("printf chained >> chain-effects").then(()=>undefined);await new Promise(r=>setTimeout(r,100));f.done('success');`], + ])('refuses %s rather than reporting terminal success', (_name,body) => { + const f=fixture(body);const result=f.run();expect(result.status).toBe(1); + expect(result.stdout+result.stderr).toContain('unawaited_step'); + expect(existsSync(join(f.directory,'forbidden-effects'))).toBe(false); + }, 60_000); + + it('refuses missing Node before body effects or root admission', () => { + const f=fixture(`writeFileSync('body-started','bad');${sequential}f.done('success');`); + const result=f.invoke(['run','case.flow.ts','--input','{}'],{FLOWS_AUTHORED_NODE:join(f.directory,'missing-node')}); + expect(result.status).toBe(2);expect(result.stdout+result.stderr).toContain('unsupported_promise_lifecycle'); + expect(existsSync(join(f.directory,'body-started'))).toBe(false); + expect(existsSync(join(f.directory,'agent-effects'))).toBe(false); + }, 30_000); + + it('refuses an old Node candidate before body effects', () => { + expect(()=>assertAuthoredNodeVersion('22.13.0')).toThrow('Node >=22.14'); + const f=fixture(`writeFileSync('body-started','bad');f.done('success');`); + const old=join(f.directory,'old-node'); + writeFileSync(old, '#!/bin/sh\nprintf \'%s\' \''+JSON.stringify({path:old,version:'20.19.0'})+'\'\n');chmodSync(old,0o755); + const result=f.invoke(['run','case.flow.ts','--input','{}'],{FLOWS_AUTHORED_NODE:old}); + expect(result.status).toBe(2);expect(existsSync(join(f.directory,'body-started'))).toBe(false); + }, 30_000); +}); diff --git a/scripts/build-standalone-cli.mjs b/scripts/build-standalone-cli.mjs new file mode 100644 index 000000000..73c0b2db2 --- /dev/null +++ b/scripts/build-standalone-cli.mjs @@ -0,0 +1,30 @@ +#!/usr/bin/env node +import { mkdtemp, readFile, rm, writeFile } from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import { dirname, join, resolve } from 'node:path'; +import { fileURLToPath } from 'node:url'; +import { spawnSync } from 'node:child_process'; +const root = resolve(dirname(fileURLToPath(import.meta.url)), '..'); +const args = process.argv.slice(2); +const target = args[0], outfile = args[1]; +if (!/^bun-(linux-x64|darwin-arm64)$/.test(target ?? '') || !outfile) { + throw new Error('usage: build-standalone-cli.mjs '); +} +const stage = await mkdtemp(join(tmpdir(), 'flows-standalone-')); +function bun(args) { + const result = spawnSync(process.env.FLOWS_BUILD_BUN ?? 'bun', args, { cwd: root, stdio: 'inherit' }); + if (result.error || result.status !== 0) throw new Error('standalone CLI build failed'); +} +try { + const payload = join(stage, 'authored-node.mjs'); + bun(['build', 'packages/sdk/src/authored-node-entry.ts', '--target=node', '--outfile='+payload]); + const entry = join(stage, 'standalone.ts'); + await writeFile(entry, ` +import { installAuthoredNodeSource } from ${JSON.stringify(join(root, 'packages/sdk/src/authored-node-runner.ts'))}; +import { runCli } from ${JSON.stringify(join(root, 'packages/sdk/src/cli.ts'))}; +installAuthoredNodeSource(${JSON.stringify(await readFile(payload, 'utf8'))}); +process.exitCode = await runCli(process.argv.slice(2)); +`); + bun(['build', entry, '--compile', '--target='+target, '--outfile='+resolve(outfile), + '--env=disable', '--no-compile-autoload-dotenv', '--no-compile-autoload-bunfig']); +} finally { await rm(stage, { recursive: true, force: true }); } From fbb4dec417c76494962fec593ff7c581b09b535a Mon Sep 17 00:00:00 2001 From: Miya Date: Tue, 15 Sep 2026 19:32:26 +0200 Subject: [PATCH 2/6] test: cover authored parent loss and same-root child replay Session-Id: 01a09c40-ce3b-7f11-a7df-b6b7ccab6fd9 --- .../sdk/tests/authored-node-runtime.test.ts | 53 ++++++++++++++++++- 1 file changed, 52 insertions(+), 1 deletion(-) diff --git a/packages/sdk/tests/authored-node-runtime.test.ts b/packages/sdk/tests/authored-node-runtime.test.ts index bccb5d814..da558dbe0 100644 --- a/packages/sdk/tests/authored-node-runtime.test.ts +++ b/packages/sdk/tests/authored-node-runtime.test.ts @@ -1,5 +1,5 @@ import { spawn, spawnSync } from 'node:child_process'; -import { chmodSync, cpSync, existsSync, mkdtempSync, readFileSync, rmSync, writeFileSync } from 'node:fs'; +import { chmodSync, cpSync, existsSync, mkdtempSync, readFileSync, readdirSync, rmSync, writeFileSync } from 'node:fs'; import { tmpdir } from 'node:os'; import { join, resolve } from 'node:path'; import { afterAll, beforeAll, describe, expect, it } from 'vitest'; @@ -106,6 +106,57 @@ describe('Bun 1.4.0 standalone → native Node authored lifecycle', () => { expect(output.journalSteps.map((s:{id:string})=>s.id)).toEqual(['agent-1','run-2','run-3','run-4','complete-5']); }, 90_000); + it('stops on parent death and replays completed children under the same unfinished root', async () => { + const f=fixture(sequential + ` +if(!existsSync('resume-ready')){ + writeFileSync('node-pid',String(process.pid)); + writeFileSync('resume-ready','yes'); + await new Promise(()=>{}); +} +f.done('success');`); + const child=spawn(cli,['run','case.flow.ts','--input','{}',...f.flags],{ + cwd:f.directory,env:f.env,stdio:['ignore','pipe','pipe'], + }); + let logs='';child.stdout.on('data',bytes=>{logs+=bytes});child.stderr.on('data',bytes=>{logs+=bytes}); + const closed=new Promise(resolve=>child.once('close',()=>resolve())); + let nodePid:number|undefined; + try { + const deadline=Date.now()+15_000; + while(!existsSync(join(f.directory,'resume-ready')) && Date.now()setTimeout(r,50)); + } + expect(existsSync(join(f.directory,'resume-ready')),logs).toBe(true); + nodePid=Number(readFileSync(join(f.directory,'node-pid'),'utf8')); + let rootId:string|undefined; + for(const file of readdirSync(join(f.directory,'data/runs')).filter(name=>name.endsWith('.sqlite3'))){ + const id=file.slice(0,-8);const journal=await entries(f.directory,id); + if(journal[0]?.payload['spec']?.steps?.[0]?.id==='authored-root')rootId=id; + } + expect(rootId).toBeDefined(); + child.kill('SIGKILL');await closed; + let alive=true; + for(let attempt=0;attempt<100;attempt++){ + try{process.kill(nodePid,0);}catch(error){ + if((error as NodeJS.ErrnoException).code==='ESRCH'){alive=false;break;}throw error; + } + await new Promise(r=>setTimeout(r,20)); + } + expect(alive,'lease-owning parent loss must terminate Node body').toBe(false);nodePid=undefined; + const resumed=f.invoke(['resume',rootId!]);expect(resumed.status,resumed.stderr+resumed.stdout).toBe(0); + expect(JSON.parse(resumed.stdout).runId).toBe(rootId); + expect(readFileSync(join(f.directory,'agent-effects'),'utf8')).toBe('once\n'); + expect(readFileSync(join(f.directory,'run-effects'),'utf8')).toBe('onetwothree'); + const rootJournal=await entries(f.directory,rootId!); + expect(rootJournal.filter(e=>e.entry_type==='step.attempt.started')).toHaveLength(2); + expect(rootJournal.filter(e=>e.entry_type==='run.completed')).toHaveLength(1); + expect(rootJournal.at(-1)?.payload['completionReason']).toBe('success'); + } finally { + child.kill('SIGKILL'); + if(nodePid!==undefined){try{process.kill(nodePid,'SIGKILL');}catch(error){if((error as NodeJS.ErrnoException).code!=='ESRCH')throw error;}} + } + },60_000); + it.each([ ['unawaited', `f.run("printf ignored >> forbidden-effects");f.done('success');`], ['manual then', `f.run("printf chained >> chain-effects").then(()=>undefined);await new Promise(r=>setTimeout(r,100));f.done('success');`], From 2e44424120357f3321513f68eab545744d6043bd Mon Sep 17 00:00:00 2001 From: Miya Date: Tue, 15 Sep 2026 19:39:12 +0200 Subject: [PATCH 3/6] fix: preserve parent signal interruption for authored resume Session-Id: 01a09c40-ce3b-7f11-a7df-b6b7ccab6fd9 --- packages/sdk/src/authored-node-runner.ts | 6 ++---- packages/sdk/tests/authored-node-runtime.test.ts | 4 ++-- 2 files changed, 4 insertions(+), 6 deletions(-) diff --git a/packages/sdk/src/authored-node-runner.ts b/packages/sdk/src/authored-node-runner.ts index 0fc542b46..f646ac690 100644 --- a/packages/sdk/src/authored-node-runner.ts +++ b/packages/sdk/src/authored-node-runner.ts @@ -87,9 +87,8 @@ export async function runAuthoredInNode( const abort = (): void => stop(new Error('authored root execution was aborted')); options.signal?.addEventListener('abort', abort, { once: true }); if (options.signal?.aborted) abort(); - const onSigint = (): void => { child.kill('SIGINT'); }; - const onSigterm = (): void => { child.kill('SIGTERM'); }; - process.on('SIGINT', onSigint); process.on('SIGTERM', onSigterm); + // Do not intercept the parent's signals: its normal exit keeps the root + // recoverable. The child observes parent EOF (and shared process-group signals). const startupTimer = setTimeout(() => stop(new Error('authored runtime readiness timed out')), 10_000); startupTimer.unref(); const pipe = child.stdio[3] as Readable; @@ -129,7 +128,6 @@ export async function runAuthoredInNode( clearTimeout(startupTimer); if (killTimer !== undefined) clearTimeout(killTimer); options.signal?.removeEventListener('abort', abort); - process.off('SIGINT', onSigint); process.off('SIGTERM', onSigterm); if (failure) reject(failure); else if (code !== 0 || !ready || result === undefined || buffer.length > 0) reject(new Error('authored runtime exited without a complete result')); else resolve(result); diff --git a/packages/sdk/tests/authored-node-runtime.test.ts b/packages/sdk/tests/authored-node-runtime.test.ts index da558dbe0..1467c1a88 100644 --- a/packages/sdk/tests/authored-node-runtime.test.ts +++ b/packages/sdk/tests/authored-node-runtime.test.ts @@ -106,7 +106,7 @@ describe('Bun 1.4.0 standalone → native Node authored lifecycle', () => { expect(output.journalSteps.map((s:{id:string})=>s.id)).toEqual(['agent-1','run-2','run-3','run-4','complete-5']); }, 90_000); - it('stops on parent death and replays completed children under the same unfinished root', async () => { + it.each(['SIGKILL', 'SIGTERM'] as const)('stops on parent %s and replays completed children under the same unfinished root', async signal => { const f=fixture(sequential + ` if(!existsSync('resume-ready')){ writeFileSync('node-pid',String(process.pid)); @@ -134,7 +134,7 @@ f.done('success');`); if(journal[0]?.payload['spec']?.steps?.[0]?.id==='authored-root')rootId=id; } expect(rootId).toBeDefined(); - child.kill('SIGKILL');await closed; + child.kill(signal);await closed; let alive=true; for(let attempt=0;attempt<100;attempt++){ try{process.kill(nodePid,0);}catch(error){ From b64b07b753cee7f3810ec9617112bcc3df439d9d Mon Sep 17 00:00:00 2001 From: Miya Date: Tue, 15 Sep 2026 19:56:15 +0200 Subject: [PATCH 4/6] fix: verify authored runtime receipts and parent-loss custody Session-Id: 01a09c40-ce3b-7f11-a7df-b6b7ccab6fd9 --- docs/SURFACE.md | 7 +- packages/sdk/src/authored-node-entry.ts | 36 ++++++-- packages/sdk/src/authored-node-runner.ts | 82 +++++++++++++++++-- .../sdk/src/authored-runtime-capability.ts | 4 +- packages/sdk/src/authored-source-authority.ts | 33 +++++++- .../sdk/tests/authored-node-result.test.ts | 52 ++++++++++++ .../sdk/tests/authored-node-runtime.test.ts | 41 ++++++++-- 7 files changed, 232 insertions(+), 23 deletions(-) create mode 100644 packages/sdk/tests/authored-node-result.test.ts diff --git a/docs/SURFACE.md b/docs/SURFACE.md index c5a5e0ee5..c05b2efaa 100644 --- a/docs/SURFACE.md +++ b/docs/SURFACE.md @@ -821,7 +821,12 @@ root worker lease, agent workers, declarative execution, daemon lookup, and resume protocol. The child inherits the working directory/environment and uses the same journal socket and child admission identities. It verifies the root's pinned source graph and Surface package before executing the body. Root aborts -and signals stop the child; loss of the parent pipe stops it immediately. +and signals stop the child. Parent-pipe loss exits a responsive child; an independent +watchdog thread checks parent identity every 50 ms and kills the process even if +the authored body blocks its event loop. This is bounded scheduling, not an +instantaneous termination guarantee. Result frames are authenticated and checked +against durable child completion before root success. The Node loader executes +captured, hash-verified authored source bytes at their original module URLs. Successful root output records the Node version, executable SHA256, and embedded payload SHA256 as `executionRuntime`. The payload remains part of the existing artifact hash, and completed child effects remain journal results on resume. diff --git a/packages/sdk/src/authored-node-entry.ts b/packages/sdk/src/authored-node-entry.ts index dce9e2a8b..d01114a88 100644 --- a/packages/sdk/src/authored-node-entry.ts +++ b/packages/sdk/src/authored-node-entry.ts @@ -1,4 +1,5 @@ -import { createHash } from 'node:crypto'; +import { Worker } from 'node:worker_threads'; +import { createHash, createHmac } from 'node:crypto'; import { readFileSync, writeSync } from 'node:fs'; import { JournalClient } from './journal-client.js'; import { executeAuthoredFlow } from './authored-flow-executor.js'; @@ -7,7 +8,16 @@ import { assertAuthoredNodeVersion } from './authored-runtime-capability.js'; import { AuthoredFlowExecutionError } from './authored-flow-error.js'; import type { AuthoredRootMetadata } from './authored-root.js'; -const send = (message: unknown): void => { writeSync(3, JSON.stringify(message) + '\n'); }; +let channelKey: string | undefined, sequence = 0; +// Capture writers before loading authored modules; credentials never enter env. +const writeFrame = writeSync, mac = createHmac; +const send = (message: unknown): void => { + const payload = JSON.stringify(message); + if (channelKey === undefined) { writeFrame(3, payload + '\n'); return; } + const seq = ++sequence; + writeFrame(3, JSON.stringify({ seq, payload, + mac: mac('sha256', channelKey).update(`${seq}\0${payload}`).digest('hex') }) + '\n'); +}; const hash = (path: string): string => createHash('sha256').update(readFileSync(path)).digest('hex'); const controller = new AbortController(); let finished = false; @@ -19,12 +29,22 @@ process.on('SIGINT', abort); process.on('SIGTERM', abort); // The parent owns the root lease. Do not continue authored effects after it dies. process.stdin.on('end', () => { if (!finished) process.exit(1); }); process.stdin.on('error', () => process.exit(1)); +// A separate event loop enforces parent loss even while authored JS is blocked. +// Capture the expected PID in the parent's spawn arguments, before any child work. +const parentPid = Number(process.argv[2]); +if (!Number.isSafeInteger(parentPid) || parentPid <= 1) throw new Error('invalid authored parent identity'); +const watchdog = new Worker(` + const {parentPort,workerData}=require('node:worker_threads'); + function check(){if(process.ppid!==workerData.parentPid)process.kill(process.pid,'SIGKILL');} + check();setInterval(check,50);parentPort.postMessage('ready'); +`, { eval: true, workerData: { parentPid } }); let client: JournalClient | undefined; try { assertAuthoredNodeVersion(); + await new Promise((resolve,reject)=>{watchdog.once('message',()=>resolve());watchdog.once('error',reject);}); send({ type: 'ready', runtime: { kind: 'node', version: process.versions.node, executableSha256: hash(process.execPath), payloadSha256: hash(process.argv[1]!) } }); - const request = await new Promise<{ metadata: AuthoredRootMetadata; socketPath: string; + const request = await new Promise<{ channelKey: string; metadata: AuthoredRootMetadata; socketPath: string; rootRunId: string; dataDir: string; localAgentStream?: string }>((resolve, reject) => { let buffer = ''; process.stdin.setEncoding('utf8'); @@ -38,8 +58,10 @@ try { }; process.stdin.on('data', onData); }); + if (typeof request.channelKey !== 'string' || !/^[a-f0-9]{64}$/.test(request.channelKey)) throw new Error('invalid authored channel key'); + channelKey = request.channelKey; controller.signal.throwIfAborted(); - const loaded = await loadPinnedAuthoredSource(request.metadata); + const loaded = await loadPinnedAuthoredSource(request.metadata, true); if (request.localAgentStream !== request.metadata.localAgentStream) throw new Error('authored root local agent surface mismatch'); client = new JournalClient(request.socketPath); await client.connect(); await client.hello('flows-authored-node'); @@ -53,11 +75,13 @@ try { }); send({ type: 'result', result }); } catch (error) { - send({ type: 'error', message: error instanceof Error ? error.message : 'authored body failed', + const message = error instanceof Error ? error.message : 'authored body failed'; + const prefix = error instanceof AuthoredFlowExecutionError ? `${error.code}: ` : ''; + send({ type: 'error', message: prefix && message.startsWith(prefix) ? message.slice(prefix.length) : message, ...(error instanceof AuthoredFlowExecutionError ? { code: error.code, completionReason: error.completionReason, runId: error.runId } : {}) }); process.exitCode = 1; } finally { - finished = true; client?.close(); process.stdin.destroy(); + finished = true; await watchdog.terminate(); client?.close(); process.stdin.destroy(); process.off('SIGINT', abort); process.off('SIGTERM', abort); } diff --git a/packages/sdk/src/authored-node-runner.ts b/packages/sdk/src/authored-node-runner.ts index f646ac690..97d7dfb48 100644 --- a/packages/sdk/src/authored-node-runner.ts +++ b/packages/sdk/src/authored-node-runner.ts @@ -1,5 +1,6 @@ +import { JournalClient } from './journal-client.js'; import { spawn, spawnSync } from 'node:child_process'; -import { createHash } from 'node:crypto'; +import { createHash, createHmac, randomBytes, timingSafeEqual } from 'node:crypto'; import { mkdtemp, readFile, rm, writeFile } from 'node:fs/promises'; import { readFileSync, realpathSync } from 'node:fs'; import { tmpdir } from 'node:os'; @@ -28,11 +29,12 @@ function nodeAuthority(): NodeAuthority { if (override !== undefined && !isAbsolute(override)) throw refusal(); const probe = spawnSync(override ?? 'node', ['--input-type=module', '--eval', ` import {createHook} from 'node:async_hooks'; +import {stripTypeScriptTypes,register} from 'node:module'; const [major,minor]=process.versions.node.split('.').map(Number); let init=false,resolve=false; const h=createHook({init(_id,type){if(type==='PROMISE')init=true},promiseResolve(){resolve=true}}).enable(); new Promise(r=>r());h.disable(); -if(process.versions.bun || major<22 || (major===22 && minor<14) || !init || !resolve) process.exit(2); +if(process.versions.bun || major<22 || (major===22 && minor<14) || !init || !resolve || typeof stripTypeScriptTypes!=='function' || typeof register!=='function') process.exit(2); process.stdout.write(JSON.stringify({path:process.execPath,version:process.versions.node})); `], { encoding: 'utf8', timeout: 10_000, maxBuffer: 4096, env: process.env }); if (probe.error || probe.status !== 0) throw refusal(); @@ -71,11 +73,13 @@ export async function runAuthoredInNode( try { await writeFile(entry, source, { mode: 0o400, flag: 'wx' }); if (hash(await readFile(entry)) !== runtime.payloadSha256) throw refusal(); - const child = spawn(authority.path, ['--experimental-transform-types', entry], { + const child = spawn(authority.path, ['--experimental-transform-types', entry, String(process.pid)], { cwd: process.cwd(), env: process.env, stdio: ['pipe', 'inherit', 'inherit', 'pipe'], }); - return await new Promise((resolve, reject) => { + const result = await new Promise((resolve, reject) => { + const channelKey = randomBytes(32).toString('hex'); + let expectedSequence = 1; let buffer = '', ready = false, result: AuthoredFlowExecutionResult | undefined; let failure: Error | undefined, killTimer: ReturnType | undefined; const stop = (error: Error): void => { @@ -100,13 +104,26 @@ export async function runAuthoredInNode( const end = buffer.indexOf('\n'); if (end < 0) break; const line = buffer.slice(0, end); buffer = buffer.slice(end + 1); try { - const message = JSON.parse(line); + let message = JSON.parse(line); + if (ready) { + if (message.seq !== expectedSequence || typeof message.payload !== 'string' + || typeof message.mac !== 'string' || !/^[a-f0-9]{64}$/.test(message.mac)) { + throw new Error('invalid authored runtime authenticated frame'); + } + const expected = createHmac('sha256', channelKey) + .update(`${expectedSequence}\0${message.payload}`).digest(); + if (!timingSafeEqual(expected, Buffer.from(message.mac, 'hex'))) { + throw new Error('invalid authored runtime authenticated frame'); + } + expectedSequence++; + message = JSON.parse(message.payload); + } if (message.type === 'ready' && !ready && !result) { if (JSON.stringify(message.runtime) !== JSON.stringify(runtime)) throw refusal(); ready = true; clearTimeout(startupTimer); // Keep stdin open: EOF tells the child its lease-owning parent died. - child.stdin!.write(JSON.stringify({ metadata, socketPath, rootRunId, + child.stdin!.write(JSON.stringify({ channelKey, metadata, socketPath, rootRunId, dataDir: options.dataDir, localAgentStream: options.localAgentStream }) + '\n'); } else if (!ready || result) throw new Error('unexpected authored runtime message'); else if (message.type === 'progress') options.onProgress?.(message.event); @@ -133,5 +150,58 @@ export async function runAuthoredInNode( else resolve(result); }); }); + await verifyAuthoredNodeResult(result, metadata, rootRunId, socketPath); + return result; } finally { await rm(directory, { recursive: true, force: true }); } } + +/** The IPC frame is a claim, not a durable terminal fact or a sandbox boundary. */ +export async function verifyAuthoredNodeResult( + result: AuthoredFlowExecutionResult, metadata: AuthoredRootMetadata, + rootRunId: string, socketPath: string, +): Promise { + const invalid = (): never => { throw new Error('authored runtime result has no matching durable completion'); }; + if (result.rootRunId !== rootRunId || result.name !== metadata.flowName + || !['success', 'needs_human'].includes(result.completionReason) + || !Array.isArray(result.journalSteps) || result.journalSteps.length === 0) invalid(); + const terminal = result.journalSteps.at(-1)!; + if (!terminal || !/^complete-[1-9][0-9]*$/.test(terminal.id)) invalid(); + const count = Number(terminal.id.slice('complete-'.length)); + if (!Number.isSafeInteger(count) || result.journalSteps.length !== count) invalid(); + const ordinal = (id: string): number => Number(/-([1-9][0-9]*)$/.exec(id)?.[1]); + // Parallel awaits may finish in either order; validate a copy in declaration order. + const ordered = [...result.journalSteps].sort((a,b)=>ordinal(a.id)-ordinal(b.id)); + const runs = new Set(); + const journal = new JournalClient(socketPath); + await journal.connect(); + try { + await journal.hello('flows-authored-result-verifier'); + for (const [index, claimed] of ordered.entries()) { + if (!claimed || typeof claimed.id !== 'string' || typeof claimed.runId !== 'string' + || ordinal(claimed.id) !== index+1 || claimed.completionReason !== 'success' + || runs.has(claimed.runId)) invalid(); + runs.add(claimed.runId); + const state = await journal.runGet(claimed.runId); + if (state.run_id !== claimed.runId || state.status !== 'completed' + || state.steps[claimed.id]?.state !== 'done') invalid(); + const entries = (await journal.journalRead(claimed.runId, 1)).entries as Array<{ + entry_type: string; step_id?: string; payload?: { + completionReason?: string; spec?: { name?: string; steps?: Array<{id?:string;type?:string;command?:string}> }; + }; + }>; + const spec = entries.find(entry => entry.entry_type === 'run.spawned')?.payload?.spec; + const step = spec?.steps?.[0]; + const completed = entries.filter(entry => entry.entry_type === 'step.completed' && entry.step_id === claimed.id); + const terminalFacts = entries.filter(entry => entry.entry_type === 'run.completed'); + if (spec?.name !== `${metadata.flowName}/${claimed.id}` || spec?.steps?.length !== 1 + || step?.id !== claimed.id || completed.length !== 1 + || completed[0]?.payload?.completionReason !== 'success' + || terminalFacts.length !== 1 || terminalFacts[0]?.payload?.completionReason !== 'success') invalid(); + if (claimed === terminal) { + const command = result.completionReason === 'needs_human' + ? `printf '%s' '{"completionReason":"needs_human"}'` : ':'; + if (step?.type !== 'deterministic' || step.command !== command) invalid(); + } + } + } finally { journal.close(); } +} diff --git a/packages/sdk/src/authored-runtime-capability.ts b/packages/sdk/src/authored-runtime-capability.ts index 7297196a2..a70560ce0 100644 --- a/packages/sdk/src/authored-runtime-capability.ts +++ b/packages/sdk/src/authored-runtime-capability.ts @@ -1,3 +1,4 @@ +import { stripTypeScriptTypes, register } from 'node:module'; import { createHook } from 'node:async_hooks'; import { AuthoredFlowExecutionError } from './authored-flow-error.js'; @@ -23,7 +24,8 @@ export function assertAuthoredNodeVersion(version = process.versions.node): void const [major, minor] = version.split('.').map(Number); if (!Number.isInteger(major) || !Number.isInteger(minor) || major! < 22 || (major === 22 && minor! < 14) - || process.versions['bun'] !== undefined) { + || process.versions['bun'] !== undefined + || typeof stripTypeScriptTypes !== 'function' || typeof register !== 'function') { throw new AuthoredFlowExecutionError('unsupported_promise_lifecycle', 'the authored runner requires Node >=22.14; no workflow body was executed'); } diff --git a/packages/sdk/src/authored-source-authority.ts b/packages/sdk/src/authored-source-authority.ts index b0cf7f35f..3275cbc23 100644 --- a/packages/sdk/src/authored-source-authority.ts +++ b/packages/sdk/src/authored-source-authority.ts @@ -1,21 +1,27 @@ import { createHash } from 'node:crypto'; import { readFile } from 'node:fs/promises'; +import { register } from 'node:module'; +import { pathToFileURL } from 'node:url'; import { canonicalize } from './canonical.js'; import { loadAuthoredFlow } from './authored-flow-loader.js'; import type { AuthoredRootMetadata } from './authored-root.js'; const sha256 = (bytes: Uint8Array): string => createHash('sha256').update(bytes).digest('hex'); /** Re-open the exact source graph and Surface instance pinned by the durable root. */ -export async function loadPinnedAuthoredSource(metadata: AuthoredRootMetadata) { +export async function loadPinnedAuthoredSource(metadata: AuthoredRootMetadata, immutableNodeImport = false) { const source = await readFile(metadata.flowPath); if (sha256(source) !== metadata.sourceSha256) { throw new Error('authored root source authority mismatch'); } + const sources = new Map([[pathToFileURL(metadata.flowPath).href, source.toString('utf8')]]); for (const pinned of metadata.sources) { - if (sha256(await readFile(pinned.path)) !== pinned.sourceSha256) { + const bytes = await readFile(pinned.path); + if (sha256(bytes) !== pinned.sourceSha256) { throw new Error(`authored root source authority mismatch for "${pinned.path}"`); } + sources.set(pathToFileURL(pinned.path).href, bytes.toString('utf8')); } + if (immutableNodeImport) installPinnedSourceLoader(sources); const loaded = await loadAuthoredFlow(metadata.flowPath); if (canonicalize(loaded.surfaceAuthority) !== canonicalize(metadata.surface) || loaded.getDefinition(loaded.handle).name !== metadata.flowName) { @@ -31,3 +37,26 @@ export async function loadPinnedAuthoredSource(metadata: AuthoredRootMetadata) { } return loaded; } + +let pinnedSourceLoaderInstalled = false; + +/** Isolated Node child only: preserve original URLs/resolution, load verified bytes. */ +function installPinnedSourceLoader(sources: Map): void { + if (pinnedSourceLoaderInstalled) throw new Error('authored Node process already owns a source graph'); + pinnedSourceLoaderInstalled = true; + // Node >=22.14 supports register and stripTypeScriptTypes. A loader thread + // receives the captured bytes, never reopening mutable authored paths. + const loader = ` + import {stripTypeScriptTypes} from 'node:module'; + let sources; + export function initialize(data){sources=new Map(data);} + export async function load(url,context,next){ + if(sources.has(url))return {format:'module',shortCircuit:true, + source:stripTypeScriptTypes(sources.get(url),{mode:'transform',sourceUrl:url})}; + if(url.startsWith('file:') && new URL(url).pathname.endsWith('.flow.ts')) + throw new Error('authored root declared source graph authority mismatch'); + return next(url,context); + } + `; + register('data:text/javascript,' + encodeURIComponent(loader), { data: [...sources] }); +} diff --git a/packages/sdk/tests/authored-node-result.test.ts b/packages/sdk/tests/authored-node-result.test.ts new file mode 100644 index 000000000..aef419ff1 --- /dev/null +++ b/packages/sdk/tests/authored-node-result.test.ts @@ -0,0 +1,52 @@ +import { beforeEach, describe, expect, it, vi } from 'vitest'; +import { verifyAuthoredNodeResult } from '../src/authored-node-runner.js'; +import type { AuthoredFlowExecutionResult } from '../src/authored-flow-executor.js'; +import type { AuthoredRootMetadata } from '../src/authored-root.js'; + +const mocks=vi.hoisted(()=>({connect:vi.fn(),hello:vi.fn(),close:vi.fn(),runGet:vi.fn(),journalRead:vi.fn()})); +vi.mock('../src/journal-client.js',()=>({JournalClient:class { + connect=mocks.connect;hello=mocks.hello;close=mocks.close; + runGet=mocks.runGet;journalRead=mocks.journalRead; +}})); +const metadata={flowName:'example'} as AuthoredRootMetadata; +function result():AuthoredFlowExecutionResult { + return {rootRunId:'root',name:'example',completionReason:'success',journalSteps:[ + {id:'run-1',runId:'child-1',completionReason:'success'}, + {id:'complete-2',runId:'child-2',completionReason:'success'}, + ]}; +} +function entries(command=':', id='complete-2') { + return {entries:[ + {entry_type:'run.spawned',payload:{spec:{name:`example/${id}`,steps:[{id,type:'deterministic',command}]}}}, + {entry_type:'step.completed',step_id:id,payload:{completionReason:'success'}}, + {entry_type:'run.completed',payload:{completionReason:'success'}}, + ]}; +} +beforeEach(()=>{ + vi.clearAllMocks(); + mocks.runGet.mockImplementation(async (run_id:string)=>({run_id,status:'completed',steps:{ + [run_id==='child-1'?'run-1':'complete-2']:{state:'done'}, + }})); + mocks.journalRead.mockImplementation(async(runId:string)=>entries(':',runId==='child-1'?'run-1':'complete-2')); +}); +describe('authored IPC result durable verification',()=>{ + it('requires independently read completed children and exact terminal marker',async()=>{ + await verifyAuthoredNodeResult(result(),metadata,'root','socket'); + expect(mocks.runGet).toHaveBeenCalledTimes(2);expect(mocks.journalRead).toHaveBeenCalledWith('child-2',1); + expect(mocks.close).toHaveBeenCalledOnce(); + }); + it.each(['wrong-root','missing-terminal','duplicate-child','unfinished-child','wrong-terminal-command','missing-terminal-fact','omitted-child','unrelated-child','duplicate-terminal-fact'])( + 'refuses %s even when a result frame claims success',async kind=>{ + const claimed=result(); + if(kind==='omitted-child')Object.assign(claimed,{journalSteps:[claimed.journalSteps[1]]}); + if(kind==='unrelated-child')mocks.journalRead.mockResolvedValue(entries(':','other-1')); + if(kind==='duplicate-terminal-fact')mocks.journalRead.mockImplementation(async(runId:string)=>runId==='child-1'?entries(':','run-1'):{entries:[...entries().entries,entries().entries[2]]}); + if(kind==='wrong-root')Object.assign(claimed,{rootRunId:'other'}); + if(kind==='missing-terminal')Object.assign(claimed,{journalSteps:claimed.journalSteps.slice(0,1)}); + if(kind==='duplicate-child')Object.assign(claimed,{journalSteps:[...claimed.journalSteps,claimed.journalSteps[1]]}); + if(kind==='unfinished-child')mocks.runGet.mockResolvedValue({run_id:'child-1',status:'running',steps:{'run-1':{state:'running'}}}); + if(kind==='wrong-terminal-command')mocks.journalRead.mockImplementation(async(runId:string)=>entries(runId==='child-2'?'echo forged':':',runId==='child-1'?'run-1':'complete-2')); + if(kind==='missing-terminal-fact')mocks.journalRead.mockImplementation(async(runId:string)=>runId==='child-1'?entries(':','run-1'):{entries:entries().entries.slice(0,-1)}); + await expect(verifyAuthoredNodeResult(claimed,metadata,'root','socket')).rejects.toThrow('no matching durable completion'); + }); +}); diff --git a/packages/sdk/tests/authored-node-runtime.test.ts b/packages/sdk/tests/authored-node-runtime.test.ts index 1467c1a88..9c0b6c628 100644 --- a/packages/sdk/tests/authored-node-runtime.test.ts +++ b/packages/sdk/tests/authored-node-runtime.test.ts @@ -1,5 +1,5 @@ import { spawn, spawnSync } from 'node:child_process'; -import { chmodSync, cpSync, existsSync, mkdtempSync, readFileSync, readdirSync, rmSync, writeFileSync } from 'node:fs'; +import { chmodSync, cpSync, existsSync, lstatSync, mkdtempSync, readFileSync, readdirSync, rmSync, symlinkSync, writeFileSync } from 'node:fs'; import { tmpdir } from 'node:os'; import { join, resolve } from 'node:path'; import { afterAll, beforeAll, describe, expect, it } from 'vitest'; @@ -43,7 +43,10 @@ afterAll(() => { function fixture(body: string) { const directory = mkdtempSync(join(tmpdir(), 'authored-node-runtime-')); fixtures.push(directory); const surface = join(directory, 'node_modules/@relayflows/surface'); - cpSync(join(sdk, 'node_modules/@relayflows/surface'), surface, { recursive: true }); + const surfaceLink = join(directory, 'surface-link'); + symlinkSync(join(sdk, 'node_modules/@relayflows/surface'), surfaceLink, 'dir'); + cpSync(surfaceLink, surface, { recursive: true, dereference: true }); + expect(lstatSync(surface).isSymbolicLink(), 'fixture must never rewrite the shared Surface symlink').toBe(false); const manifest = JSON.parse(readFileSync(join(surface, 'package.json'), 'utf8')); // Same verified execution envelope Cloud uses around the unchanged Surface payload. writeFileSync(join(surface, 'package.json'), JSON.stringify({ name: manifest.name, @@ -64,7 +67,7 @@ if(request){appendFileSync('agent-effects','once\\n');await new Promise(r=>setTi chmodSync(wrapper, 0o755); writeFileSync(join(directory, 'flows.json'), JSON.stringify({ cli: wrapper })); writeFileSync(join(directory, 'case.flow.ts'), `import {flow} from '@relayflows/surface'; -import {appendFileSync,existsSync,writeFileSync} from 'node:fs'; +import {appendFileSync,existsSync,writeFileSync,writeSync} from 'node:fs'; export default flow('runtime-case',async f=>{${body}}); `); const env = { ...process.env, FLOWS_AUTHORED_NODE: process.execPath, @@ -106,12 +109,12 @@ describe('Bun 1.4.0 standalone → native Node authored lifecycle', () => { expect(output.journalSteps.map((s:{id:string})=>s.id)).toEqual(['agent-1','run-2','run-3','run-4','complete-5']); }, 90_000); - it.each(['SIGKILL', 'SIGTERM'] as const)('stops on parent %s and replays completed children under the same unfinished root', async signal => { + it.each(['SIGKILL', 'SIGTERM', 'blocked-SIGKILL'] as const)('stops on parent %s and replays completed children under the same unfinished root', async signal => { const f=fixture(sequential + ` if(!existsSync('resume-ready')){ writeFileSync('node-pid',String(process.pid)); writeFileSync('resume-ready','yes'); - await new Promise(()=>{}); + ${signal === 'blocked-SIGKILL' ? 'while(true){}' : 'await new Promise(()=>{});'} } f.done('success');`); const child=spawn(cli,['run','case.flow.ts','--input','{}',...f.flags],{ @@ -120,6 +123,7 @@ f.done('success');`); let logs='';child.stdout.on('data',bytes=>{logs+=bytes});child.stderr.on('data',bytes=>{logs+=bytes}); const closed=new Promise(resolve=>child.once('close',()=>resolve())); let nodePid:number|undefined; + let bodySucceeded = false; try { const deadline=Date.now()+15_000; while(!existsSync(join(f.directory,'resume-ready')) && Date.now()e.entry_type==='step.attempt.started')).toHaveLength(2); expect(rootJournal.filter(e=>e.entry_type==='run.completed')).toHaveLength(1); expect(rootJournal.at(-1)?.payload['completionReason']).toBe('success'); + bodySucceeded = true; } finally { child.kill('SIGKILL'); - if(nodePid!==undefined){try{process.kill(nodePid,'SIGKILL');}catch(error){if((error as NodeJS.ErrnoException).code!=='ESRCH')throw error;}} + if(nodePid!==undefined){try{process.kill(nodePid,'SIGKILL');}catch(error){if(bodySucceeded && (error as NodeJS.ErrnoException).code!=='ESRCH')throw error;}} } },60_000); @@ -163,9 +168,31 @@ f.done('success');`); ])('refuses %s rather than reporting terminal success', (_name,body) => { const f=fixture(body);const result=f.run();expect(result.status).toBe(1); expect(result.stdout+result.stderr).toContain('unawaited_step'); + expect(result.stdout+result.stderr).not.toContain('unawaited_step: unawaited_step:'); expect(existsSync(join(f.directory,'forbidden-effects'))).toBe(false); }, 60_000); + it('loads captured graph bytes before preserving the unsupported-use refusal', () => { + const f=fixture(`f.done('success');`); + const original=`import {flow} from '@relayflows/surface';import {writeFileSync} from 'node:fs';if(!process.versions.bun)writeFileSync('loaded-source','original');export default flow('child',async f=>{f.done('success')});`; + const modified=original.replace("'original'", "'modified'"); + writeFileSync(join(f.directory,'child.flow.ts'),original); + const root=readFileSync(join(f.directory,'case.flow.ts'),'utf8') + .replace("flow('runtime-case',async", "flow('runtime-case',{use:['./child.flow.ts']},async"); + writeFileSync(join(f.directory,'case.flow.ts'), + `if(!process.versions.bun)writeFileSync('child.flow.ts',${JSON.stringify(modified)});\n`+root); + const result=f.run();expect(result.status,result.stderr+result.stdout).toBe(2); + expect(result.stderr+result.stdout).toContain('unsupported_header'); + expect(readFileSync(join(f.directory,'child.flow.ts'),'utf8')).toBe(modified); + expect(readFileSync(join(f.directory,'loaded-source'),'utf8')).toBe('original'); + },30_000); + + it('rejects a forged result frame without durable completion', () => { + const f=fixture(`writeSync(3,JSON.stringify({type:'result',result:{name:'runtime-case',completionReason:'success',journalSteps:[]}})+'\\n');process.exit(0);`); + const result=f.run();expect(result.status).toBe(1); + expect(result.stdout+result.stderr).toContain('invalid authored runtime authenticated frame'); + },30_000); + it('refuses missing Node before body effects or root admission', () => { const f=fixture(`writeFileSync('body-started','bad');${sequential}f.done('success');`); const result=f.invoke(['run','case.flow.ts','--input','{}'],{FLOWS_AUTHORED_NODE:join(f.directory,'missing-node')}); From 51fb0345a8c6c68377c5407750e546587a0401a3 Mon Sep 17 00:00:00 2001 From: Miya Date: Tue, 15 Sep 2026 19:58:00 +0200 Subject: [PATCH 5/6] fix: fail closed when authored parent watchdog exits Session-Id: 01a09c40-ce3b-7f11-a7df-b6b7ccab6fd9 --- packages/sdk/src/authored-node-entry.ts | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/packages/sdk/src/authored-node-entry.ts b/packages/sdk/src/authored-node-entry.ts index d01114a88..484f689cc 100644 --- a/packages/sdk/src/authored-node-entry.ts +++ b/packages/sdk/src/authored-node-entry.ts @@ -38,6 +38,10 @@ const watchdog = new Worker(` function check(){if(process.ppid!==workerData.parentPid)process.kill(process.pid,'SIGKILL');} check();setInterval(check,50);parentPort.postMessage('ready'); `, { eval: true, workerData: { parentPid } }); +// Enforcement must not silently disappear after the readiness Promise settles. +const watchdogLost = (): void => { if (!finished) process.kill(process.pid, 'SIGKILL'); }; +watchdog.on('error', watchdogLost); +watchdog.on('exit', watchdogLost); let client: JournalClient | undefined; try { assertAuthoredNodeVersion(); From d379cb1f0b113ef5e7ce383b081fa1c19ebea295 Mon Sep 17 00:00:00 2001 From: Miya Date: Tue, 15 Sep 2026 20:11:52 +0200 Subject: [PATCH 6/6] fix: verify complete child journals and allow container init parent Session-Id: 01a09c40-ce3b-7f11-a7df-b6b7ccab6fd9 --- packages/sdk/src/authored-node-entry.ts | 5 +-- packages/sdk/src/authored-node-runner.ts | 18 ++++++-- .../sdk/src/authored-runtime-capability.ts | 7 +++ .../sdk/tests/authored-node-result.test.ts | 43 +++++++++++++++---- 4 files changed, 58 insertions(+), 15 deletions(-) diff --git a/packages/sdk/src/authored-node-entry.ts b/packages/sdk/src/authored-node-entry.ts index 484f689cc..3259d697f 100644 --- a/packages/sdk/src/authored-node-entry.ts +++ b/packages/sdk/src/authored-node-entry.ts @@ -4,7 +4,7 @@ import { readFileSync, writeSync } from 'node:fs'; import { JournalClient } from './journal-client.js'; import { executeAuthoredFlow } from './authored-flow-executor.js'; import { loadPinnedAuthoredSource } from './authored-source-authority.js'; -import { assertAuthoredNodeVersion } from './authored-runtime-capability.js'; +import { assertAuthoredNodeVersion, parseAuthoredParentPid } from './authored-runtime-capability.js'; import { AuthoredFlowExecutionError } from './authored-flow-error.js'; import type { AuthoredRootMetadata } from './authored-root.js'; @@ -31,8 +31,7 @@ process.stdin.on('end', () => { if (!finished) process.exit(1); }); process.stdin.on('error', () => process.exit(1)); // A separate event loop enforces parent loss even while authored JS is blocked. // Capture the expected PID in the parent's spawn arguments, before any child work. -const parentPid = Number(process.argv[2]); -if (!Number.isSafeInteger(parentPid) || parentPid <= 1) throw new Error('invalid authored parent identity'); +const parentPid = parseAuthoredParentPid(process.argv[2]); const watchdog = new Worker(` const {parentPort,workerData}=require('node:worker_threads'); function check(){if(process.ppid!==workerData.parentPid)process.kill(process.pid,'SIGKILL');} diff --git a/packages/sdk/src/authored-node-runner.ts b/packages/sdk/src/authored-node-runner.ts index 97d7dfb48..05ce203d0 100644 --- a/packages/sdk/src/authored-node-runner.ts +++ b/packages/sdk/src/authored-node-runner.ts @@ -184,11 +184,23 @@ export async function verifyAuthoredNodeResult( const state = await journal.runGet(claimed.runId); if (state.run_id !== claimed.runId || state.status !== 'completed' || state.steps[claimed.id]?.state !== 'done') invalid(); - const entries = (await journal.journalRead(claimed.runId, 1)).entries as Array<{ - entry_type: string; step_id?: string; payload?: { + const entries: Array<{ + seq: number; entry_type: string; step_id?: string; payload?: { completionReason?: string; spec?: { name?: string; steps?: Array<{id?:string;type?:string;command?:string}> }; }; - }>; + }> = []; + let fromSeq = 1; + for (;;) { + const page = (await journal.journalRead(claimed.runId, fromSeq, 100)).entries; + if (page.length === 0) break; + for (const raw of page) { + const entry = raw as typeof entries[number]; + if (!Number.isSafeInteger(entry?.seq) || entry.seq < fromSeq) invalid(); + fromSeq = entry.seq + 1; + // Keep only completion evidence; streaming logs can span many pages. + if (['run.spawned', 'step.completed', 'run.completed'].includes(entry.entry_type)) entries.push(entry); + } + } const spec = entries.find(entry => entry.entry_type === 'run.spawned')?.payload?.spec; const step = spec?.steps?.[0]; const completed = entries.filter(entry => entry.entry_type === 'step.completed' && entry.step_id === claimed.id); diff --git a/packages/sdk/src/authored-runtime-capability.ts b/packages/sdk/src/authored-runtime-capability.ts index a70560ce0..201068877 100644 --- a/packages/sdk/src/authored-runtime-capability.ts +++ b/packages/sdk/src/authored-runtime-capability.ts @@ -31,3 +31,10 @@ export function assertAuthoredNodeVersion(version = process.versions.node): void } assertAuthoredPromiseHooks(); } + +/** PID 1 is valid when the standalone CLI is a container's init process. */ +export function parseAuthoredParentPid(value: string | undefined): number { + const pid = Number(value); + if (!Number.isSafeInteger(pid) || pid <= 0) throw new Error('invalid authored parent identity'); + return pid; +} diff --git a/packages/sdk/tests/authored-node-result.test.ts b/packages/sdk/tests/authored-node-result.test.ts index aef419ff1..752c4d71d 100644 --- a/packages/sdk/tests/authored-node-result.test.ts +++ b/packages/sdk/tests/authored-node-result.test.ts @@ -1,5 +1,6 @@ import { beforeEach, describe, expect, it, vi } from 'vitest'; import { verifyAuthoredNodeResult } from '../src/authored-node-runner.js'; +import { parseAuthoredParentPid } from '../src/authored-runtime-capability.js'; import type { AuthoredFlowExecutionResult } from '../src/authored-flow-executor.js'; import type { AuthoredRootMetadata } from '../src/authored-root.js'; @@ -15,38 +16,62 @@ function result():AuthoredFlowExecutionResult { {id:'complete-2',runId:'child-2',completionReason:'success'}, ]}; } -function entries(command=':', id='complete-2') { - return {entries:[ +function entries(command=':', id='complete-2'):Array> { + return [ {entry_type:'run.spawned',payload:{spec:{name:`example/${id}`,steps:[{id,type:'deterministic',command}]}}}, {entry_type:'step.completed',step_id:id,payload:{completionReason:'success'}}, {entry_type:'run.completed',payload:{completionReason:'success'}}, - ]}; + ]; } +let records:Map>>; beforeEach(()=>{ vi.clearAllMocks(); + records=new Map([['child-1',entries(':','run-1')],['child-2',entries()]]); mocks.runGet.mockImplementation(async (run_id:string)=>({run_id,status:'completed',steps:{ [run_id==='child-1'?'run-1':'complete-2']:{state:'done'}, }})); - mocks.journalRead.mockImplementation(async(runId:string)=>entries(':',runId==='child-1'?'run-1':'complete-2')); + mocks.journalRead.mockImplementation(async(runId:string,fromSeq:number,limit=100)=>({ + entries:records.get(runId)!.map((entry,index)=>({...entry,seq:index+1})) + .filter(entry=>entry.seq>=fromSeq).slice(0,limit), + })); }); describe('authored IPC result durable verification',()=>{ it('requires independently read completed children and exact terminal marker',async()=>{ await verifyAuthoredNodeResult(result(),metadata,'root','socket'); - expect(mocks.runGet).toHaveBeenCalledTimes(2);expect(mocks.journalRead).toHaveBeenCalledWith('child-2',1); + expect(mocks.runGet).toHaveBeenCalledTimes(2);expect(mocks.journalRead).toHaveBeenCalledWith('child-2',1,100); expect(mocks.close).toHaveBeenCalledOnce(); }); + it('reads successful terminal evidence beyond 100-entry journal pages',async()=>{ + const original=records.get('child-1')!; + records.set('child-1',[original[0]!,...Array.from({length:248},()=>({entry_type:'worker.stream'})),...original.slice(1)]); + await verifyAuthoredNodeResult(result(),metadata,'root','socket'); + expect(mocks.journalRead).toHaveBeenCalledWith('child-1',101,100); + expect(mocks.journalRead).toHaveBeenCalledWith('child-1',201,100); + expect(mocks.journalRead).toHaveBeenCalledWith('child-1',252,100); + }); + it('refuses a stalled or backwards journal page instead of looping',async()=>{ + mocks.journalRead.mockResolvedValue({entries:[{seq:1,entry_type:'worker.stream'}]}); + await expect(verifyAuthoredNodeResult(result(),metadata,'root','socket')).rejects.toThrow('no matching durable completion'); + expect(mocks.journalRead).toHaveBeenCalledTimes(2); + }); it.each(['wrong-root','missing-terminal','duplicate-child','unfinished-child','wrong-terminal-command','missing-terminal-fact','omitted-child','unrelated-child','duplicate-terminal-fact'])( 'refuses %s even when a result frame claims success',async kind=>{ const claimed=result(); if(kind==='omitted-child')Object.assign(claimed,{journalSteps:[claimed.journalSteps[1]]}); - if(kind==='unrelated-child')mocks.journalRead.mockResolvedValue(entries(':','other-1')); - if(kind==='duplicate-terminal-fact')mocks.journalRead.mockImplementation(async(runId:string)=>runId==='child-1'?entries(':','run-1'):{entries:[...entries().entries,entries().entries[2]]}); + if(kind==='unrelated-child')records.set('child-1',entries(':','other-1')); + if(kind==='duplicate-terminal-fact')records.set('child-2',[...entries(),entries()[2]!]); if(kind==='wrong-root')Object.assign(claimed,{rootRunId:'other'}); if(kind==='missing-terminal')Object.assign(claimed,{journalSteps:claimed.journalSteps.slice(0,1)}); if(kind==='duplicate-child')Object.assign(claimed,{journalSteps:[...claimed.journalSteps,claimed.journalSteps[1]]}); if(kind==='unfinished-child')mocks.runGet.mockResolvedValue({run_id:'child-1',status:'running',steps:{'run-1':{state:'running'}}}); - if(kind==='wrong-terminal-command')mocks.journalRead.mockImplementation(async(runId:string)=>entries(runId==='child-2'?'echo forged':':',runId==='child-1'?'run-1':'complete-2')); - if(kind==='missing-terminal-fact')mocks.journalRead.mockImplementation(async(runId:string)=>runId==='child-1'?entries(':','run-1'):{entries:entries().entries.slice(0,-1)}); + if(kind==='wrong-terminal-command')records.set('child-2',entries('echo forged')); + if(kind==='missing-terminal-fact')records.set('child-2',entries().slice(0,-1)); await expect(verifyAuthoredNodeResult(claimed,metadata,'root','socket')).rejects.toThrow('no matching durable completion'); }); }); +describe('authored parent identity',()=>{ + it('accepts a container PID 1 parent',()=>expect(parseAuthoredParentPid('1')).toBe(1)); + it.each([undefined,'0','-1','1.5','NaN','9007199254740992'])('refuses invalid identity %s',value=>{ + expect(()=>parseAuthoredParentPid(value)).toThrow('invalid authored parent identity'); + }); +});