Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 2 additions & 4 deletions .github/workflows/cloud-runtime-artifact.yml
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ on:
- "kernel/**"
- "packages/sdk/**"
- "scripts/cloud-artifact.mjs"
- "scripts/build-standalone-cli.mjs"
- "scripts/cloud-artifact.test.mjs"
- "testdata/**"

Expand Down Expand Up @@ -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:
Expand Down
6 changes: 2 additions & 4 deletions .github/workflows/publish.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
28 changes: 28 additions & 0 deletions docs/SURFACE.md
Original file line number Diff line number Diff line change
Expand Up @@ -803,3 +803,31 @@ 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. 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.
10 changes: 10 additions & 0 deletions packages/sdk/src/authored-flow-executor.ts
Original file line number Diff line number Diff line change
@@ -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';
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -133,6 +142,7 @@ export async function executeAuthoredFlow<Input = undefined>(
input?: Input,
options: ExecuteAuthoredFlowOptions = {},
): Promise<AuthoredFlowExecutionResult> {
assertAuthoredPromiseHooks();
const getDefinition = options.getDefinition ?? getAuthoredFlowDefinition;
const localAgentStream = options.localAgentStream;
const onProgress = options.onProgress;
Expand Down
90 changes: 90 additions & 0 deletions packages/sdk/src/authored-node-entry.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,90 @@
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';
import { loadPinnedAuthoredSource } from './authored-source-authority.js';
import { assertAuthoredNodeVersion, parseAuthoredParentPid } from './authored-runtime-capability.js';
import { AuthoredFlowExecutionError } from './authored-flow-error.js';
import type { AuthoredRootMetadata } from './authored-root.js';

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;
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); });
Comment thread
coderabbitai[bot] marked this conversation as resolved.
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 = 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');}
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();
await new Promise<void>((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<{ channelKey: string; 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);
});
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, 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');
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) {
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; await watchdog.terminate(); client?.close(); process.stdin.destroy();
process.off('SIGINT', abort); process.off('SIGTERM', abort);
}
Loading
Loading