diff --git a/docs/BABYSITTER-CATALOG-HANDOFF.md b/docs/BABYSITTER-CATALOG-HANDOFF.md index 414493cb6..485edec47 100644 --- a/docs/BABYSITTER-CATALOG-HANDOFF.md +++ b/docs/BABYSITTER-CATALOG-HANDOFF.md @@ -13,12 +13,20 @@ The reviewed native artifact shipped in Flows 2.0.26. The existing `examples/babysitter` declares Claude, GitHub comment writes, and a merge-gate hook; it is not the native existing-session package. Ordinary authored dispatch still refuses matched extension handlers with `plugin_unsupported`, because it -imports tenant JavaScript in the host. The external exact-target Relay/Flows -runtime must instead call `runHostedSoftwareGardenBabysitter`, which captures -the reviewed Software Garden base and complete lock-backed installation as one -generation before it runs the exact matched native handler in the capability -sandbox. Keep enabled activation blocked with zero writes until that runtime -calls the entrypoint. +imports tenant JavaScript in the host. An owning exact-target Relay/Flows action +must instead use the normal embedded `runCli(["run", flowPath, "--input", ...])` +surface and inject `RunCliOptions.hostedSoftwareGardenBabysitter`. That +non-serializable option carries the host-verified dispatch and the single queue +capability; no CLI flag or flow input can mint either. The canonical authored +run preflights the reviewed Software Garden base and complete lock-backed +installation as one generation before trigger inspection, tenant import, or +daemon attachment. It then admits one journaled effect step under a delivery- +and-pin-bound key and runs the exact matched native handler in the capability +sandbox. The queue write uses the journal's record/perform/confirm protocol; +the returned run ID, terminal reason, and completed-step count come from that +journal rather than an in-memory synthetic result. Every other command or path +surface refuses the hosted authority. Keep enabled activation blocked with zero +writes until that action is released and deployed. The native handler must consume host-verified delivery authority, normalize the repository/PR event, and call the Cloud lineage path. Cloud must recheck the @@ -31,8 +39,9 @@ through the generic executor: #549 still refuses it. The SDK has a separate Linu capability sandbox that injects exactly `capabilities.cloud.babysitterTurn.queue` without exposing the base context, workspace, environment credentials, network, helpers, MCP, or harnesses. It is -reached by the canonical composition entrypoint, but no deployed runtime calls -that entrypoint and this must not be treated as enablement. The package's +reached by the canonical authored `run` surface when the owning hosted action +injects verified authority, but no deployed runtime consumes this contract and +this must not be treated as enablement. The package's `compat` requires the published 2.0.26 Surface/SDK release that routes `labeled`, `unlabeled`, and `ready_for_review`. Export it only from the reviewed release commit pinned below. The Software Factory flow's own independently @@ -59,11 +68,12 @@ normalized input against non-serializable verified dispatch authority, and permits one queue call. The parent capability adapter receives that original authority plus immutable extension provenance; the capability request never carries workspace, activation, listener, session, -lineage, label, head, prompt, merge, route, or config authority. Cloud PR #4002 -at `25412782bf148ff8dd8018bdbafd719d4e8347fa` deliberately supplies no execution +lineage, label, head, prompt, merge, route, or config authority. Merged Cloud +PR #4002 at merge commit `ced414ab40424c7bbd4cd780ad01751d7fc85685` +(reviewed head `6201470228b23c225290d0eee356eb1c0006e31d`) deliberately supplies no execution authority: it emits the exact-target `relay:hosted-flow-extension:v1` action for -an external Relay/Flows runtime. That runtime owns this entrypoint; only its -validated capability call reaches `relay:native-existing-session:v1` +an external Relay/Flows runtime. That runtime owns the embedded hosted-run +option; only its validated capability call reaches `relay:native-existing-session:v1` downstream. Merged Cloud PR #3942 owns that downstream lineage/authority core and must inject workspace, activation, and listener from persisted dispatch context, re-read live PR/label/head state, and return only `{ receiptId, status: @@ -80,14 +90,19 @@ the branded dispatch, and waits for the adapter's authoritative outcome before settling any premature child terminal frame. An authoritative adapter rejection settles immediately with its original typed error even if the child hangs. -Before replacing #549's refusal, the hosted caller must obtain an opaque base -and installation as one generation with `loadHostedExtensionRuntime`, then call -`runHostedCapabilityExtension` with both values. Every dispatch rechecks the -current extension declarations and complete project source tree against that -generation. Directory entries are streamed beneath a shared entry bound; +The normal embedded run deliberately does not replace #549's standalone +refusal. Its hosted option calls the canonical composition boundary, which +obtains an opaque base and installation as one generation with +`loadHostedExtensionRuntime`, then calls `runHostedCapabilityExtension` with +both values from its journal-attached worker. Unsupported local-worker flags +are refused instead of ignored. Every dispatch rechecks the current extension +declarations and complete project source tree against that generation. +Directory entries are streamed beneath a shared entry bound; nonblocking no-follow descriptors and explicitly bounded reads enforce the -cumulative-byte limit before source contents are buffered. The loader never imports -tenant base code to derive authority. It +cumulative-byte limit before source contents are buffered. Only runtime/control +directories (`.flows`, `.git`, `.relayflowd`, and `node_modules`) are excluded, +so starting the standard journal daemon cannot invalidate the generation it is +executing. The loader never imports tenant base code to derive authority. It requires the exact reviewed Software Factory flow-file SHA-256 and assigns its pinned name/version in the parent; project `node_modules`, relative imports, stdout, process termination, globals, and module caches therefore cannot forge diff --git a/packages/sdk/src/cli.ts b/packages/sdk/src/cli.ts index e5be59c1e..2888bc7e3 100644 --- a/packages/sdk/src/cli.ts +++ b/packages/sdk/src/cli.ts @@ -30,7 +30,10 @@ import { import { answerFlow } from './cli/answer.js'; import { checkAuthoredTriggers } from './cli/check-triggers.js'; import { parseWebhookArgs, runServeWebhook } from './cli/serve-webhook.js'; -import { runDirectFlow } from './cli/direct-run.js'; +import { + runDirectFlow, + type HostedSoftwareGardenRunOptions, +} from './cli/direct-run.js'; import { parseReplayArgs, replayJournal, type ReplayArgs } from './cli/replay.js'; import { parseStatusArgs, runStatus, type StatusArgs } from './cli/status.js'; import { @@ -182,6 +185,13 @@ export interface RunCliOptions { * SIGINT/SIGTERM are handled here, for the duration of that verb only. */ signal?: AbortSignal; + + /** + * Verified delivery authority and the only capability exposed to the + * canonical hosted Software Garden + Babysitter run. The standalone binary + * never constructs this option; an owning hosted action must inject it. + */ + hostedSoftwareGardenBabysitter?: HostedSoftwareGardenRunOptions; } /** @@ -213,11 +223,13 @@ export async function runCli( io: CliIo = PROCESS_IO, options: RunCliOptions = {}, ): Promise { - if (args.length === 1 && (args[0] === '--version' || args[0] === '-V')) { + if (options.hostedSoftwareGardenBabysitter === undefined + && args.length === 1 && (args[0] === '--version' || args[0] === '-V')) { io.stdout(options.version ?? packageVersion()); return 0; } - if (args.length === 1 && (args[0] === '--help' || args[0] === '-h')) { + if (options.hostedSoftwareGardenBabysitter === undefined + && args.length === 1 && (args[0] === '--help' || args[0] === '-h')) { io.stdout(USAGE); return 0; } @@ -229,6 +241,33 @@ export async function runCli( return 2; } + if (options.hostedSoftwareGardenBabysitter !== undefined + && (parsed.command !== 'run' || !isAuthoredFlowPath(parsed.value))) { + const report = inputFailureReport({ + kind: 'invalid_invocation', + message: 'Hosted Software Garden authority is accepted only by an authored flow run.', + }, 'value' in parsed && typeof parsed.value === 'string' ? parsed.value : undefined); + emitCheckReport(report, 'json' in parsed && parsed.json === true, io); + return 2; + } + if (options.hostedSoftwareGardenBabysitter !== undefined && parsed.command === 'run') { + const ignoredFlag = parsed.localAgent + ? '--local-agent' + : parsed.agentCapacity !== undefined + ? '--agent-capacity' + : parsed.allowHumanInfluenced + ? '--allow-human-influenced' + : undefined; + if (ignoredFlag !== undefined) { + const report = inputFailureReport({ + kind: 'invalid_invocation', + message: `${ignoredFlag} is not supported by a hosted Software Garden run.`, + }, parsed.value); + emitCheckReport(report, parsed.json, io); + return 2; + } + } + if (parsed.command === 'add') return addPlugin(parsed.value, io); if (parsed.command === 'plugin') return runPluginCommand(parsed, io); @@ -389,6 +428,9 @@ export async function runCli( elapsedMs: now - startedSteps.get(progress.stepId)! }); }, daemon: { spawn: parsed.spawn && spawnAllowedByEnv() }, + ...(options.hostedSoftwareGardenBabysitter === undefined ? {} : { + hostedSoftwareGardenBabysitter: options.hostedSoftwareGardenBabysitter, + }), }; const execution = parsed.command === 'run' ? isAuthoredFlowPath(parsed.value) diff --git a/packages/sdk/src/cli/direct-run.ts b/packages/sdk/src/cli/direct-run.ts index 2d1fcaf0a..a08f9c779 100644 --- a/packages/sdk/src/cli/direct-run.ts +++ b/packages/sdk/src/cli/direct-run.ts @@ -16,6 +16,10 @@ import { DirectInputError, parseDirectInput } from '../direct-input.js'; import { JournalClient } from '../journal-client.js'; import { inputFailureReport } from './check.js'; import { checkAuthoredTriggers } from './check-triggers.js'; +import { + runHostedSoftwareGardenFlow, + type HostedSoftwareGardenRunOptions, +} from './hosted-software-garden-run.js'; import { authoredInput, authoredWorkerRemedy, localAgentRemedy } from './local-agent-remedy.js'; import { authoredCompletion, @@ -31,11 +35,22 @@ import { type RunReport, } from './run.js'; +export type { HostedSoftwareGardenRunOptions } from './hosted-software-garden-run.js'; + +export interface RunDirectFlowOptions extends RunLifecycleOptions { + /** + * Host-verified authority for the canonical Software Garden + Babysitter + * delivery path. This is deliberately an in-process option: neither flow + * input nor CLI flags can mint the branded dispatch or queue capability. + */ + hostedSoftwareGardenBabysitter?: HostedSoftwareGardenRunOptions; +} + export async function runDirectFlow( path: string, inputArgument: string | undefined, dataDir: string, - options: RunLifecycleOptions = {}, + options: RunDirectFlowOptions = {}, ): Promise { let input: unknown; try { @@ -51,6 +66,17 @@ export async function runDirectFlow( }; } + // A hosted Software Garden delivery is still a normal authored `run`, but + // it must branch before the ordinary trigger checker imports tenant code or + // a daemon is attached. The caller supplies only authority that was minted + // from its verified delivery and its exact queue capability; the loader + // independently resolves the reviewed base plus installed, lock-backed + // Babysitter generation and the sandbox selects the exact matched handler. + const hostedSoftwareGarden = options.hostedSoftwareGardenBabysitter; + if (hostedSoftwareGarden !== undefined) { + return runHostedSoftwareGardenFlow(path, input, dataDir, hostedSoftwareGarden, options); + } + // Declared triggers are knowable before any daemon or step is started. // Importing the authored module is unavoidable here — trigger sources // are only observable after `flow(...).on(webhook(...))` has run — but diff --git a/packages/sdk/src/cli/hosted-software-garden-run.ts b/packages/sdk/src/cli/hosted-software-garden-run.ts new file mode 100644 index 000000000..b124db278 --- /dev/null +++ b/packages/sdk/src/cli/hosted-software-garden-run.ts @@ -0,0 +1,351 @@ +import { createHash, randomUUID } from 'node:crypto'; +import { join } from 'node:path'; +import { compileSpec, toKernelSpec } from '../compile.js'; +import { runtimeVersions } from '../flow-extension-compat.js'; +import { + hostedExtensionDispatchIdentity, + type HostedEventIdentity, +} from '../flow-extension-loader.js'; +import { + loadHostedExtensionRuntime, + runHostedCapabilityExtension, + selectHostedExtensionForRuntime, + type RunHostedSoftwareGardenBabysitterOptions, +} from '../hosted-extension-isolation.js'; +import { atomicJson, readHelperReceipt } from '../helper-storage.js'; +import { JournalClient } from '../journal-client.js'; +import { PluginError } from '../plugin-manifest.js'; +import type { RunOutcome, StepDispatchEvent } from '../protocol.js'; +import { SPEC_SCHEMA_VERSION } from '../spec.js'; +import { withWorkerLease } from '../worker-lease.js'; +import { + classifyOutcome, + connect, + emptyReport, + protocolFailure, + socketFor, + type RunExecution, + type RunLifecycleOptions, + type RunReport, +} from './run.js'; + +const STEP_ID = 'babysitter-turn'; +const SURFACE_PATH = '/cloud/babysitter-turn'; + +export type HostedSoftwareGardenRunOptions = Omit< + RunHostedSoftwareGardenBabysitterOptions, + 'flowPath' | 'input' | 'bubblewrapPath' | 'nodePath' | 'prlimitPath' +>; + +/** + * Execute the pinned hosted composition as one journaled effect step. Plugin, + * pin, route, and base refusals are resolved before a run exists; + * every error after admission is a terminal step failure because the external + * queue outcome may already be in doubt. + */ +export async function runHostedSoftwareGardenFlow( + path: string, + input: unknown, + dataDir: string, + hosted: HostedSoftwareGardenRunOptions, + lifecycle: RunLifecycleOptions, +): Promise { + const base: RunReport = { ...emptyReport('run'), path }; + let identity: HostedEventIdentity; + let runtime: Awaited>; + let selected: Awaited>; + try { + identity = hostedExtensionDispatchIdentity(hosted.dispatch); + runtime = await loadHostedExtensionRuntime(path); + selected = await selectHostedExtensionForRuntime( + runtime.installation, + runtime.base, + identity, + runtimeVersions(), + ); + } catch (error) { + if (error instanceof PluginError) return pluginRefusal(base, error); + return hostedFailure(base, error); + } + + const socketPath = socketFor(dataDir); + const client = new JournalClient(socketPath); + const connected = await connect(client, 'run', dataDir, base, lifecycle); + if (connected !== undefined) return connected; + + const admissionIdentity = { + provider: identity.provider, + eventType: hosted.dispatch.eventType, + deliveryId: hosted.dispatch.deliveryId, + base: runtime.base, + extension: { + name: selected.manifest.name, + version: selected.manifest.version, + ref: selected.artifact.ref, + digest: selected.artifact.digest, + manifestSha256: selected.artifact.manifestSha256, + }, + }; + const admissionDigest = createHash('sha256') + .update(JSON.stringify(admissionIdentity)) + .digest('hex'); + const stream = `hosted-babysitter-${admissionDigest}`; + const instruction = JSON.stringify({ + type: 'effect', + provider: 'cloud', + verb: 'babysitter-turn', + identity, + authority: admissionIdentity, + input, + }); + const spec = toKernelSpec(compileSpec({ + version: SPEC_SCHEMA_VERSION, + name: 'software-factory/hosted-babysitter', + steps: [{ + id: STEP_ID, + type: 'agent', + instruction, + maxIterations: 1, + recoveryMode: 'reset', + surfaces: { + streams: [{ stream }], + external: [SURFACE_PATH], + }, + }], + })); + const peer = client.createPeer(); + let work: Promise | undefined; + let completedWork: HostedCompletion | undefined; + let capabilityFailed = false; + let capabilityFailure: unknown; + let workerFailure: unknown; + + const dispatch = (event: StepDispatchEvent): void => { + const dispatched = event.spec as { instruction?: string; surfaces?: { streams?: { stream: string }[] } }; + if (work !== undefined || event.step_id !== STEP_ID || event.step_type !== 'agent' + || dispatched.instruction !== instruction + || !dispatched.surfaces?.streams?.some(pin => pin.stream === stream)) { + workerFailure = new Error('Hosted Babysitter worker received an unexpected dispatch.'); + peer.close(workerFailure); + return; + } + const attempt = completeHostedDispatch(peer, event, runtime, input, hosted, dataDir); + work = attempt; + void attempt.then(completion => { + completedWork = completion; + if (completion.failed) { + capabilityFailed = true; + capabilityFailure = completion.failure; + } + }, error => { + workerFailure ??= error; + }).finally(() => { + if (work === attempt) work = undefined; + }); + }; + + peer.on('step.dispatch', dispatch); + peer.on('error', error => { workerFailure ??= error; }); + if (lifecycle.onJournalEntry !== undefined) client.on('entry', lifecycle.onJournalEntry); + try { + await peer.connect(); + await peer.hello('flows-hosted-babysitter'); + const pins = { workspace: [], streams: [{ stream, read_offset: 0 }] }; + await peer.workerAttach(`hosted-babysitter-${randomUUID()}`, ['agent'], pins, 1, [stream]); + const admissionKey = `hosted-babysitter:${admissionDigest}`; + const started = await client.runStart(spec, undefined, admissionKey, lifecycle.onJournalEntry !== undefined); + lifecycle.onRunStarted?.({ runId: started.run_id, flow: path }); + // An idempotent concurrent start can observe the original admission while + // its leased attempt is still running. Resume is live-lease-aware and + // converts that snapshot into the same parked/terminal shape handled by + // every other CLI run before classification. + const outcome = started.status === 'running' + ? await client.runResume(started.run_id, true) + : started; + let execution = await classifyOutcome( + client, + 'run', + outcome, + base, + socketPath, + { ...lifecycle, dataDir }, + ); + // The generic classifier may observe a transient parked/protocol shape + // after the isolated child is killed even though this dedicated worker is + // still holding the authoritative, uncancellable capability call. Never + // close the peer or return that intermediate report while its dispatch is + // live. The worker completion is the journal boundary for this one-step + // hosted run; classify its resulting kernel outcome instead. + const completion = work === undefined ? completedWork : await work; + if (completion !== undefined) { + execution = await classifyOutcome( + client, + 'run', + completion.outcome, + base, + socketPath, + { ...lifecycle, dataDir }, + ); + } + if (execution.exitCode !== 1 || !capabilityFailed) return execution; + return { + ...execution, + report: { + ...execution.report, + diagnostics: execution.report.diagnostics.map(diagnostic => diagnostic.kind === 'step_failed' + ? { ...diagnostic, message: errorMessage(capabilityFailure) } + : diagnostic), + }, + }; + } catch (error) { + return protocolFailure('run', base, socketPath, workerFailure ?? error); + } finally { + if (lifecycle.onJournalEntry !== undefined) client.off('entry', lifecycle.onJournalEntry); + peer.off('step.dispatch', dispatch); + peer.close(); + client.close(); + } +} + +type HostedCompletion = + | { readonly outcome: RunOutcome; readonly failed: false } + | { readonly outcome: RunOutcome; readonly failed: true; readonly failure: unknown }; + +async function completeHostedDispatch( + peer: JournalClient, + dispatch: StepDispatchEvent, + runtime: Awaited>, + input: unknown, + hosted: HostedSoftwareGardenRunOptions, + dataDir: string, +): Promise { + let result: unknown; + let effectRecorded = false; + let failed = false; + let failure: unknown; + try { + result = await withWorkerLease(peer, dispatch, async signal => { + return await runHostedCapabilityExtension({ + installation: runtime.installation, + base: runtime.base, + dispatch: hosted.dispatch, + input, + ...(hosted.timeoutMs === undefined ? {} : { timeoutMs: hosted.timeoutMs }), + babysitterTurn: { + queue: async (request, authority) => { + const file = hostedReceiptPath(dataDir, dispatch); + let receipt: unknown; + const performed = await peer.performEffect({ + runId: dispatch.run_id, + stepId: dispatch.step_id, + attempt: dispatch.attempt, + idempotencyKey: dispatch.idempotency_key, + surfacePath: SURFACE_PATH, + revisionBefore: 'pending', + revisionAfter: hosted.dispatch.deliveryId, + }, async () => { + // performEffect reaches this callback only after effect.record + // succeeded, so a provider rejection still has a journal fact + // the failing step completion must carry. + effectRecorded = true; + signal.throwIfAborted(); + // A crash after the durable receipt write but before + // effect.confirm leaves this election reclaimable. Consume that + // receipt on the reclaimed attempt instead of calling the + // external provider a second time. + try { + receipt = await readHelperReceipt(file); + } catch (error) { + if ((error as NodeJS.ErrnoException).code !== 'ENOENT') throw error; + receipt = await hosted.babysitterTurn.queue(request, authority); + await atomicJson(file, receipt); + } + signal.throwIfAborted(); + }); + // Also required after a confirmed election followed by a crash + // before step.complete. An elected recovery callback may already + // have populated receipt above, so this is deliberately based on + // the value rather than only on performEffect's return flag. + if (!performed) effectRecorded = true; + if (receipt === undefined) receipt = await readHelperReceipt(file); + return receipt; + }, + }, + }); + }); + } catch (error) { + failed = true; + failure = error; + } + + const output = failed + ? { type: 'hosted-flow-extension', diagnostic: errorMessage(failure) } + : { type: 'hosted-flow-extension', result }; + const outcome = await peer.stepComplete( + dispatch.run_id, + dispatch.step_id, + dispatch.attempt, + dispatch.idempotency_key, + failed ? 'worker_error' : 'success', + { + output, + started_pins: dispatch.pins, + end_pins: dispatch.pins, + ...(failed ? { trajectory_tail: output } : {}), + effects: effectRecorded + ? [{ surface_path: SURFACE_PATH, idempotency_key: dispatch.idempotency_key }] + : [], + }, + ); + return failed ? { outcome, failed: true, failure } : { outcome, failed: false }; +} + +function hostedReceiptPath(dataDir: string, dispatch: StepDispatchEvent): string { + const name = createHash('sha256') + .update(`${dispatch.run_id}:${dispatch.step_id}:${dispatch.idempotency_key}`) + .digest('hex'); + return join(dataDir, 'hosted-extension-receipts', `${name}.json`); +} + +function pluginRefusal(base: RunReport, error: PluginError): RunExecution { + return { + exitCode: 2, + report: { + ...base, + diagnostics: [...base.diagnostics, { + severity: 'refusal', + kind: error.code, + message: error.message, + }], + }, + }; +} + +function hostedFailure( + base: RunReport, + error: unknown, + runId?: string, + socketPath?: string, + completedSteps?: number, +): RunExecution { + return { + exitCode: 1, + report: { + ...base, + ...(runId === undefined ? {} : { runId }), + ...(socketPath === undefined ? {} : { socketPath }), + status: 'failed', + completionReason: 'step_failed', + ...(completedSteps === undefined ? {} : { completedSteps }), + diagnostics: [...base.diagnostics, { + severity: 'failure', + kind: 'step_failed', + message: errorMessage(error), + }], + }, + }; +} + +function errorMessage(error: unknown): string { + return error instanceof Error ? error.message : 'Hosted Babysitter capability failed.'; +} diff --git a/packages/sdk/src/hosted-base-snapshot.ts b/packages/sdk/src/hosted-base-snapshot.ts index 284d7f4a2..747f4adde 100644 --- a/packages/sdk/src/hosted-base-snapshot.ts +++ b/packages/sdk/src/hosted-base-snapshot.ts @@ -27,7 +27,10 @@ import { } from './hosted-promise-safety.js'; import { PluginError } from './plugin-manifest.js'; -const EXCLUDED_DIRECTORIES = new Set(['.flows', '.git', 'node_modules']); +// Runtime journals are not authored base source. Keeping the standard data +// directory in the generation would make daemon startup invalidate the exact +// generation it was started to execute. +const EXCLUDED_DIRECTORIES = new Set(['.flows', '.git', '.relayflowd', 'node_modules']); const CHMOD = chmod; const MKDIR = mkdir; const MKDTEMP = mkdtemp; diff --git a/packages/sdk/src/hosted-extension-protocol.ts b/packages/sdk/src/hosted-extension-protocol.ts index a468587d7..17ac1ccd6 100644 --- a/packages/sdk/src/hosted-extension-protocol.ts +++ b/packages/sdk/src/hosted-extension-protocol.ts @@ -141,22 +141,41 @@ export async function exchangeHostedExtension( finish(capabilityError ?? error); }; const timeout = SET_TIMEOUT(() => { - CHILD_PROCESS_KILL(child, 'SIGKILL'); - finish(new PluginError( + const error = new PluginError( 'plugin_unsupported', capabilityState === 'pending' ? 'Hosted capability outcome is in doubt after the sandbox timeout.' : 'Hosted extension sandbox timed out.', - )); + ); + // The host capability is not cancellable at this boundary. Reporting a + // terminal result while it is still running would let the provider write + // land after the journal had already completed the step. Kill the tenant + // sandbox immediately, but defer the terminal protocol result until the + // in-flight capability has settled and its effect can be journaled. + if (capabilityState === 'pending') { + deferredProtocolError ??= error; + CHILD_PROCESS_KILL(child, 'SIGKILL'); + return; + } + CHILD_PROCESS_KILL(child, 'SIGKILL'); + finish(error); }, timeoutMs); TIMER_UNREF(timeout); EVENT_ON(stdin, 'error', () => refuse('Hosted extension capability channel closed.')); EVENT_ON(protocol, 'error', () => refuse('Hosted extension protocol channel failed.')); EVENT_ON(protocol, 'data', chunk => { + // Streams may still deliver bytes buffered before finish killed the + // sandbox. A terminal protocol exchange must never start a capability + // from one of those late frames. + if (settled) return; buffer += STRING(chunk); if (BUFFER_BYTE_LENGTH(buffer) > MAX_FRAME_BYTES) return refuse('Hosted extension protocol exceeded its size limit.'); for (;;) { + // A frame earlier in this same chunk may have called finish. Stop at + // the terminal frame boundary before parsing or invoking anything + // else already buffered behind it. + if (settled) return; const end = STRING_INDEX_OF(buffer, '\n'); if (end < 0) break; const line = STRING_SLICE(buffer, 0, end); @@ -211,7 +230,7 @@ export async function exchangeHostedExtension( || message.completionReason !== 'success' || message.capabilityCalls !== 1) { return refuse('Hosted extension reported a completion without exactly one capability call.'); } - finish(undefined, frozenHostedPromiseValue({ completionReason: 'success', capabilityCalls: 1 })); + return finish(undefined, frozenHostedPromiseValue({ completionReason: 'success', capabilityCalls: 1 })); } else if (message.type === 'error') { if (!hasExactKeys(message, ['type', 'message']) || typeof message.message !== 'string') { return refuse('Hosted extension emitted a malformed error frame.'); @@ -224,7 +243,7 @@ export async function exchangeHostedExtension( deferredProtocolError ??= error; return; } - finish(capabilityError ?? error); + return finish(capabilityError ?? error); } else return refuse('Hosted extension emitted an unknown protocol message.'); } }); diff --git a/packages/sdk/tests/hosted-base-snapshot.test.ts b/packages/sdk/tests/hosted-base-snapshot.test.ts index 9bc975849..1ea16b372 100644 --- a/packages/sdk/tests/hosted-base-snapshot.test.ts +++ b/packages/sdk/tests/hosted-base-snapshot.test.ts @@ -192,6 +192,22 @@ describe('hosted base private snapshot', () => { } }); + it('excludes the standard relayflowd data directory from the admitted generation', async () => { + const { project, flowPath } = fixture(); + const data = join(project, '.relayflowd'); + mkdirSync(join(data, 'runs'), { recursive: true }); + writeFileSync(join(data, 'relayflowd.sqlite3'), 'before'); + const snapshot = await createHostedBaseSnapshot(flowPath); + try { + expect(() => readFileSync(join(snapshot.snapshotRoot, '.relayflowd/relayflowd.sqlite3'))).toThrow(); + const before = await hostedBaseSourceDigest(snapshot.liveSources); + writeFileSync(join(data, 'runs/new.sqlite3'), 'after'); + expect(await hostedBaseSourceDigest(snapshot.liveSources)).toBe(before); + } finally { + await removeHostedBaseSnapshot(snapshot); + } + }); + it('refuses an oversized source file before buffering its contents', async () => { const { project, flowPath } = fixture(); const oversized = join(project, 'oversized.bin'); diff --git a/packages/sdk/tests/hosted-extension-protocol.test.ts b/packages/sdk/tests/hosted-extension-protocol.test.ts index d5a86b8d0..818a05f98 100644 --- a/packages/sdk/tests/hosted-extension-protocol.test.ts +++ b/packages/sdk/tests/hosted-extension-protocol.test.ts @@ -360,6 +360,99 @@ describe('hosted extension hostile protocol', () => { expect(calls).toBe(1); }, 15_000); + it('waits for a pending adapter to settle after the sandbox timeout', async () => { + const protocol = new PassThrough(); + const stdin = new PassThrough(); + const stderr = new PassThrough(); + const child = Object.assign(new EventEmitter(), { + exitCode: null, + signalCode: null, + kill: () => true, + }) as unknown as ChildProcess; + let settle!: (value: unknown) => void; + const adapter = new Promise(resolve => { settle = resolve; }); + let markInvoked!: () => void; + const invoked = new Promise(resolve => { markInvoked = resolve; }); + let completed = false; + const run = exchangeHostedExtension( + child, + protocol, + stdin, + stderr, + 10, + { type: 'run' }, + async () => { markInvoked(); return await adapter; }, + ).finally(() => { completed = true; }); + protocol.write(`${JSON.stringify(capabilityFrame())}\n`); + await waitForInvocation(invoked); + await new Promise(resolve => setTimeout(resolve, 30)); + expect(completed).toBe(false); + settle({ receiptId: 'settled-after-timeout', status: 'queued' }); + await expect(run).rejects.toMatchObject({ + code: 'plugin_unsupported', + message: expect.stringContaining('outcome is in doubt'), + }); + }, 15_000); + + it('does not invoke a capability frame delivered after timeout settlement', async () => { + const protocol = new PassThrough(); + const stdin = new PassThrough(); + const stderr = new PassThrough(); + const child = Object.assign(new EventEmitter(), { + exitCode: null, + signalCode: null, + kill: () => true, + }) as unknown as ChildProcess; + let calls = 0; + const run = exchangeHostedExtension( + child, + protocol, + stdin, + stderr, + 10, + { type: 'run' }, + async () => { calls += 1; return { receiptId: 'late', status: 'queued' }; }, + ); + const observed = run.then(() => undefined, error => error as Error); + await new Promise(resolve => setTimeout(resolve, 30)); + expect(await observed).toMatchObject({ + code: 'plugin_unsupported', + message: expect.stringContaining('sandbox timed out'), + }); + protocol.write(`${JSON.stringify(capabilityFrame())}\n`); + await new Promise(resolve => setTimeout(resolve, 0)); + expect(calls).toBe(0); + }); + + it('does not invoke a capability buffered after a terminal frame in the same chunk', async () => { + const protocol = new PassThrough(); + const stdin = new PassThrough(); + const stderr = new PassThrough(); + const child = Object.assign(new EventEmitter(), { + exitCode: null, + signalCode: null, + kill: () => true, + }) as unknown as ChildProcess; + let calls = 0; + const run = exchangeHostedExtension( + child, + protocol, + stdin, + stderr, + 10_000, + { type: 'run' }, + async () => { calls += 1; return { receiptId: 'late', status: 'queued' }; }, + ); + const observed = run.then(() => undefined, error => error as Error); + protocol.write(`${JSON.stringify({ type: 'error', message: 'terminal' })}\n${JSON.stringify(capabilityFrame())}\n`); + expect(await observed).toMatchObject({ + code: 'plugin_unsupported', + message: expect.stringContaining('Hosted extension failed: terminal'), + }); + await new Promise(resolve => setTimeout(resolve, 0)); + expect(calls).toBe(0); + }); + it('settles successful completion with captured intrinsics', async () => { const protocol = new PassThrough(); const stdin = new PassThrough(); diff --git a/packages/sdk/tests/software-garden-babysitter-canonical-run.test.ts b/packages/sdk/tests/software-garden-babysitter-canonical-run.test.ts new file mode 100644 index 000000000..1196ab975 --- /dev/null +++ b/packages/sdk/tests/software-garden-babysitter-canonical-run.test.ts @@ -0,0 +1,486 @@ +import { + appendFileSync, + existsSync, + mkdtempSync, + readFileSync, + readdirSync, + rmSync, + writeFileSync, +} from 'node:fs'; +import { tmpdir } from 'node:os'; +import { join, resolve } from 'node:path'; +import { setTimeout as delay } from 'node:timers/promises'; +import { afterAll, describe, expect, it, vi } from 'vitest'; +import { runCli, type RunCliOptions } from '../src/cli.js'; +import { addExtensionPlugin } from '../src/cli/add-extension.js'; +import { hostedExtensionDispatchFromVerifiedDelivery } from '../src/flow-extension-loader.js'; +import { JournalClient } from '../src/journal-client.js'; +import { entriesFromDirectory, fakeGithub } from './fake-github.js'; + +const NATIVE_SHA = '8b33ebab8347514f80d9da5a81206a087f641714'; +const REF = `github:AgentWorkforce/flows@${NATIVE_SHA}#extensions/babysitter`; +const DIGEST = 'bdf2187b9a242667d34bbc63e7a744753e146dc8cd6f4047047f2aed28f406ee'; +const MANIFEST_SHA256 = '5631a06bbdc8186f4ee0ff955610ead24d001c5197b59fb1fe81fe422c44f226'; +const entries = entriesFromDirectory(resolve('../..', 'extensions/babysitter'), 'extensions/babysitter'); +const versions = { sdk: '2.0.33', surface: '2.0.33' }; +const roots: string[] = []; + +afterAll(() => roots.splice(0).forEach(root => rmSync(root, { recursive: true, force: true }))); + +function github() { + return fakeGithub({ 'AgentWorkforce/flows': { refs: {}, commits: { [NATIVE_SHA]: { entries } } } }); +} + +async function project(install = true) { + const root = mkdtempSync(join(tmpdir(), 'software-garden-canonical-run-')); + roots.push(root); + const flowPath = join(root, 'software-factory.flow.ts'); + writeFileSync(flowPath, readFileSync(resolve('../..', 'examples/software-factory/software-factory.flow.ts'))); + writeFileSync(join(root, 'flows.json'), JSON.stringify({ cli: 'codex', executors: ['github'] })); + if (install) { + const io = { stdout: () => {}, stderr: (message: string) => { throw new Error(message); } }; + expect(await addExtensionPlugin(REF, io, { + cwd: root, + fetch: github().fetch, + now: () => new Date('2026-09-28T00:00:00Z'), + versions, + })).toBe(0); + } else { + writeFileSync(join(root, 'flows.lock.json'), JSON.stringify({ version: 2, plugins: [] })); + } + return { root, flowPath }; +} + +function input(eventType = 'pull_request.labeled', deliveryId = 'delivery-canonical') { + return { + event: { provider: 'github', eventType, deliveryId }, + pullRequest: { + host: 'github', + owner: 'AgentWorkforce', + repo: 'flows', + number: 584, + headSha: 'a'.repeat(40), + }, + }; +} + +function hosted( + eventType: string, + deliveryId: string, + queue: NonNullable['babysitterTurn']['queue'], +): NonNullable { + return { + dispatch: hostedExtensionDispatchFromVerifiedDelivery({ + provider: 'github', + eventType, + deliveryId, + }), + babysitterTurn: { queue }, + }; +} + +async function run( + flowPath: string, + descriptor: unknown, + authority: NonNullable, +) { + const stdout: string[] = []; + const stderr: string[] = []; + const exitCode = await runCli([ + 'run', + flowPath, + '--input', + JSON.stringify(descriptor), + '--data-dir', + join(flowPath, '..', '.relayflowd'), + '--json', + '--no-observer-link', + ], { + stdout: line => stdout.push(line), + stderr: line => stderr.push(line), + }, { hostedSoftwareGardenBabysitter: authority }); + return { + exitCode, + report: JSON.parse(stdout.at(-1) ?? '{}') as Record, + stderr, + }; +} + +describe('canonical run dispatches the installed Software Garden Babysitter', () => { + it('refuses hosted authority on any command surface other than an authored run', async () => { + const stdout: string[] = []; + const exitCode = await runCli(['run', 'software-factory.flow.yaml', '--json'], { + stdout: line => stdout.push(line), + stderr: () => {}, + }, { + hostedSoftwareGardenBabysitter: hosted( + 'pull_request.labeled', + 'delivery-canonical', + async () => ({ receiptId: 'never', status: 'queued' }), + ), + }); + expect(exitCode).toBe(2); + expect(JSON.parse(stdout.at(-1) ?? '{}')).toMatchObject({ + ok: false, + path: 'software-factory.flow.yaml', + diagnostics: [expect.objectContaining({ + severity: 'refusal', + kind: 'invalid_invocation', + })], + }); + }); + + it('records one deterministic, exact reviewed plugin inventory member', async () => { + const installed = await project(); + const lock = JSON.parse(readFileSync(join(installed.root, 'flows.lock.json'), 'utf8')) as { + plugins: unknown[]; + }; + expect(lock.plugins).toEqual([{ + name: 'babysitter', + version: '0.2.0', + kind: 'flow-extension', + source: { + host: 'github', + owner: 'AgentWorkforce', + repo: 'flows', + sha: NATIVE_SHA, + path: 'extensions/babysitter', + }, + digest: DIGEST, + manifestSha256: MANIFEST_SHA256, + resolvedAt: '2026-09-28T00:00:00.000Z', + order: 1, + }]); + }); + + it.each(['--local-agent', '--agent-capacity', '--allow-human-influenced'])( + 'refuses the unsupported %s flag instead of silently ignoring it', async (flag) => { + const installed = await project(); + const stdout: string[] = []; + const args = [ + 'run', installed.flowPath, '--input', JSON.stringify(input()), flag, + ...(flag === '--agent-capacity' ? ['2'] : []), '--json', '--no-observer-link', + ]; + const exitCode = await runCli(args, { + stdout: line => stdout.push(line), + stderr: () => {}, + }, { + hostedSoftwareGardenBabysitter: hosted( + 'pull_request.labeled', + 'delivery-canonical', + async () => ({ receiptId: 'never', status: 'queued' }), + ), + }); + expect(exitCode).toBe(2); + expect(JSON.parse(stdout.at(-1) ?? '{}')).toMatchObject({ + ok: false, + diagnostics: [expect.objectContaining({ + severity: 'refusal', + kind: 'invalid_invocation', + message: expect.stringContaining(flag), + })], + }); + }, + ); + + it.runIf(process.platform === 'linux')( + 'refuses a canonical hosted run with no installed Babysitter before the capability is called', async () => { + const empty = await project(false); + let calls = 0; + const result = await run(empty.flowPath, input(), hosted( + 'pull_request.labeled', + 'delivery-canonical', + async () => { calls += 1; return { receiptId: 'never', status: 'queued' }; }, + )); + expect(result).toMatchObject({ + exitCode: 2, + report: { ok: false, command: 'run', path: empty.flowPath }, + }); + expect(result.report.diagnostics).toEqual([ + expect.objectContaining({ severity: 'refusal', kind: 'plugin_event_unroutable' }), + ]); + expect(calls).toBe(0); + }, + ); + + it.runIf(process.platform === 'linux')( + 'refuses installed-store drift before the canonical run calls the capability', async () => { + const changed = await project(); + appendFileSync(join(changed.root, '.flows/plugins', `babysitter@sha256:${DIGEST}`, 'turn.ts'), '\n// drift\n'); + let calls = 0; + const result = await run(changed.flowPath, input(), hosted( + 'pull_request.labeled', + 'delivery-canonical', + async () => { calls += 1; return { receiptId: 'never', status: 'queued' }; }, + )); + expect(result).toMatchObject({ exitCode: 2, report: { ok: false } }); + expect(result.report.diagnostics).toEqual([ + expect.objectContaining({ severity: 'refusal', kind: 'plugin_source_drift' }), + ]); + expect(calls).toBe(0); + }, + ); + + it.runIf(process.platform === 'linux')( + 'refuses an unmatched verified route before the canonical run calls the capability', async () => { + const installed = await project(); + let calls = 0; + const result = await run(installed.flowPath, input('push'), hosted( + 'push', + 'delivery-canonical', + async () => { calls += 1; return { receiptId: 'never', status: 'queued' }; }, + )); + expect(result).toMatchObject({ exitCode: 2, report: { ok: false } }); + expect(result.report.diagnostics).toEqual([ + expect.objectContaining({ severity: 'refusal', kind: 'plugin_event_unroutable' }), + ]); + expect(calls).toBe(0); + }, + ); + + it.skipIf(process.platform !== 'linux' || !existsSync('/usr/bin/bwrap'))( + 'runs the exact matched installed handler through the sandbox from the normal run command', async () => { + const installed = await project(); + const deliveryId = 'delivery-canonical-e2e'; + let calls = 0; + const authority = hosted('pull_request.labeled', deliveryId, async (request, received) => { + calls += 1; + expect(received.dispatch).toBe(authority.dispatch); + expect(received.extension).toEqual({ + name: 'babysitter', + version: '0.2.0', + ref: REF, + digest: DIGEST, + }); + expect(request).toEqual({ delivery: { + deliveryId, + provider: 'github', + eventType: 'pull_request.labeled', + pullRequest: { owner: 'AgentWorkforce', repository: 'flows', number: 584 }, + } }); + return { receiptId: `bst_${'1'.repeat(64)}`, status: 'queued' }; + }); + + const first = await run(installed.flowPath, input('pull_request.labeled', deliveryId), authority); + expect(first).toMatchObject({ + exitCode: 0, + report: { + ok: true, + command: 'run', + path: installed.flowPath, + runId: expect.any(String), + socketPath: expect.any(String), + status: 'completed', + completionReason: 'success', + completedSteps: 1, + diagnostics: [], + }, + }); + const retried = await run(installed.flowPath, input('pull_request.labeled', deliveryId), authority); + expect(retried).toMatchObject({ + exitCode: 0, + report: { + ok: true, + runId: first.report.runId, + status: 'completed', + completionReason: 'success', + completedSteps: 1, + }, + }); + expect(calls).toBe(1); + }, + ); + + it.skipIf(process.platform !== 'linux' || !existsSync('/usr/bin/bwrap'))( + 'does not complete the journal while a timed-out capability is still pending', async () => { + const installed = await project(); + const deliveryId = 'delivery-timeout-pending'; + let calls = 0; + let release!: () => void; + let capabilityStarted!: () => void; + const released = new Promise(resolveReleased => { release = resolveReleased; }); + const started = new Promise(resolveStarted => { capabilityStarted = resolveStarted; }); + const authority = { + ...hosted('pull_request.labeled', deliveryId, async () => { + calls += 1; + capabilityStarted(); + await released; + return { receiptId: `bst_${'2'.repeat(64)}`, status: 'queued' }; + }), + // Sandbox startup is part of this deadline. Leave enough time for the + // capability to begin, then deliberately hold it beyond the deadline. + timeoutMs: 3_000, + }; + + let settled = false; + const pending = run(installed.flowPath, input('pull_request.labeled', deliveryId), authority); + void pending.then(() => { settled = true; }, () => { settled = true; }); + await started; + await delay(3_050); + expect(settled).toBe(false); + expect(calls).toBe(1); + + release(); + const result = await pending; + expect(result).toMatchObject({ + exitCode: 1, + report: { + ok: false, + status: 'failed', + completionReason: 'step_failed', + diagnostics: [expect.objectContaining({ + severity: 'failure', + kind: 'step_failed', + message: expect.stringContaining('outcome is in doubt'), + })], + }, + }); + const receipts = join(installed.root, '.relayflowd', 'hosted-extension-receipts'); + const afterCompletion = readdirSync(receipts).map(file => readFileSync(join(receipts, file), 'utf8')); + await delay(50); + expect(readdirSync(receipts).map(file => readFileSync(join(receipts, file), 'utf8'))).toEqual(afterCompletion); + expect(calls).toBe(1); + }, + 10_000, + ); + + it.skipIf(process.platform !== 'linux' || !existsSync('/usr/bin/bwrap'))( + 'reuses the durable receipt when an elected effect is replayed before confirmation', async () => { + const installed = await project(); + const deliveryId = 'delivery-reclaimed-receipt'; + let calls = 0; + const replay = vi.spyOn(JournalClient.prototype, 'performEffect').mockImplementationOnce( + async function (effect, perform) { + const { deduped } = await this.effectRecord( + effect.runId, + effect.stepId, + effect.attempt, + effect.idempotencyKey, + effect.surfacePath, + effect.revisionBefore, + effect.revisionAfter, + ); + expect(deduped).toBe(false); + await perform(); + // Inject the crash boundary: the provider receipt is durable, but + // effect.confirm has not happened and a reclaimed election reruns + // the callback. Recovery must consume the receipt, not write twice. + await perform(); + await this.effectConfirm( + effect.runId, + effect.stepId, + effect.attempt, + effect.idempotencyKey, + effect.surfacePath, + ); + return true; + }, + ); + try { + const result = await run(installed.flowPath, input('pull_request.labeled', deliveryId), hosted( + 'pull_request.labeled', + deliveryId, + async () => { + calls += 1; + return { receiptId: `bst_${'3'.repeat(64)}`, status: 'queued' }; + }, + )); + expect(result).toMatchObject({ + exitCode: 0, + report: { + ok: true, + status: 'completed', + completionReason: 'success', + completedSteps: 1, + }, + }); + expect(calls).toBe(1); + } finally { + replay.mockRestore(); + } + }, + ); + + it.skipIf(process.platform !== 'linux' || !existsSync('/usr/bin/bwrap'))( + 'reports capability denial once as a terminal failure without fallback', async () => { + const installed = await project(); + const refusal = new Error('live babysit label is absent'); + let calls = 0; + const result = await run(installed.flowPath, input(), hosted( + 'pull_request.labeled', + 'delivery-canonical', + async () => { calls += 1; throw refusal; }, + )); + expect(result).toMatchObject({ + exitCode: 1, + report: { + ok: false, + runId: expect.any(String), + socketPath: expect.any(String), + status: 'failed', + completionReason: 'step_failed', + diagnostics: [expect.objectContaining({ + severity: 'failure', + kind: 'step_failed', + message: refusal.message, + })], + }, + }); + expect(calls).toBe(1); + }, + ); + + it.skipIf(process.platform !== 'linux' || !existsSync('/usr/bin/bwrap'))( + 'does not misclassify a post-capability receipt error as a clean refusal', async () => { + const installed = await project(); + let calls = 0; + const result = await run(installed.flowPath, input(), hosted( + 'pull_request.labeled', + 'delivery-canonical', + async () => { calls += 1; return { receiptId: '', status: 'queued' }; }, + )); + expect(result).toMatchObject({ + exitCode: 1, + report: { + ok: false, + runId: expect.any(String), + socketPath: expect.any(String), + status: 'failed', + completionReason: 'step_failed', + diagnostics: [expect.objectContaining({ + severity: 'failure', + kind: 'step_failed', + })], + }, + }); + expect(calls).toBe(1); + }, + ); + + it.skipIf(process.platform !== 'linux' || !existsSync('/usr/bin/bwrap'))( + 'fails closed when the capability rejects without an Error value', async () => { + const installed = await project(); + let calls = 0; + const result = await run(installed.flowPath, input(), hosted( + 'pull_request.labeled', + 'delivery-canonical', + async () => { calls += 1; return await Promise.reject(undefined); }, + )); + expect(result).toMatchObject({ + exitCode: 1, + report: { + ok: false, + runId: expect.any(String), + status: 'failed', + completionReason: 'step_failed', + diagnostics: [expect.objectContaining({ + severity: 'failure', + kind: 'step_failed', + message: expect.stringContaining('non-error value'), + })], + }, + }); + expect(calls).toBe(1); + }, + ); +}); diff --git a/packages/sdk/type-tests/hosted-software-garden-run.ts b/packages/sdk/type-tests/hosted-software-garden-run.ts new file mode 100644 index 000000000..917b96905 --- /dev/null +++ b/packages/sdk/type-tests/hosted-software-garden-run.ts @@ -0,0 +1,11 @@ +import type { RunCliOptions } from '../src/cli.js'; + +type Hosted = NonNullable; + +// Production callers can supply only the verified authority/capability pair +// and a timeout. Sandbox executable paths remain internal test seams. +const noBubblewrapOverride: 'bubblewrapPath' extends keyof Hosted ? true : false = false; +const noNodeOverride: 'nodePath' extends keyof Hosted ? true : false = false; +const noPrlimitOverride: 'prlimitPath' extends keyof Hosted ? true : false = false; + +void [noBubblewrapOverride, noNodeOverride, noPrlimitOverride];