From 562a478899ee11a443ee897e6e387f4efd09e584 Mon Sep 17 00:00:00 2001 From: Relayflow Lead Date: Thu, 17 Sep 2026 15:10:01 -0700 Subject: [PATCH 1/4] feat(sdk,surface): journaled agent artifacts, artifact_exists gate, predicate gates MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Three pieces that together make "the agent must have written this file" a journal-honest, replayable check for authored flows. Artifacts move into the worker. #434 measured them in the authored step, after the journal entry was written, so a gate could not read them and a resume could not reproduce them. `runAgentCli` (the process that spawns the CLI in its cwd) now snapshots that directory before the spawn and content- diffs it after, and `WorkerCliResult.artifacts` rides in the step's journaled `output` on the CliResult path. `AgentResult.artifacts` is read from that entry — no second scan, no guess about which host wrote what. Relay transport journals none (the agent ran elsewhere); an agent whose final message is a JSON object owns its output shape, which is left untouched, so such a step journals no artifacts either (documented; gate those on a deterministic check). `agent-artifacts.ts` and its unit tests are kept as-is. `artifact_exists` named gate: `{ type: 'artifact_exists', path }` end to end — surface `NamedGate`, SDK `NamedDataGate`, validation (`gate_path_invalid`: relative POSIX path, no empty/./.. segments, no NUL), lowering to a deterministic `.gate` step whose command judges only `input.output.artifacts` from FLOWS_INPUT, `unknown_gate_kind` message, preflight refusal coverage, regenerated JSON schema. Predicate gates: `.gate(fn, because?)` was refused as `unsupported_gate` although SURFACE.md §6 specified running it as runtime control flow. The executor now runs the closure once, after the step, on the journaled value, and journals the verdict as a lowered `.gate` deterministic step (`{"gate":"predicate","step","verdict","because"}`; exit 0/1), so resume and replay read the recorded verdict and never re-run author code. A false or throwing predicate fails the run as `gate_failed` naming the step and the reason; `flows run`/`resume` report it as a run failure, not a protocol one. One `.gate()` per step. The function is never serialized; `flows check` prints no gate line for it (runtime-only, by construction — documented). Tests: worker artifacts through a mocked spawn; gate validation, lowering (command exercised against FLOWS_INPUT), `flows check` inspection; #434's loopback integration adapted to the journaled contract; and a live test that invokes the built CLI with Cloud's exact argv (`run --json --data-dir … --local-agent x.flow.ts --input …`) against a real daemon and a wrapper CLI that writes `review/*.md`, asserting `output.artifacts` in the journal, both gate forms passing, and both negative cases failing the run with the gate named. Examples: pr-review-pipeline uses `artifact_exists` per lens and a predicate for consensus; examples/README no longer calls it BLOCKED on the budget header (accepted since #306). Co-Authored-By: Claude Opus 5 (1M context) --- docs/SURFACE.md | 33 ++- examples/README.md | 14 +- .../pr-review-pipeline.flow.ts | 27 ++- packages/schema/flows.schema.json | 31 ++- packages/sdk/src/authored-flow-error.ts | 1 + packages/sdk/src/authored-flow-executor.ts | 50 ++++- packages/sdk/src/authored-flow-operation.ts | 35 +++- packages/sdk/src/authored-worker-step.ts | 24 +-- packages/sdk/src/cli/direct-run.ts | 2 +- packages/sdk/src/cli/run.ts | 8 +- packages/sdk/src/failure-kinds.ts | 2 + packages/sdk/src/named-gate-lowering.ts | 9 + packages/sdk/src/named-gates.ts | 12 +- packages/sdk/src/spec.ts | 13 +- packages/sdk/src/validate.ts | 2 +- packages/sdk/src/worker-cli.ts | 28 ++- packages/sdk/src/worker.ts | 9 + .../sdk/tests/agent-artifacts-live.test.ts | 193 ++++++++++++++++++ packages/sdk/tests/artifact-gates.test.ts | 130 ++++++++++++ .../tests/authored-agent-artifacts.test.ts | 43 ++-- packages/sdk/tests/authored-flow.test.ts | 16 +- packages/sdk/tests/preflight.test.ts | 2 + packages/surface/src/context.ts | 7 + packages/surface/src/step.ts | 28 ++- 24 files changed, 628 insertions(+), 91 deletions(-) create mode 100644 packages/sdk/tests/agent-artifacts-live.test.ts create mode 100644 packages/sdk/tests/artifact-gates.test.ts diff --git a/docs/SURFACE.md b/docs/SURFACE.md index d3b115c33..e86c8c71a 100644 --- a/docs/SURFACE.md +++ b/docs/SURFACE.md @@ -845,15 +845,34 @@ author predicate. The v1 `verification:` shape remains supported and compiles to the same kernel fields; no kernel verb or verification field is added by this decision. -TypeScript may additionally accept a callback such as +TypeScript additionally accepts a callback such as `.gate(value => value.length < 200, "keep the summary short")`. That callback is author code: `flows check` cannot prove it, YAML cannot serialize it, and -the journal cannot replay the closure. A TypeScript runtime must execute it as -runtime control flow and journal the resulting step outcome before dependents -continue. It must never stringify the function into a spec or silently label -it preflightable. Authors who need portable, inspectable gates use a named data -check; plugins may contribute named checks only by compiling them to existing -kernel primitives. +the journal cannot replay the closure. The authored runtime executes it as +runtime control flow — once, in the authoring process, on the value read back +from the step's `step.completed` — and journals the verdict as a lowered +`.gate` deterministic step: a passing predicate journals +`{"gate":"predicate","step":"","verdict":"pass","because":…}` as that +step's stdout with exit 0; a failing one (or one that throws) journals +`"verdict":"fail"` on stderr with exit 1, and the run fails as `gate_failed` +naming the step and the author's reason. Dependents therefore wait on a +journaled fact, and resume/replay read that fact rather than re-running the +closure. The function is never stringified into a spec, and `flows check` +prints no gate line for it — a predicate is runtime-only and unprovable +before execution, by construction. A step takes one `.gate()`. Authors who +need portable, inspectable gates use a named data check; plugins may +contribute named checks only by compiling them to existing kernel primitives. + +`artifact_exists` is the named gate for "the agent wrote this file": +`.gate({ type: 'artifact_exists', path: 'review/security.md' })`. The worker +that spawned the agent CLI snapshots the agent's working directory before the +run and content-diffs it after, and journals the changed paths as +`output.artifacts` on the agent's `step.completed`; `AgentResult.artifacts` +is read from that journal entry, never from a later look at the disk, and the +gate lowers to a deterministic step that checks the journaled list. An agent +whose final message is a JSON object owns its output shape and journals no +artifacts; gate such a step on a deterministic check instead. The relay +transport journals none, because the agent ran on another host. - Are YAML helper verbs (`slack:`, `mcp:`) core spec vocabulary or compile-time expansion into `run`/effect steps? Leaning: expansion — the kernel spec stays seven words; helpers stay a surface concern. - Helper generation cadence: generated from relayfile adapter manifests at build time vs published per-adapter packages. Leaning: generated, with hand-tuned verb names for the top providers. diff --git a/examples/README.md b/examples/README.md index 606dad24d..f2069fdb2 100644 --- a/examples/README.md +++ b/examples/README.md @@ -9,7 +9,7 @@ For a working local starting point, use the [small agent starter](../README.md) | Example | Status | Observed result | Elapsed | |---|---|---|---:| | [dependency-upgrade-bot](dependency-upgrade-bot/) | **BLOCKED** | SDK refuses unsupported `budget` header before entering the body; exit 2 | [5.138s](../docs/evidence/ws13/review/gallery/gallery-dependency-upgrade-bot.txt) | -| [pr-review-pipeline](pr-review-pipeline/) | **BLOCKED** | SDK refuses unsupported `budget` header before entering the body; exit 2 | [3.539s](../docs/evidence/ws13/review/gallery/gallery-pr-review-pipeline.txt) | +| [pr-review-pipeline](pr-review-pipeline/) | **RUNNABLE** | `budget:` headers have been accepted since #306; agent artifacts are journaled by the worker and both gate forms (`artifact_exists`, predicate) are lowered — see `packages/sdk/tests/agent-artifacts-live.test.ts` for the same shape through the built CLI and a real daemon. Needs a `flows.json` naming an authenticated agent CLI. | — | | [research](research/) | **PASS** | All model probes passed; three lane reports and synthesis produced; exit 0, `completionReason: synthesized` | [690.935s](../docs/evidence/ws13/followup/default-budget/gallery-research.txt) | **Correction:** the previously listed 0.138s and 0.143s captures used a stale @@ -22,12 +22,12 @@ These are individual runs from a separate clone on an authenticated macOS host, against the packed candidate CLI. Research uses its documented source shim. These timings are not clean-machine measurements. -**Dependency-upgrade-bot and pr-review-pipeline need the SDK/kernel capability -owner.** Their authored budgets are currently rejected. Their postfix artifact -gates and workspace permission declarations also require runtime support. -Removing those requirements would weaken what the examples promise; this -branch leaves them intact. The local agent worker handles stream-only steps -and cannot supply workspace isolation. +**Dependency-upgrade-bot still needs the SDK/kernel capability owner** for +its workspace permission declarations: the local agent worker handles +stream-only steps and cannot supply workspace isolation, and a +`"...: readwrite"` annotation is refused because nothing enforces it. +Budget headers and postfix artifact gates are supported; pr-review-pipeline +uses them without workspace scoping. **Research now prints provider preflight activity.** Each CLI/model probe names its timeout on stderr, while stdout remains the final structured result. diff --git a/examples/pr-review-pipeline/pr-review-pipeline.flow.ts b/examples/pr-review-pipeline/pr-review-pipeline.flow.ts index 7b41e7ac5..febd232aa 100644 --- a/examples/pr-review-pipeline/pr-review-pipeline.flow.ts +++ b/examples/pr-review-pipeline/pr-review-pipeline.flow.ts @@ -9,11 +9,16 @@ // looks specifically for conflicting verdicts rather than just concatenating // opinions. // -// STATUS: typechecks against the real `@relayflows/surface` package (see -// ../tsconfig.json / `npm --prefix packages/surface run typecheck:examples`) -// but does not run yet — `f.agent` parks without an attached worker. Unlike -// the other two examples, this one never calls `f.human`, so nothing here -// depends on that verb being wired up. +// STATUS: runs with `flows run pr-review-pipeline.flow.ts --local-agent +// --input '{"diffRange":"main...HEAD"}'` from a checkout with a `flows.json` +// naming the agent CLI. Each lens is gated on a journaled `artifact_exists` +// check — the worker that spawned the agent journals the files it wrote, and +// the gate reads that journal — and the consensus step on a predicate whose +// verdict is journaled as `agent-N.gate`. The steps run in the invoking +// directory (no `workspace:` scoping: the local agent worker accepts +// stream-only steps, and a "...: readwrite" annotation is refused because +// nothing enforces it). Unlike the other two examples, this one never calls +// `f.human`, so nothing here depends on that verb being wired up. import { flow } from "@relayflows/surface"; @@ -36,6 +41,7 @@ export default flow( const diff = await f .run(`git diff ${input.diffRange}`) .gate((out) => out.trim().length > 0, "nothing to review — the diff is empty"); + await f.run("mkdir -p review"); await Promise.all( LENSES.map((lens) => @@ -45,12 +51,10 @@ export default flow( `Review this diff for ${lens} issues ONLY — ignore everything else. ` + `Write every finding, or an explicit "no issues found", to ` + `${findingsPath(lens)}.\n\n${diff}`, - workspace: "review/: readwrite", }) - .gate( - (r) => r.artifacts.includes(findingsPath(lens)), - `the ${lens} reviewer must write ${findingsPath(lens)}, even to report nothing`, - ), + // A named gate: preflightable by `flows check`, evaluated against + // the artifacts the worker journaled for this step, never the disk. + .gate({ type: "artifact_exists", path: findingsPath(lens) }), ), ); @@ -65,8 +69,9 @@ export default flow( `reached opposite verdicts on the same spot in the diff, resolve it ` + `or mark it UNRESOLVED with both positions. Write your reconciled ` + `verdict to review/consensus.json.`, - workspace: "review/: readwrite", }) + // A predicate gate: author code, run once on the journaled result; its + // verdict is journaled as `agent-N.gate` so resume/replay never re-run it. .gate( (r) => r.artifacts.includes("review/consensus.json"), "the consensus step must write review/consensus.json", diff --git a/packages/schema/flows.schema.json b/packages/schema/flows.schema.json index ad83d6004..2b867d06e 100644 --- a/packages/schema/flows.schema.json +++ b/packages/schema/flows.schema.json @@ -153,7 +153,8 @@ "references_input", "subprocess_gate", "word_count_bounds", - "regex_match" + "regex_match", + "artifact_exists" ] }, "ExitCodeGate": { @@ -390,6 +391,29 @@ ], "additionalProperties": false }, + "ArtifactExistsGate": { + "title": "ArtifactExistsGate", + "description": "Passes when the step's journaled `output.artifacts` lists `path`: a file the\nagent's worker measured as created or changed under its working directory.\nReads the journal, never the disk, so replay and resume see the same verdict.", + "type": "object", + "properties": { + "type": { + "title": "type", + "description": "type in the Relayflows spec.", + "type": "string", + "const": "artifact_exists" + }, + "path": { + "title": "path", + "description": "Working-directory-relative POSIX path, as the worker journals it.", + "type": "string" + } + }, + "required": [ + "type", + "path" + ], + "additionalProperties": false + }, "NamedDataGate": { "title": "NamedDataGate", "description": "NamedDataGate in the Relayflows spec.", @@ -413,6 +437,11 @@ "$ref": "#/$defs/RegexMatchGate", "title": "NamedDataGate alternative 4", "description": "See RegexMatchGate." + }, + { + "$ref": "#/$defs/ArtifactExistsGate", + "title": "NamedDataGate alternative 5", + "description": "See ArtifactExistsGate." } ] }, diff --git a/packages/sdk/src/authored-flow-error.ts b/packages/sdk/src/authored-flow-error.ts index e372f24a5..5bc794a42 100644 --- a/packages/sdk/src/authored-flow-error.ts +++ b/packages/sdk/src/authored-flow-error.ts @@ -23,6 +23,7 @@ export type AuthoredFlowExecutionErrorCode = | 'lease_exceeded' | 'unsupported_completion' | 'unsupported_gate' + | 'gate_failed' | 'unsupported_header' | 'unsettled_derived_work' | 'unsupported_promise_lifecycle' diff --git a/packages/sdk/src/authored-flow-executor.ts b/packages/sdk/src/authored-flow-executor.ts index d23c4c514..464a71eb9 100644 --- a/packages/sdk/src/authored-flow-executor.ts +++ b/packages/sdk/src/authored-flow-executor.ts @@ -211,6 +211,47 @@ export async function executeAuthoredFlow( localAgentStream, budget, definition.header.budget, options.rootRunId, ); + /** + * Predicate gates (docs/SURFACE.md §6). The closure runs here, once, on the + * value the journal handed back; the VERDICT is then journaled as a lowered + * `.gate` deterministic step that succeeds or fails, so a resume or + * replay reads the recorded verdict and never re-runs author code. A false + * verdict fails the step as `gate_failed`, carrying the author's reason. + */ + async function applyPredicateGate(operation: AuthoredFlowOperation, id: string, value: T): Promise { + const gate = operation.predicateGate; + if (gate === undefined) return value; + let verdict: boolean; + let detail: string | undefined; + try { + verdict = gate.predicate(value) === true; + } catch (error) { + verdict = false; + detail = error instanceof Error ? error.message : String(error); + } + const record = JSON.stringify({ + gate: 'predicate', step: id, verdict: verdict ? 'pass' : 'fail', + ...(gate.because === undefined ? {} : { because: gate.because }), + ...(detail === undefined ? {} : { threw: detail }), + }); + const literal = `'${record.replaceAll("'", "'\\''")}'`; + const command = verdict ? `printf '%s' ${literal}` : `printf '%s' ${literal} >&2; exit 1`; + try { + await observeStep(`${id}.gate`, 'deterministic', () => lowerDeterministic(`${id}.gate`, command, false), options.onProgress); + } catch (error) { + if (verdict) throw error; + throw new AuthoredFlowExecutionError( + 'gate_failed', + `step "${id}" failed its predicate gate` + + (gate.because === undefined ? '' : `: ${gate.because}`) + + (detail === undefined ? '' : ` (predicate threw: ${detail})`), + 'verification_failed', + error instanceof AuthoredFlowExecutionError ? error.runId : undefined, + ); + } + return value; + } + function llmOperation(strings: TemplateStringsArray, ...values: unknown[]): Step; function llmOperation(prompt: string, options: LlmOptions): Step; function llmOperation(prompt: string | TemplateStringsArray, ...values: unknown[]): Step { @@ -223,7 +264,7 @@ export async function executeAuthoredFlow( llmOp = new AuthoredFlowOperation( id, 'llm', () => assertOperationAllowed('llm', definition.name, requestedCompletion), - () => observeStep(id, 'llm', () => { + async () => applyPredicateGate(llmOp, id, await observeStep(id, 'llm', () => { if (typeof prompt === 'string') { if (values.length !== 1 || values[0] === undefined) { throw new AuthoredFlowExecutionError('llm_cli_unresolved', 'f.llm(prompt, options) requires an output JSON Schema.'); @@ -233,7 +274,7 @@ export async function executeAuthoredFlow( const text = prompt.reduce((result, part, index) => result + part + (index < values.length ? String(values[index]) : ''), ''); return worker.llm(id, text, undefined, llmOp.namedGate); - }, onProgress), + }, onProgress)), lifecycle, ); return trackStep(authoredSteps, llmOp); @@ -299,7 +340,8 @@ export async function executeAuthoredFlow( id, 'run', () => assertOperationAllowed('run', definition.name, requestedCompletion), - () => observeStep(id, 'deterministic', () => lowerDeterministic(id, command, false, leaseMs, runOp.namedGate), options.onProgress), + async () => applyPredicateGate(runOp, id, + await observeStep(id, 'deterministic', () => lowerDeterministic(id, command, false, leaseMs, runOp.namedGate), options.onProgress)), lifecycle, ); return trackStep(authoredSteps, runOp); @@ -314,7 +356,7 @@ export async function executeAuthoredFlow( id, 'agent', () => assertOperationAllowed('agent', definition.name, requestedCompletion), - () => observeStep(id, 'agent', () => worker.agent(id, options, agentOp.namedGate), onProgress), + async () => applyPredicateGate(agentOp, id, await observeStep(id, 'agent', () => worker.agent(id, options, agentOp.namedGate), onProgress)), lifecycle, ); return trackStep(authoredSteps, agentOp); diff --git a/packages/sdk/src/authored-flow-operation.ts b/packages/sdk/src/authored-flow-operation.ts index faee0776f..4bec64aac 100644 --- a/packages/sdk/src/authored-flow-operation.ts +++ b/packages/sdk/src/authored-flow-operation.ts @@ -8,7 +8,7 @@ import { /** Slice-P kinds the surface `.gate(config)` accepts and the SDK lowers. */ const NAMED_GATE_KINDS = new Set([ - 'references_input', 'subprocess_gate', 'word_count_bounds', 'regex_match', + 'references_input', 'subprocess_gate', 'word_count_bounds', 'regex_match', 'artifact_exists', ]); function isNamedGateConfig(candidate: unknown): candidate is NamedGate { @@ -31,6 +31,14 @@ export class AuthoredFlowOperation { * still throws `unsupported_gate` and never sets this field. */ namedGate: NamedGate | undefined = undefined; + /** + * Predicate gate attached via `.gate(fn, because?)`. Author code: it runs + * in this process after the step completes, and its verdict is journaled as + * a lowered `.gate` deterministic step (docs/SURFACE.md §6), so replay + * and resume see the recorded verdict and never re-run the closure. The + * function itself is never serialized. `flows check` cannot prove it. + */ + predicateGate: { predicate: (value: T) => boolean; because?: string } | undefined = undefined; private state: OperationState = 'created'; private thenInvoked = false; private rootFailureRecorded = false; @@ -61,20 +69,29 @@ export class AuthoredFlowOperation { const operation = this; const step: Step = { gate(configOrPredicate: NamedGate | ((value: T) => boolean), _because?: string): Step { + if (operation.namedGate !== undefined || operation.predicateGate !== undefined) { + throw new AuthoredFlowExecutionError('unsupported_gate', 'a step takes one .gate().'); + } if (isNamedGateConfig(configOrPredicate)) { // Config-object gate: lowers into the compiled StepSpec's // `verification:` field via slice-P named-gate lowering. operation.namedGate = configOrPredicate; return step; } - // Predicate gate: closures cannot be journaled (covenant 1 - // journal-as-truth). Refuse — authors should use a config-object - // gate or the declarative `verification:` block. - throw new AuthoredFlowExecutionError( - 'unsupported_gate', - 'postfix .gate(predicate) closures cannot be journaled; ' - + 'use .gate({type: "…", …}) with a slice-P named gate instead.', - ); + // Predicate gate: runtime control flow. The closure cannot be + // journaled, but its VERDICT can — the executor runs it once after + // the step completes and records pass/fail as a `.gate` step. + if (typeof configOrPredicate !== 'function') { + throw new AuthoredFlowExecutionError( + 'unsupported_gate', + '.gate() takes a named gate config ({type: "…", …}) or a predicate function.', + ); + } + if (_because !== undefined && typeof _because !== 'string') { + throw new AuthoredFlowExecutionError('unsupported_gate', '.gate(predicate, because) takes a string reason.'); + } + operation.predicateGate = { predicate: configOrPredicate, ...(_because === undefined ? {} : { because: _because }) }; + return step; }, then( onfulfilled?: ((value: T) => TResult1 | PromiseLike) | null, diff --git a/packages/sdk/src/authored-worker-step.ts b/packages/sdk/src/authored-worker-step.ts index defa8ff51..63198e1e6 100644 --- a/packages/sdk/src/authored-worker-step.ts +++ b/packages/sdk/src/authored-worker-step.ts @@ -12,7 +12,6 @@ import { isSurfaceCompletionReason, readCompletedStepOutput, readSuccessfulOutpu import type { AuthoredFlowJournalStep } from './authored-flow-executor.js'; import { snapshotJsonValue } from './json-value.js'; import { authoredChildAdmissionKey } from './authored-admission.js'; -import { diffWorkspaceFiles, snapshotWorkspaceFiles } from './agent-artifacts.js'; const WORKSPACE_PERMISSION_ANNOTATION = /:\s*(readonly|readwrite)\s*$/i; @@ -130,20 +129,6 @@ export function authoredWorkerRunner( `f.agent options.transport must be 'direct' or 'relay' (got ${JSON.stringify(options.transport)}).`, ); } - // Artifact detection only tells the truth for the local-agent DIRECT - // path: that is the only case that runs in this same process, on this - // same filesystem, so `options.cwd` (or `process.cwd()`) is provably - // where the CLI actually wrote — a workspace-scoped step never reaches - // here with a local agent attached (refused above). `transport: 'relay'` - // dispatches to agent-relay, which executes on a remote host even - // though a local agent stream is still attached, so it gets no local - // snapshot either. Any other worker attachment may execute on a - // different host entirely; snapshotting this process's filesystem for - // that case would be a guess, not a fact, so `artifacts` stays `[]` - // there, exactly as before this fix. - const artifactRoot = localAgentStream === undefined || options.transport === 'relay' - ? undefined : options.cwd ?? process.cwd(); - const before = artifactRoot === undefined ? undefined : await snapshotWorkspaceFiles(artifactRoot); const output = await run({ id, type: 'agent', instruction: options.task, ...(localAgentStream === undefined ? {} : { surfaces: { streams: [{ stream: localAgentStream }] } }), @@ -158,9 +143,12 @@ export function authoredWorkerRunner( throw new AuthoredFlowExecutionError('journal_protocol_violation', `step "${id}" produced a non-object output`); } const stdout = 'stdout_tail' in output ? output.stdout_tail : undefined; - const artifacts = before === undefined || artifactRoot === undefined - ? [] - : diffWorkspaceFiles(before, await snapshotWorkspaceFiles(artifactRoot)); + // The worker that ran the CLI measured the artifacts and journaled them + // in the step's output; read that fact back rather than re-scanning a + // directory this process may not even share with the agent. + const journaled = 'artifacts' in output ? output.artifacts : undefined; + const artifacts = Array.isArray(journaled) && journaled.every(entry => typeof entry === 'string') + ? [...journaled] : []; return { summary: typeof stdout === 'string' ? stdout : JSON.stringify(output), artifacts }; }, async llm(id: string, prompt: string, options?: LlmOptions, verification?: NamedGate): Promise { diff --git a/packages/sdk/src/cli/direct-run.ts b/packages/sdk/src/cli/direct-run.ts index f1cfccf98..2afa1a313 100644 --- a/packages/sdk/src/cli/direct-run.ts +++ b/packages/sdk/src/cli/direct-run.ts @@ -172,7 +172,7 @@ export async function runDirectFlow( // diagnostic carried up from `classifyOutcome` already names the step, its // exit code and its output tail; this branch is what lets it reach the // terminal. `resumeFlow` takes the same branch, through the same helper. - if (error instanceof AuthoredFlowExecutionError && error.code === 'step_failed') { + if (error instanceof AuthoredFlowExecutionError && (error.code === 'step_failed' || error.code === 'gate_failed')) { return authoredStepFailure('run', base, socketPath, error); } const runId = error instanceof AuthoredFlowExecutionError ? error.runId : undefined; diff --git a/packages/sdk/src/cli/run.ts b/packages/sdk/src/cli/run.ts index 6d0110aa1..2420987ff 100644 --- a/packages/sdk/src/cli/run.ts +++ b/packages/sdk/src/cli/run.ts @@ -223,7 +223,7 @@ export async function resumeFlow( // leaving it on `protocolFailure` meant `flows run` printed the evidence // while `flows resume` still printed `protocol_error` and // `RUN unknown` for the identical failure. - if (error instanceof AuthoredFlowExecutionError && error.code === 'step_failed') { + if (error instanceof AuthoredFlowExecutionError && (error.code === 'step_failed' || error.code === 'gate_failed')) { return authoredStepFailure('resume', base, socketPath, error, runId); } if (!(error instanceof JournalProtocolError) || error.code !== 'run_not_found') { @@ -287,10 +287,12 @@ export function authoredStepFailure( completionReason: 'step_failed', diagnostics: [...base.diagnostics, { severity: 'failure', - kind: 'step_failed', + // A predicate gate that judged false is a run failure with its own + // name, so the report says which kind of check the body did not pass. + kind: error.code === 'gate_failed' ? 'gate_failed' : 'step_failed', // The `step_failed: ` prefix `AuthoredFlowExecutionError` adds is // redundant once the diagnostic is labelled `[step_failed]`. - message: error.message.replace(/^step_failed: /, ''), + message: error.message.replace(/^(?:step_failed|gate_failed): /, ''), }], }, }; diff --git a/packages/sdk/src/failure-kinds.ts b/packages/sdk/src/failure-kinds.ts index 6fbea19ed..5cb48c759 100644 --- a/packages/sdk/src/failure-kinds.ts +++ b/packages/sdk/src/failure-kinds.ts @@ -111,6 +111,8 @@ export const RUN_FAILURE_KINDS = [ 'run_parked', 'run_declined', 'run_unavailable', + /** A predicate `.gate(fn)` judged false; the verdict is journaled as `.gate`. */ + 'gate_failed', ] as const; /** diff --git a/packages/sdk/src/named-gate-lowering.ts b/packages/sdk/src/named-gate-lowering.ts index 750ce9289..479136c28 100644 --- a/packages/sdk/src/named-gate-lowering.ts +++ b/packages/sdk/src/named-gate-lowering.ts @@ -60,6 +60,15 @@ export function lowerNamedGates(steps: readonly StepSpec[]): StepSpec[] { } function gateCommand(gate: NamedDataGate, deterministic: boolean): string { + if (gate.type === 'artifact_exists') { + // The whole output envelope is bound; the verdict is whether the worker's + // journaled `artifacts` list names the path. Nothing on disk is consulted, + // so the journal alone reproduces the verdict on replay and resume. + return `node -e ${quote(`const input=JSON.parse(process.env.FLOWS_INPUT); +const output=input.output; +const artifacts=output!==null&&typeof output==='object'&&Array.isArray(output.artifacts)?output.artifacts:[]; +process.exit(artifacts.includes(${JSON.stringify(gate.path)})?0:1);`)}`; + } const path = gate.type === 'subprocess_gate' ? gate.from_output : gate.type === 'word_count_bounds' ? undefined : gate.in_output_at; // Only compiler-owned code is serialized. Author strings are JSON literals; diff --git a/packages/sdk/src/named-gates.ts b/packages/sdk/src/named-gates.ts index 01d61c18e..20a50b1ad 100644 --- a/packages/sdk/src/named-gates.ts +++ b/packages/sdk/src/named-gates.ts @@ -2,7 +2,7 @@ import { RE2JS } from 're2js'; import type { NamedDataGate, VerificationSpec } from './spec.js'; export const NAMED_GATE_FAILURE_KINDS = [ - 'unknown_gate_kind', 'gate_pattern_invalid', 'gate_command_missing', 'gate_bound_invalid', + 'unknown_gate_kind', 'gate_pattern_invalid', 'gate_command_missing', 'gate_bound_invalid', 'gate_path_invalid', ] as const; export type NamedGateFailureKind = typeof NAMED_GATE_FAILURE_KINDS[number]; @@ -11,6 +11,7 @@ export const NAMED_GATE_KEYS: Record = subprocess_gate: ['type', 'command', 'from_output'], word_count_bounds: ['type', 'min', 'max'], regex_match: ['type', 'pattern', 'in_output_at', 'flags'], + artifact_exists: ['type', 'path'], }; export function isNamedGate(gate: VerificationSpec | undefined): gate is NamedDataGate { @@ -45,6 +46,15 @@ export function namedGateErrors(gate: Record, input: unknown, a errors.push(`${at}.command: gate_command_missing: expected a non-empty shell command`); } break; + case 'artifact_exists': { + const path = gate.path; + const segments = typeof path === 'string' ? path.split('/') : []; + if (typeof path !== 'string' || !path.trim() || path.includes('\0') || path.startsWith('/') + || path !== path.trim() || segments.some(segment => segment === '' || segment === '.' || segment === '..')) { + errors.push(`${at}.path: gate_path_invalid: expected a relative POSIX path without empty, "." or ".." segments`); + } + break; + } case 'word_count_bounds': { const { min, max } = gate; if ([min, max].some(n => n !== undefined && (!Number.isSafeInteger(n) || (n as number) < 0)) diff --git a/packages/sdk/src/spec.ts b/packages/sdk/src/spec.ts index 6ea5798b1..662c0cb2e 100644 --- a/packages/sdk/src/spec.ts +++ b/packages/sdk/src/spec.ts @@ -82,7 +82,18 @@ export interface RegexMatchGate { flags?: string; } -export type NamedDataGate = ReferencesInputGate | SubprocessGate | WordCountBoundsGate | RegexMatchGate; +/** + * Passes when the step's journaled `output.artifacts` lists `path`: a file the + * agent's worker measured as created or changed under its working directory. + * Reads the journal, never the disk, so replay and resume see the same verdict. + */ +export interface ArtifactExistsGate { + type: 'artifact_exists'; + /** Working-directory-relative POSIX path, as the worker journals it. */ + path: string; +} + +export type NamedDataGate = ReferencesInputGate | SubprocessGate | WordCountBoundsGate | RegexMatchGate | ArtifactExistsGate; export type OutputVerificationSpec = OutputContainsGate | JsonSchemaGate | NamedDataGate; export type VerificationSpec = ExitCodeGate | OutputVerificationSpec; diff --git a/packages/sdk/src/validate.ts b/packages/sdk/src/validate.ts index 4b7d44ef9..09b733c1a 100644 --- a/packages/sdk/src/validate.ts +++ b/packages/sdk/src/validate.ts @@ -446,7 +446,7 @@ class Validator { } else if (Object.hasOwn(NAMED_GATE_KEYS, gate.type)) { for (const error of namedGateErrors(v, input, at)) this.fail(error); } else { - this.fail(`${at}.type: unknown_gate_kind: expected exit_code | output_contains | json_schema | references_input | subprocess_gate | word_count_bounds | regex_match`); + this.fail(`${at}.type: unknown_gate_kind: expected exit_code | output_contains | json_schema | references_input | subprocess_gate | word_count_bounds | regex_match | artifact_exists`); } } diff --git a/packages/sdk/src/worker-cli.ts b/packages/sdk/src/worker-cli.ts index 0752964e2..970514ba7 100644 --- a/packages/sdk/src/worker-cli.ts +++ b/packages/sdk/src/worker-cli.ts @@ -1,3 +1,4 @@ +import { diffWorkspaceFiles, snapshotWorkspaceFiles } from './agent-artifacts.js'; import { decodeProviderResult, decodeWrapperResult, requirePricedUsage } from './worker-usage.js'; import { openSidechannel, type SidechannelContext } from './pty-sidechannel.js'; import { spawn } from 'node:child_process'; @@ -38,6 +39,16 @@ export interface WorkerCliResult { exit_code: number | null; stdout_tail: string; stderr_tail: string; + /** + * Files the agent created or changed under its working directory, + * cwd-relative POSIX paths, sorted. Measured by the worker that spawned the + * CLI — the one process provably sharing the agent's filesystem — as a + * content-hash diff of the directory before and after the run, so it is a + * journaled fact rather than a later guess. Present for every direct agent + * execution; absent for `llm` mode and for the relay transport, where the + * agent runs on another host. + */ + artifacts?: string[]; } /** @@ -75,8 +86,19 @@ export async function runAgentCli( return runViaAgentRelay(kind, instruction, wakeContext, effectiveModel, relayContext, cwd, signal); } + // Artifact detection brackets the spawn: the directory the CLI runs in is + // snapshotted before and diffed after, by this process, on this + // filesystem. Only an agent execution writes artifacts; an llm step has no + // workspace to change. + const artifactRoot = mode === 'agent' ? (cwd ?? process.cwd()) : undefined; + const before = artifactRoot === undefined ? undefined : await snapshotWorkspaceFiles(artifactRoot); + const withArtifacts = async (result: WorkerCliResult): Promise => { + if (before === undefined || artifactRoot === undefined) return result; + return { ...result, artifacts: diffWorkspaceFiles(before, await snapshotWorkspaceFiles(artifactRoot)) }; + }; + if (kind === 'relayflows-wrapper-v1') { - return requirePricedUsage(decodeWrapperResult(await runWrapperSession( + return withArtifacts(requirePricedUsage(decodeWrapperResult(await runWrapperSession( cli, instruction, wakeContext, @@ -84,7 +106,7 @@ export async function runAgentCli( wrapperEnvironment(process.env), wrapperLimits, signal, - )), effectiveModel); + )), effectiveModel)); } const env: NodeJS.ProcessEnv = { ...process.env }; @@ -108,7 +130,7 @@ export async function runAgentCli( // Structured provider output carries the authoritative token counts. const args = [...invocation.args]; args.splice(args.length - 1, 0, ...(kind === 'claude' ? ['--output-format', 'json'] : ['--json'])); - return requirePricedUsage(decodeProviderResult(await spawnInvocation(cli, { ...invocation, args }, env, signal, sidechannel, cwd), kind), effectiveModel); + return withArtifacts(requirePricedUsage(decodeProviderResult(await spawnInvocation(cli, { ...invocation, args }, env, signal, sidechannel, cwd), kind), effectiveModel)); } /** Wait under the same worker lease for an authoritative task receipt. */ diff --git a/packages/sdk/src/worker.ts b/packages/sdk/src/worker.ts index 0dae443e4..5b6a124d8 100644 --- a/packages/sdk/src/worker.ts +++ b/packages/sdk/src/worker.ts @@ -133,6 +133,15 @@ export class AgentWorker extends EventEmitter { // emitting an error JSON with exit 0. `completionReason` is // derived from exit code, so a CLI that exits 0 while emitting // `{"error":...}` will report success with an error payload. + // + // `artifacts` (worker-cli.ts) rides inside the CliResult wrapper, so on + // the wrapper path it is journaled in `output`, where the kernel's output + // binding — and therefore an `artifact_exists` gate — can read it. On the + // JSON path the author's object IS the output and is left exactly as the + // agent emitted it: no key is added that a schema, a downstream binding + // or `summary` never asked for. An agent that answers with a JSON object + // therefore journals no artifacts; gate such a step on a deterministic + // check instead. The relay transport reports none (the agent ran elsewhere). const output = result.relay_task?.status === 'completed' && result.exit_code === 0 ? result.relay_task.output : parseJsonOutput(result.stdout_tail) ?? result; diff --git a/packages/sdk/tests/agent-artifacts-live.test.ts b/packages/sdk/tests/agent-artifacts-live.test.ts new file mode 100644 index 000000000..454f73698 --- /dev/null +++ b/packages/sdk/tests/agent-artifacts-live.test.ts @@ -0,0 +1,193 @@ +import { execFileSync, spawnSync } from 'node:child_process'; +import { chmodSync, existsSync, mkdtempSync, readFileSync, rmSync, symlinkSync, writeFileSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import { join, resolve } from 'node:path'; +import { afterEach, describe, expect, it } from 'vitest'; +import { JournalClient } from '../src/journal-client.js'; +import { socketPathFor } from '../src/daemon-connection.js'; + +/** + * The Cloud path, end to end: the built CLI is invoked with the exact argv + * Cloud's relayflow-v2-executor builds (`run --json --data-dir … --local-agent + * --input …`), against a real daemon, with a wrapper CLI that + * actually writes files into its cwd. What is asserted is read back from the + * journal, because that is what a gate, a resume and a replay read. + */ + +const roots: string[] = []; +const sdk = resolve('.'); +const cli = process.env['FLOWS_TEST_CLI'] ?? join(sdk, 'dist/cli.js'); +const wrapperHelper = resolve('../../testdata/preflight/wrapper-session.mjs'); + +function resolveDaemon(): string { + if (process.env['RELAYFLOWD_BIN']) return process.env['RELAYFLOWD_BIN']; + try { + return join(JSON.parse(execFileSync('sh', [ + resolve('../../ops/cargo.sh'), 'metadata', '--format-version=1', '--no-deps', '--locked', '--offline', + ], { cwd: resolve('../../kernel'), encoding: 'utf8', + env: { ...process.env, RELAYFLOWS_NO_TOOLCHAIN_INSTALL: '1' }, + })).target_directory, 'debug', 'relayflowd'); + } catch (cause) { + throw new Error('Live CLI tests require npm run test:prep or an explicit RELAYFLOWD_BIN.', { cause }); + } +} + +afterEach(() => { + for (const root of roots.splice(0)) { + const connection = join(root, '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(root, { recursive: true, force: true }); + } +}); + +/** A wrapper CLI that writes `review/.md` for each lens named in its instruction. */ +function fixture(body: string) { + const relayflowd = resolveDaemon(); + const root = mkdtempSync(join(tmpdir(), 'flows-artifacts-live-')); + roots.push(root); + symlinkSync(join(sdk, 'node_modules'), join(root, 'node_modules')); + const wrapper = join(root, 'agent.mjs'); + writeFileSync(wrapper, `#!/usr/bin/env node +import { receiveWrapperRequest } from ${JSON.stringify(wrapperHelper)}; +import { mkdirSync, writeFileSync } from 'node:fs'; +if (process.argv[2] === 'auth') process.exit(0); +const request = await receiveWrapperRequest(); +if (request) { + for (const lens of request.instruction.match(/write:(\\S+)/g) ?? []) { + const path = lens.slice('write:'.length); + mkdirSync(path.split('/').slice(0, -1).join('/') || '.', { recursive: true }); + writeFileSync(path, 'findings for ' + path + '\\n'); + } + console.log('reviewed'); + process.exit(0); +} +`); + chmodSync(wrapper, 0o755); + writeFileSync(join(root, 'flows.json'), JSON.stringify({ cli: wrapper })); + writeFileSync(join(root, 'package.json'), '{"type":"module"}'); + writeFileSync(join(root, 'review.flow.ts'), `import { flow } from '@relayflows/surface';\nexport default flow('review', async f => {\n${body}\n});\n`); + return { + root, + invoke: () => spawnSync(process.execPath, + [cli, 'run', '--json', '--data-dir', join(root, 'data'), '--local-agent', 'review.flow.ts', '--input', '{}'], + { cwd: root, encoding: 'utf8', timeout: 120_000, env: { ...process.env, RELAYFLOWD_BIN: relayflowd } }), + }; +} + +interface Entry { entry_type: string; step_id?: string; payload: { completionReason?: string; output?: unknown } } + +/** + * Every `step.completed` across the runs named in `runIds`, keyed by step id. + * A successful root lists its child runs in `journalSteps`; a failed root does + * not, so callers hand over the run ids the report named. + */ +async function readJournal(root: string, rootRunId: string, extraRunIds: string[] = []) { + const client = new JournalClient(socketPathFor(join(root, 'data')), { requestTimeoutMs: 5000 }); + await client.connect(); + await client.hello('artifacts-live-test'); + try { + const rootEntries = (await client.journalRead(rootRunId, 1, 1000)).entries as Entry[]; + const rootDone = rootEntries.find(e => e.entry_type === 'step.completed' && e.step_id === 'authored-root'); + const journalSteps = (rootDone?.payload.output as { journalSteps?: Array<{ id: string; runId: string; completionReason: string }> } | undefined)?.journalSteps ?? []; + const completed = new Map(); + for (const runId of [...journalSteps.map(step => step.runId), ...extraRunIds]) { + const entries = (await client.journalRead(runId, 1, 1000)).entries as Entry[]; + for (const entry of entries) { + if (entry.entry_type === 'step.completed' && entry.step_id !== undefined) completed.set(entry.step_id, entry); + } + } + return { journalSteps, completed }; + } finally { + client.close(); + } +} + +/** Run ids a failure report names: the failing step's own run, and any it quotes. */ +function reportedRunIds(report: { runId?: string; diagnostics: Array<{ message: string }> }): string[] { + const ids = new Set(); + if (report.runId) ids.add(report.runId); + for (const d of report.diagnostics) for (const m of d.message.matchAll(/\b(01[0-9A-HJKMNP-TV-Z]{24})\b/g)) ids.add(m[1]!); + return [...ids]; +} + +describe('agent artifacts and gates through the built CLI, a real daemon and the local agent', () => { + it('journals the files the agent wrote, and both artifact gates pass on that journal', async () => { + const f = fixture(` + const security = await f.agent('security-reviewer', { task: 'review write:review/security.md' }) + .gate({ type: 'artifact_exists', path: 'review/security.md' }); + const correctness = await f.agent('correctness-reviewer', { task: 'review write:review/correctness.md' }) + .gate(r => r.artifacts.includes('review/correctness.md'), 'the reviewer must write its findings'); + const both = await f.run('ls review'); + if (!security.artifacts.includes('review/security.md')) throw new Error('security artifacts missing'); + if (!correctness.artifacts.includes('review/correctness.md')) throw new Error('correctness artifacts missing'); + if (!both.includes('security.md') || !both.includes('correctness.md')) throw new Error('files not on disk'); + f.done('success');`); + const result = f.invoke(); + expect(result.status, result.stderr + result.stdout).toBe(0); + const report = JSON.parse(result.stdout) as { runId: string; completionReason: string }; + expect(report.completionReason).toBe('success'); + + const { journalSteps, completed } = await readJournal(f.root, report.runId); + // The named gate is lowered INTO the agent's own kernel run (`agent-1` + + // `agent-1.gate` in one spec); the predicate gate is its own lowered run; + // then the ls and the terminal marker. Every one is a journaled kernel step. + expect(journalSteps.map(s => s.id)).toEqual(['agent-1', 'agent-2', 'agent-2.gate', 'run-3', 'complete-4']); + expect(completed.has('agent-1.gate')).toBe(true); + + // The worker journaled the artifacts in the step output — the fact the gates read. + const agent1 = completed.get('agent-1')!.payload.output as { artifacts?: string[]; stdout_tail?: string }; + expect(agent1.artifacts).toEqual(['review/security.md']); + expect(agent1.stdout_tail).toContain('reviewed'); + expect((completed.get('agent-2')!.payload.output as { artifacts?: string[] }).artifacts).toEqual(['review/correctness.md']); + + // Named gate: a deterministic step reading FLOWS_INPUT, passed. + expect(completed.get('agent-1.gate')!.payload.completionReason).toBe('success'); + // Predicate gate: the recorded verdict, with the author's reason. + const predicate = completed.get('agent-2.gate')!.payload; + expect(predicate.completionReason).toBe('success'); + expect((predicate.output as { stdout_tail: string }).stdout_tail) + .toBe(JSON.stringify({ gate: 'predicate', step: 'agent-2', verdict: 'pass', because: 'the reviewer must write its findings' })); + }, 120_000); + + it('fails the run when the artifact_exists gate names a file the agent did not write', async () => { + const f = fixture(` + await f.agent('lazy-reviewer', { task: 'review write:review/security.md' }) + .gate({ type: 'artifact_exists', path: 'review/performance.md' }); + f.done('success');`); + const result = f.invoke(); + expect(result.status, result.stderr + result.stdout).toBe(1); + const report = JSON.parse(result.stdout) as { runId: string; diagnostics: Array<{ kind: string; message: string }> }; + expect(report.diagnostics.map(d => d.kind)).toContain('step_failed'); + expect(JSON.stringify(report.diagnostics)).toContain('agent-1.gate'); + // The report's run id is the failing step's own kernel run: the agent step + // completed there with its artifacts journaled, and the gate step failed. + const { completed } = await readJournal(f.root, report.runId, reportedRunIds(report)); + expect((completed.get('agent-1')!.payload.output as { artifacts?: string[] }).artifacts).toEqual(['review/security.md']); + expect(completed.get('agent-1.gate')!.payload.completionReason).not.toBe('success'); + }, 120_000); + + it('fails the run with the author reason when a predicate gate returns false, journaling the verdict', async () => { + const f = fixture(` + await f.agent('quiet-reviewer', { task: 'review write:review/security.md' }) + .gate(r => r.artifacts.includes('review/consensus.md'), 'consensus must be written'); + f.done('success');`); + const result = f.invoke(); + expect(result.status, result.stderr + result.stdout).toBe(1); + const report = JSON.parse(result.stdout) as { runId: string; diagnostics: Array<{ kind: string; message: string }> }; + const messages = report.diagnostics.map(d => `${d.kind}: ${d.message}`).join('\n'); + expect(messages).toContain('gate_failed'); + expect(messages).toContain('consensus must be written'); + const { completed } = await readJournal(f.root, report.runId, reportedRunIds(report)); + const verdict = completed.get('agent-1.gate')!.payload; + expect(verdict.completionReason).not.toBe('success'); + expect((verdict.output as { stderr_tail: string }).stderr_tail) + .toContain('"verdict":"fail"'); + }, 120_000); +}); diff --git a/packages/sdk/tests/artifact-gates.test.ts b/packages/sdk/tests/artifact-gates.test.ts new file mode 100644 index 000000000..89da438dc --- /dev/null +++ b/packages/sdk/tests/artifact-gates.test.ts @@ -0,0 +1,130 @@ +import { EventEmitter } from 'node:events'; +import { mkdirSync, mkdtempSync, rmSync, writeFileSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import { afterEach, describe, expect, it, vi } from 'vitest'; + +// A fake `claude` whose "work" is whatever the test's `onSpawn` hook writes +// into the cwd it was spawned in, so artifact detection is exercised without +// a real CLI. +let onSpawn: ((cwd: string | undefined) => void) | undefined; +vi.mock('node:child_process', async () => { + const actual = await vi.importActual('node:child_process'); + return { + ...actual, + spawn: (cli: string, _args: string[], options: Record) => { + const child = new EventEmitter() as EventEmitter & Record; + const stdin = new EventEmitter() as EventEmitter & Record; + stdin.end = () => {}; + stdin.write = (_: unknown, cb: (e?: Error) => void) => { cb(); return true; }; + stdin.destroyed = false; + stdin.writableEnded = false; + const stdout = new EventEmitter(); + child.stdin = stdin; child.stdout = stdout; child.stderr = new EventEmitter(); + child.kill = () => true; + setImmediate(() => { + onSpawn?.(options['cwd'] as string | undefined); + stdout.emit('data', Buffer.from(cli === 'claude' + ? JSON.stringify({ type: 'result', result: 'done', usage: { input_tokens: 1, output_tokens: 1 }, total_cost_usd: 0 }) + : JSON.stringify({ type: 'usage', usage: { input_tokens: 1, output_tokens: 1, total_cost_usd: 0 } }))); + child.emit('close', 0); + }); + return child as unknown as ReturnType; + }, + }; +}); + +import { runAgentCli } from '../src/worker-cli.js'; +import { compileSpec } from '../src/compile.js'; +import { lowerNamedGates } from '../src/named-gate-lowering.js'; +import { namedGateErrors, namedGateFailure } from '../src/named-gates.js'; +import { validateSpec } from '../src/validate.js'; +import { checkFlow } from '../src/cli/check.js'; +import { execFileSync } from 'node:child_process'; + +const dirs: string[] = []; +afterEach(() => { onSpawn = undefined; for (const d of dirs.splice(0)) rmSync(d, { recursive: true, force: true }); }); +function tempDir(): string { const d = mkdtempSync(join(tmpdir(), 'artifact-gates-')); dirs.push(d); return d; } + +describe('worker-side artifacts', () => { + it('journals the files the CLI created or changed in its cwd, content-hashed, dotdirs and node_modules excluded', async () => { + const cwd = tempDir(); + mkdirSync(join(cwd, 'review')); + writeFileSync(join(cwd, 'review/existing.md'), 'v1'); + writeFileSync(join(cwd, 'untouched.txt'), 'same'); + onSpawn = (dir) => { + mkdirSync(join(dir!, 'review'), { recursive: true }); + writeFileSync(join(dir!, 'review/security.md'), 'findings'); + writeFileSync(join(dir!, 'review/existing.md'), 'v2'); + mkdirSync(join(dir!, '.cache'), { recursive: true }); + writeFileSync(join(dir!, '.cache/tmp'), 'x'); + mkdirSync(join(dir!, 'node_modules/dep'), { recursive: true }); + writeFileSync(join(dir!, 'node_modules/dep/index.js'), 'x'); + }; + const result = await runAgentCli('claude', 'review', undefined, 'claude-opus-5', undefined, undefined, 'agent', undefined, cwd); + expect(result.exit_code, result.stderr_tail).toBe(0); + expect(result.artifacts).toEqual(['review/existing.md', 'review/security.md']); + }); + + it('reports an empty list when nothing changed, and none at all for llm mode', async () => { + const cwd = tempDir(); + writeFileSync(join(cwd, 'a.txt'), 'a'); + const agent = await runAgentCli('claude', 'noop', undefined, 'claude-opus-5', undefined, undefined, 'agent', undefined, cwd); + expect(agent.artifacts).toEqual([]); + const llm = await runAgentCli('claude', 'answer', undefined, 'claude-opus-5', undefined, undefined, 'llm', undefined, cwd); + expect(llm.artifacts).toBeUndefined(); + }); +}); + +describe('artifact_exists named gate', () => { + it('validates a relative POSIX path and refuses escapes, absolute paths and NUL', () => { + expect(namedGateErrors({ type: 'artifact_exists', path: 'review/security.md' }, undefined, 'g')).toEqual([]); + for (const path of ['', ' ', '/etc/passwd', '../x', 'a/../b', './x', 'a//b', 'a\0b', 42]) { + const errors = namedGateErrors({ type: 'artifact_exists', path }, undefined, 'g'); + expect(errors, String(path)).toHaveLength(1); + expect(namedGateFailure(errors)).toBe('gate_path_invalid'); + } + const rejected = validateSpec({ version: '0.1.0', name: 'x', steps: [{ id: 'a', type: 'agent', instruction: 'i', cli: 'c', + verification: { type: 'artifact_exists', path: '../x' } }] }); + expect(rejected.ok).toBe(false); + expect(JSON.stringify(rejected)).toContain('gate_path_invalid'); + expect(JSON.stringify(validateSpec({ version: '0.1.0', name: 'x', steps: [{ id: 'a', type: 'agent', instruction: 'i', cli: 'c', + verification: { type: 'nope' } }] }))).toContain('artifact_exists'); + }); + + it('lowers to a deterministic gate step that reads the journaled artifacts, never the disk', () => { + const spec = compileSpec({ version: '0.1.0', name: 'x', steps: [ + { id: 'review', type: 'agent', instruction: 'i', cli: 'c', verification: { type: 'artifact_exists', path: 'review/security.md' } }, + { id: 'after', type: 'deterministic', command: 'true', dependsOn: ['review'] }, + ] }); + const lowered = lowerNamedGates(spec.steps); + expect(lowered.map(s => s.id)).toEqual(['review', 'review.gate', 'after']); + const gate = lowered[1]!; + expect(gate.type).toBe('deterministic'); + expect(gate.dependsOn).toEqual(['review']); + expect(gate.input).toEqual({ output: { step: 'review' } }); + expect(gate.verification).toEqual({ type: 'exit_code' }); + expect(lowered[2]!.dependsOn).toEqual(['review', 'review.gate']); + // The gate command judges FLOWS_INPUT only: present → 0, absent → 1. + const command = (gate as { command: string }).command; + expect(command).not.toContain('readdir'); + expect(command).not.toContain('existsSync'); + const run = (output: unknown) => execFileSync('sh', ['-c', command], { + env: { ...process.env, FLOWS_INPUT: JSON.stringify({ output }) }, stdio: 'pipe', + }); + expect(() => run({ exit_code: 0, stdout_tail: '', artifacts: ['review/security.md'] })).not.toThrow(); + expect(() => run({ exit_code: 0, stdout_tail: '', artifacts: ['other.md'] })).toThrow(); + expect(() => run({ exit_code: 0, stdout_tail: '' })).toThrow(); + expect(() => run({ message: 'a JSON-speaking agent owns its output' })).toThrow(); + }); + + it('is preflightable: flows check prints it as a kernel exit_code gate', () => { + const dir = tempDir(); + writeFileSync(join(dir, 'flows.json'), JSON.stringify({ cli: 'true' })); + writeFileSync(join(dir, 'flow.yaml'), JSON.stringify({ version: '0.1.0', name: 'gated', steps: [ + { id: 'review', type: 'agent', instruction: 'i', cli: 'true', verification: { type: 'artifact_exists', path: 'review/security.md' } }, + ] })); + const report = checkFlow(join(dir, 'flow.yaml')).report; + expect(report.gates).toContainEqual(expect.objectContaining({ stepId: 'review', checks: ['exit_code'], preflightable: true, replayable: true })); + }); +}); diff --git a/packages/sdk/tests/authored-agent-artifacts.test.ts b/packages/sdk/tests/authored-agent-artifacts.test.ts index 04e598526..939a96ba2 100644 --- a/packages/sdk/tests/authored-agent-artifacts.test.ts +++ b/packages/sdk/tests/authored-agent-artifacts.test.ts @@ -35,23 +35,28 @@ function setupProject(): { root: string } { } /** - * A fake kernel that completes every run it is sent immediately: an agent - * step "writes" its file as a side effect of `run.start` (standing in for a - * real CLI, which nothing here spawns — no local-agent worker is attached), - * and every other step type (the executor's own synthetic `f.done()` marker - * step included) completes as a trivial deterministic success. + * A fake kernel that completes every run it is sent immediately. For an agent + * step it journals what a real AgentWorker journals: the CLI wrapper output + * with the worker-measured `artifacts` list (the worker that spawned the CLI + * diffed its cwd; nothing here spawns anything). Every other step type (the + * executor's own synthetic `f.done()` marker step included) completes as a + * trivial deterministic success. `artifactsFor` decides what the journal says + * for the agent step — the disk is never consulted by the authored runtime. */ -function startArtifactServer(sock: string, root: string): Server { +function startArtifactServer(sock: string, root: string, artifactsFor: (transport: unknown) => string[] | undefined = () => ['research/notes.md']): Server { let nextRun = 1; - const stepByRun = new Map(); + const stepByRun = new Map(); return startLoopback(sock, { hello: ctx => sendOk(ctx), 'run.start': (ctx, params) => { const spec = params.spec as Record; const step = (spec['steps'] as Record[])[0]!; const runId = `authored-artifact-run-${nextRun++}`; - stepByRun.set(runId, { id: step['id'] as string, type: step['type'] as string }); + stepByRun.set(runId, { id: step['id'] as string, type: step['type'] as string, + artifacts: step['type'] === 'agent' ? artifactsFor(step['transport']) : undefined }); if (step['type'] === 'agent') { + // The agent "wrote" a file; the disk state is irrelevant to the + // authored runtime, which must read the journal instead. mkdirSync(join(root, 'research'), { recursive: true }); writeFileSync(join(root, 'research', 'notes.md'), 'facts, with sources'); } @@ -67,7 +72,8 @@ function startArtifactServer(sock: string, root: string): Server { completionReason: 'success', disposition: 'step_done', output: step.type === 'agent' - ? { exit_code: 0, stdout_tail: 'wrote research/notes.md', stderr_tail: '' } + ? { exit_code: 0, stdout_tail: 'wrote research/notes.md', stderr_tail: '', + ...(step.artifacts === undefined ? {} : { artifacts: step.artifacts }) } : { exit_code: 0, stdout_tail: '', stderr_tail: '' }, }, }], @@ -88,7 +94,7 @@ describe('f.agent artifacts (local-agent path)', () => { server = undefined; path = undefined; root = undefined; }); - it('reports a file the agent step wrote under its cwd', async () => { + it('reports the artifacts the worker journaled for the agent step', async () => { ({ root } = setupProject()); path = sockPath(); server = startArtifactServer(path, root); @@ -111,10 +117,10 @@ describe('f.agent artifacts (local-agent path)', () => { } }); - it('reports no artifacts for transport: relay even with a local agent stream attached', async () => { + it('reports no artifacts for transport: relay — the worker journals none for a remote agent', async () => { ({ root } = setupProject()); path = sockPath(); - server = startArtifactServer(path, root); + server = startArtifactServer(path, root, transport => transport === 'relay' ? undefined : ['research/notes.md']); const client = new JournalClient(path, { requestTimeoutMs: 2000 }); await client.connect(); await client.hello('authored-artifacts-test-relay'); @@ -138,17 +144,22 @@ describe('f.agent artifacts (local-agent path)', () => { } }); - it('reports no artifacts when no local agent is attached', async () => { + it('reports the journaled artifacts whatever worker ran the step, and [] when it journaled none', async () => { ({ root } = setupProject()); path = sockPath(); - server = startArtifactServer(path, root); + // A worker other than the local agent (a workspace-pinned one here) + // journals its own measurement; a malformed or absent list reads as []. + let calls = 0; + server = startArtifactServer(path, root, () => (calls++ === 0 ? ['research/notes.md'] : undefined)); const client = new JournalClient(path, { requestTimeoutMs: 2000 }); await client.connect(); await client.hello('authored-artifacts-test-2'); try { const handle = flow('artifact-test-no-local', async (f) => { - const result = await f.agent('writer', { task: 'write research/notes.md', cwd: root!, workspace: 'research' }); - expect(result.artifacts).toEqual([]); + const first = await f.agent('writer', { task: 'write research/notes.md', cwd: root!, workspace: 'research' }); + expect(first.artifacts).toEqual(['research/notes.md']); + const second = await f.agent('writer-again', { task: 'write nothing', cwd: root!, workspace: 'research' }); + expect(second.artifacts).toEqual([]); f.done('success'); }); const result = await executeAuthoredFlow(handle, client, undefined, { diff --git a/packages/sdk/tests/authored-flow.test.ts b/packages/sdk/tests/authored-flow.test.ts index 4e4ea80eb..56fd07246 100644 --- a/packages/sdk/tests/authored-flow.test.ts +++ b/packages/sdk/tests/authored-flow.test.ts @@ -140,10 +140,22 @@ describe('authored flow journal executor', () => { async (f) => f.done('success'), ), disconnectedJournal)).rejects.toMatchObject({ code: 'unsupported_header' }); - // Predicate .gate(fn) still refuses — JS closures can't be journaled. - await expect(executeAuthoredFlow(flow('gate-predicate-not-lowered', async (f) => { + // Predicate .gate(fn) is accepted (its VERDICT is journaled as a lowered + // `.gate` run once the step completes), so with a disconnected + // journal it refuses on the run.start path like any other step. What is + // still refused as `unsupported_gate`: a second gate on one step, and a + // `because` that is not a string. + await expect(executeAuthoredFlow(flow('gate-predicate-lowered', async (f) => { await f.run('true').gate(Boolean); f.done('success'); + }), disconnectedJournal)).rejects.not.toMatchObject({ code: 'unsupported_gate' }); + await expect(executeAuthoredFlow(flow('gate-twice', async (f) => { + await f.run('true').gate(Boolean).gate({ type: 'regex_match', pattern: 'x' }); + f.done('success'); + }), disconnectedJournal)).rejects.toMatchObject({ code: 'unsupported_gate' }); + await expect(executeAuthoredFlow(flow('gate-bad-because', async (f) => { + await f.run('true').gate(Boolean, 42 as unknown as string); + f.done('success'); }), disconnectedJournal)).rejects.toMatchObject({ code: 'unsupported_gate' }); // Config-object .gate({...}) MUST NOT throw unsupported_gate — it lowers diff --git a/packages/sdk/tests/preflight.test.ts b/packages/sdk/tests/preflight.test.ts index a0dbdef02..46e12a360 100644 --- a/packages/sdk/tests/preflight.test.ts +++ b/packages/sdk/tests/preflight.test.ts @@ -422,6 +422,8 @@ describe('preflight: CLI resolution and refusal predicates', () => { preflight(flow({ id: 'a', type: 'llm', prompt: 'p', cli: 'x', verification: { type: 'subprocess_gate', command: '/nonexistent-gate-binary-xxx' } }), { probes: probes({ command: (c) => c !== '/nonexistent-gate-binary-xxx' }) }), + preflight(flow({ id: 'a', type: 'agent', instruction: 'i', cli: 'x', + verification: { type: 'artifact_exists', path: '../escape.md' } }), { probes: probes() }), // `scope_syntax_invalid`: grant string doesn't match "mount/path: mode". preflight({ ...flow({ id: 'a', type: 'deterministic', command: 'x' }), workspace: 'not-a-grant' } as FlowSpec, { probes: probes() }), diff --git a/packages/surface/src/context.ts b/packages/surface/src/context.ts index e924ad6b9..d70ff7e8c 100644 --- a/packages/surface/src/context.ts +++ b/packages/surface/src/context.ts @@ -6,6 +6,13 @@ import type { Step } from "./step.js"; export interface AgentResult { summary: string; + /** + * Files the agent created or changed under its working directory, + * cwd-relative POSIX paths, sorted — as journaled by the worker that spawned + * the CLI on the step's `step.completed`, never re-measured later. Empty + * for the relay transport (the agent ran elsewhere) and for an agent whose + * final message is a JSON object (that object is the output, unmodified). + */ artifacts: string[]; } diff --git a/packages/surface/src/step.ts b/packages/surface/src/step.ts index e5bd4def4..a64eebfd5 100644 --- a/packages/surface/src/step.ts +++ b/packages/surface/src/step.ts @@ -8,11 +8,14 @@ export interface Step extends PromiseLike { */ gate(config: NamedGate): Step; /** - * Predicate gates cannot be journaled — the JavaScript closure would not - * survive replay — so the executor refuses this branch with - * `unsupported_gate`. Use a `NamedGate` config-object gate above, or the - * declarative `verification:` field on the compiled step spec, when you - * need journal-honest verification. + * Predicate gate: author code, run once by the runtime after the step + * completes, on the value the journal handed back. The closure itself is + * never serialized; its VERDICT is journaled as a lowered `.gate` + * deterministic step that succeeds or fails, so resume and replay read the + * recorded verdict and never re-run the function. A false verdict fails the + * run as `gate_failed`, carrying `because`. `flows check` cannot prove a + * predicate (docs/SURFACE.md §6); use a `NamedGate` when the check must be + * inspectable before execution. One `.gate()` per step. */ gate(predicate: (value: T) => boolean, because?: string): Step; } @@ -27,7 +30,20 @@ export type NamedGate = | ReferencesInputNamedGate | SubprocessNamedGate | WordCountBoundsNamedGate - | RegexMatchNamedGate; + | RegexMatchNamedGate + | ArtifactExistsNamedGate; + +/** + * Passes when the step's journaled `artifacts` lists `path` — a file the + * agent's worker measured as created or changed under its working directory. + * The check reads the journal, so replay and resume see the verdict that was + * recorded, never a fresh look at the disk. + */ +export interface ArtifactExistsNamedGate { + type: 'artifact_exists'; + /** Working-directory-relative POSIX path, e.g. `review/security.md`. */ + path: string; +} export interface ReferencesInputNamedGate { type: 'references_input'; From 5cb417af99ec75bd0cd260e6ca571472812ac3ff Mon Sep 17 00:00:00 2001 From: Relayflow Lead Date: Thu, 17 Sep 2026 15:26:00 -0700 Subject: [PATCH 2/4] fix(sdk): apply predicate gates to every operation, record verdicts for resume, isolate artifact intervals, run wrappers in cwd MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Review round on #449 (Cursor, Devin), each addressed: - A predicate accepted on a helper/MCP/plugin Step was stored and never evaluated. The applier now lives on the lifecycle and runs in `AuthoredFlowOperation.begin` for every operation kind; a runtime without an applier refuses the gate instead of skipping it. - On resume the closure re-ran and could change the gate step's spec under its admission key. The verdict is now appended to the root run's `predicate-gates` stream before the gate run opens; a resumed body finds the recorded verdict for that step and reuses it, never re-running author code, so the gate spec is identical. - Two agents in one cwd with worker capacity > 1 could attribute each other's writes. The snapshot-spawn-snapshot interval is now serialized per canonical working directory. - Wrapper sessions spawned in the worker's process directory while the scanner measured the requested `cwd`. `cwd` is threaded through `runWrapperSession`/`executePinnedWrapper` to the spawn. Tests: helper-step predicate failing the run; `predicate-gates` stream record read back from the root run; capacity-2 overlap attributing the write to the agent that made it; a real wrapper measured in the requested cwd. (`f.agent({ cwd })` is refused by the kernel spec today — unknown field — so that last one is a worker-level test.) Co-Authored-By: Claude Opus 5 (1M context) --- packages/sdk/src/authored-flow-executor.ts | 87 ++++++++++++++----- packages/sdk/src/authored-flow-lifecycle.ts | 7 ++ packages/sdk/src/authored-flow-operation.ts | 10 ++- packages/sdk/src/worker-cli.ts | 39 ++++++--- packages/sdk/src/wrapper-session.ts | 6 +- .../sdk/tests/agent-artifacts-live.test.ts | 37 ++++++++ packages/sdk/tests/artifact-gates.test.ts | 19 ++++ .../sdk/tests/wrapper-artifacts-cwd.test.ts | 37 ++++++++ 8 files changed, 206 insertions(+), 36 deletions(-) create mode 100644 packages/sdk/tests/wrapper-artifacts-cwd.test.ts diff --git a/packages/sdk/src/authored-flow-executor.ts b/packages/sdk/src/authored-flow-executor.ts index 464a71eb9..ed840f756 100644 --- a/packages/sdk/src/authored-flow-executor.ts +++ b/packages/sdk/src/authored-flow-executor.ts @@ -217,40 +217,80 @@ export async function executeAuthoredFlow( * `.gate` deterministic step that succeeds or fails, so a resume or * replay reads the recorded verdict and never re-runs author code. A false * verdict fails the step as `gate_failed`, carrying the author's reason. + * + * Durability across a resume: the verdict is appended to the root run's + * `predicate-gates` stream BEFORE the gate run is opened. A resumed body + * re-executes and reaches the same gate; it finds the recorded verdict and + * reuses it, so the gate run's spec (which embeds the verdict) is identical + * under its admission key and the closure is never re-run. Without a root + * run there is nothing to resume, and the closure simply runs. */ - async function applyPredicateGate(operation: AuthoredFlowOperation, id: string, value: T): Promise { - const gate = operation.predicateGate; + const PREDICATE_STREAM = 'predicate-gates'; + interface PredicateRecord { gate: 'predicate'; step: string; verdict: 'pass' | 'fail'; because?: string; threw?: string } + let recordedVerdicts: Map | undefined; + async function recordedVerdict(id: string): Promise { + if (options.rootRunId === undefined) return undefined; + if (recordedVerdicts === undefined) { + recordedVerdicts = new Map(); + let offset = 0; + for (;;) { + const page = await journal.streamRead(options.rootRunId, PREDICATE_STREAM, offset, 1000); + for (const message of page.messages) { + const record = (message as { message?: unknown }).message ?? message; + if (typeof record === 'object' && record !== null && (record as PredicateRecord).gate === 'predicate' + && typeof (record as PredicateRecord).step === 'string' + && ((record as PredicateRecord).verdict === 'pass' || (record as PredicateRecord).verdict === 'fail')) { + recordedVerdicts.set((record as PredicateRecord).step, record as PredicateRecord); + } + } + if (page.messages.length === 0 || page.next_offset <= offset) break; + offset = page.next_offset; + } + } + return recordedVerdicts.get(id); + } + async function applyPredicateGate(operation: { id: string; predicateGate: unknown }, value: T): Promise { + const gate = operation.predicateGate as { predicate: (value: T) => boolean; because?: string } | undefined; if (gate === undefined) return value; - let verdict: boolean; - let detail: string | undefined; - try { - verdict = gate.predicate(value) === true; - } catch (error) { - verdict = false; - detail = error instanceof Error ? error.message : String(error); + const id = operation.id; + let record = await recordedVerdict(id); + if (record === undefined) { + let verdict: boolean; + let detail: string | undefined; + try { + verdict = gate.predicate(value) === true; + } catch (error) { + verdict = false; + detail = error instanceof Error ? error.message : String(error); + } + record = { + gate: 'predicate', step: id, verdict: verdict ? 'pass' : 'fail', + ...(gate.because === undefined ? {} : { because: gate.because }), + ...(detail === undefined ? {} : { threw: detail }), + }; + if (options.rootRunId !== undefined) { + await journal.streamAppend(options.rootRunId, PREDICATE_STREAM, record); + recordedVerdicts?.set(id, record); + } } - const record = JSON.stringify({ - gate: 'predicate', step: id, verdict: verdict ? 'pass' : 'fail', - ...(gate.because === undefined ? {} : { because: gate.because }), - ...(detail === undefined ? {} : { threw: detail }), - }); - const literal = `'${record.replaceAll("'", "'\\''")}'`; - const command = verdict ? `printf '%s' ${literal}` : `printf '%s' ${literal} >&2; exit 1`; + const literal = `'${JSON.stringify(record).replaceAll("'", "'\\''")}'`; + const command = record.verdict === 'pass' ? `printf '%s' ${literal}` : `printf '%s' ${literal} >&2; exit 1`; try { await observeStep(`${id}.gate`, 'deterministic', () => lowerDeterministic(`${id}.gate`, command, false), options.onProgress); } catch (error) { - if (verdict) throw error; + if (record.verdict === 'pass') throw error; throw new AuthoredFlowExecutionError( 'gate_failed', `step "${id}" failed its predicate gate` - + (gate.because === undefined ? '' : `: ${gate.because}`) - + (detail === undefined ? '' : ` (predicate threw: ${detail})`), + + (record.because === undefined ? '' : `: ${record.because}`) + + (record.threw === undefined ? '' : ` (predicate threw: ${record.threw})`), 'verification_failed', error instanceof AuthoredFlowExecutionError ? error.runId : undefined, ); } return value; } + lifecycle.applyPredicateGate = applyPredicateGate; function llmOperation(strings: TemplateStringsArray, ...values: unknown[]): Step; function llmOperation(prompt: string, options: LlmOptions): Step; @@ -264,7 +304,7 @@ export async function executeAuthoredFlow( llmOp = new AuthoredFlowOperation( id, 'llm', () => assertOperationAllowed('llm', definition.name, requestedCompletion), - async () => applyPredicateGate(llmOp, id, await observeStep(id, 'llm', () => { + () => observeStep(id, 'llm', () => { if (typeof prompt === 'string') { if (values.length !== 1 || values[0] === undefined) { throw new AuthoredFlowExecutionError('llm_cli_unresolved', 'f.llm(prompt, options) requires an output JSON Schema.'); @@ -274,7 +314,7 @@ export async function executeAuthoredFlow( const text = prompt.reduce((result, part, index) => result + part + (index < values.length ? String(values[index]) : ''), ''); return worker.llm(id, text, undefined, llmOp.namedGate); - }, onProgress)), + }, onProgress), lifecycle, ); return trackStep(authoredSteps, llmOp); @@ -340,8 +380,7 @@ export async function executeAuthoredFlow( id, 'run', () => assertOperationAllowed('run', definition.name, requestedCompletion), - async () => applyPredicateGate(runOp, id, - await observeStep(id, 'deterministic', () => lowerDeterministic(id, command, false, leaseMs, runOp.namedGate), options.onProgress)), + () => observeStep(id, 'deterministic', () => lowerDeterministic(id, command, false, leaseMs, runOp.namedGate), options.onProgress), lifecycle, ); return trackStep(authoredSteps, runOp); @@ -356,7 +395,7 @@ export async function executeAuthoredFlow( id, 'agent', () => assertOperationAllowed('agent', definition.name, requestedCompletion), - async () => applyPredicateGate(agentOp, id, await observeStep(id, 'agent', () => worker.agent(id, options, agentOp.namedGate), onProgress)), + () => observeStep(id, 'agent', () => worker.agent(id, options, agentOp.namedGate), onProgress), lifecycle, ); return trackStep(authoredSteps, agentOp); diff --git a/packages/sdk/src/authored-flow-lifecycle.ts b/packages/sdk/src/authored-flow-lifecycle.ts index ed9a8430c..88ca4c708 100644 --- a/packages/sdk/src/authored-flow-lifecycle.ts +++ b/packages/sdk/src/authored-flow-lifecycle.ts @@ -109,6 +109,13 @@ function collectMembers(values: Iterable): unknown[] | undefined { * called the operation's then method. */ export class AuthoredFlowLifecycle { + /** + * Installed by the executor: runs an operation's predicate gate (if any) on + * its resolved value before the operation fulfills. Lives on the lifecycle + * so EVERY authored operation — core steps, helpers, MCP, plugins — passes + * through it; a gate accepted on a Step must never be silently ignored. + */ + applyPredicateGate: ((operation: { id: string; predicateGate: unknown }, value: T) => Promise) | undefined = undefined; private readonly graph: AuthoredPromiseGraph; private readonly activeResolverProbes: ResolverProbe[] = []; private readonly invocations = new Map(); diff --git a/packages/sdk/src/authored-flow-operation.ts b/packages/sdk/src/authored-flow-operation.ts index 4bec64aac..aea720092 100644 --- a/packages/sdk/src/authored-flow-operation.ts +++ b/packages/sdk/src/authored-flow-operation.ts @@ -161,7 +161,15 @@ export class AuthoredFlowOperation { try { this.assertCanStart(); this.state = 'running'; - const value = await this.start(); + const started = await this.start(); + // The predicate gate, when present, is applied here for every kind of + // operation, so a helper or plugin step cannot carry a gate that never + // runs. The executor installs the applier; without one, a predicate + // gate is refused rather than skipped. + const value = this.predicateGate === undefined ? started + : this.scope.applyPredicateGate === undefined + ? (() => { throw new AuthoredFlowExecutionError('unsupported_gate', 'this runtime cannot apply predicate gates'); })() + : await this.scope.applyPredicateGate(this as unknown as { id: string; predicateGate: unknown }, started); this.state = 'fulfilled'; this.resolve(value); } catch (error) { diff --git a/packages/sdk/src/worker-cli.ts b/packages/sdk/src/worker-cli.ts index 970514ba7..a4247680c 100644 --- a/packages/sdk/src/worker-cli.ts +++ b/packages/sdk/src/worker-cli.ts @@ -1,3 +1,4 @@ +import { resolve } from 'node:path'; import { diffWorkspaceFiles, snapshotWorkspaceFiles } from './agent-artifacts.js'; import { decodeProviderResult, decodeWrapperResult, requirePricedUsage } from './worker-usage.js'; import { openSidechannel, type SidechannelContext } from './pty-sidechannel.js'; @@ -89,16 +90,21 @@ export async function runAgentCli( // Artifact detection brackets the spawn: the directory the CLI runs in is // snapshotted before and diffed after, by this process, on this // filesystem. Only an agent execution writes artifacts; an llm step has no - // workspace to change. - const artifactRoot = mode === 'agent' ? (cwd ?? process.cwd()) : undefined; - const before = artifactRoot === undefined ? undefined : await snapshotWorkspaceFiles(artifactRoot); - const withArtifacts = async (result: WorkerCliResult): Promise => { - if (before === undefined || artifactRoot === undefined) return result; - return { ...result, artifacts: diffWorkspaceFiles(before, await snapshotWorkspaceFiles(artifactRoot)) }; - }; + // workspace to change. Executions sharing a working directory are + // serialized around their snapshot-spawn-snapshot interval, so one agent's + // writes are never attributed to a concurrent one in the same directory. + const artifactRoot = mode === 'agent' ? resolve(cwd ?? process.cwd()) : undefined; + return artifactRoot === undefined + ? execute() + : serializedByDirectory(artifactRoot, async () => { + const before = await snapshotWorkspaceFiles(artifactRoot); + const result = await execute(); + return { ...result, artifacts: diffWorkspaceFiles(before, await snapshotWorkspaceFiles(artifactRoot)) }; + }); + async function execute(): Promise { if (kind === 'relayflows-wrapper-v1') { - return withArtifacts(requirePricedUsage(decodeWrapperResult(await runWrapperSession( + return requirePricedUsage(decodeWrapperResult(await runWrapperSession( cli, instruction, wakeContext, @@ -106,7 +112,8 @@ export async function runAgentCli( wrapperEnvironment(process.env), wrapperLimits, signal, - )), effectiveModel)); + cwd, + )), effectiveModel); } const env: NodeJS.ProcessEnv = { ...process.env }; @@ -130,7 +137,19 @@ export async function runAgentCli( // Structured provider output carries the authoritative token counts. const args = [...invocation.args]; args.splice(args.length - 1, 0, ...(kind === 'claude' ? ['--output-format', 'json'] : ['--json'])); - return withArtifacts(requirePricedUsage(decodeProviderResult(await spawnInvocation(cli, { ...invocation, args }, env, signal, sidechannel, cwd), kind), effectiveModel)); + return requirePricedUsage(decodeProviderResult(await spawnInvocation(cli, { ...invocation, args }, env, signal, sidechannel, cwd), kind), effectiveModel); + } +} + +/** One agent at a time per canonical working directory, for the artifact interval. */ +const directoryQueues = new Map>(); +function serializedByDirectory(directory: string, task: () => Promise): Promise { + const previous = directoryQueues.get(directory) ?? Promise.resolve(); + const run = previous.then(task, task); + const settled = run.then(() => undefined, () => undefined); + directoryQueues.set(directory, settled); + void settled.then(() => { if (directoryQueues.get(directory) === settled) directoryQueues.delete(directory); }); + return run; } /** Wait under the same worker lease for an authoritative task receipt. */ diff --git a/packages/sdk/src/wrapper-session.ts b/packages/sdk/src/wrapper-session.ts index bff6ad018..bb268030a 100644 --- a/packages/sdk/src/wrapper-session.ts +++ b/packages/sdk/src/wrapper-session.ts @@ -49,6 +49,8 @@ export function runWrapperSession( env: NodeJS.ProcessEnv, overrides: Partial = {}, signal?: AbortSignal, + /** Working directory for the wrapper process; the artifact scanner uses the same root. */ + cwd?: string, ): Promise { if (signal?.aborted) return Promise.reject(signal.reason); if (signal !== undefined && process.platform === 'win32') { @@ -77,7 +79,7 @@ export function runWrapperSession( )); } - return executePinnedWrapper(cli, identity, request, env, limits, signal); + return executePinnedWrapper(cli, identity, request, env, limits, signal, cwd); } function executePinnedWrapper( @@ -87,6 +89,7 @@ function executePinnedWrapper( env: NodeJS.ProcessEnv, limits: WrapperSessionLimits, signal?: AbortSignal, + cwd?: string, ): Promise { return new Promise((resolve) => { const ownsGroup = ownsProcessGroup(signal); @@ -94,6 +97,7 @@ function executePinnedWrapper( stdio: ['pipe', 'pipe', 'pipe'], env, detached: ownsGroup, + ...(cwd === undefined ? {} : { cwd }), }); const stop = childStop(child, ownsGroup); const stdout: string[] = []; diff --git a/packages/sdk/tests/agent-artifacts-live.test.ts b/packages/sdk/tests/agent-artifacts-live.test.ts index 454f73698..3dbf90f3a 100644 --- a/packages/sdk/tests/agent-artifacts-live.test.ts +++ b/packages/sdk/tests/agent-artifacts-live.test.ts @@ -191,3 +191,40 @@ describe('agent artifacts and gates through the built CLI, a real daemon and the .toContain('"verdict":"fail"'); }, 120_000); }); + +describe('review follow-ups', () => { + it('applies a predicate gate on a helper step too, and journals its verdict', async () => { + // f.memory / helper steps go through the same lifecycle hook; a false + // predicate on one must fail the run rather than pass unchecked. + const f = fixture(` + await f.run('echo one').gate(out => out.trim() === 'one', 'echo says one'); + const r = await f.agent('writer', { task: 'review write:review/a.md' }); + await f.run('echo two').gate(() => r.artifacts.length === 99, 'never true'); + f.done('success');`); + const result = f.invoke(); + expect(result.status, result.stderr + result.stdout).toBe(1); + const report = JSON.parse(result.stdout) as { runId: string; diagnostics: Array<{ kind: string; message: string }> }; + expect(report.diagnostics.map(d => d.kind)).toContain('gate_failed'); + expect(report.diagnostics.map(d => d.message).join('\n')).toContain('never true'); + }, 120_000); + + it('records predicate verdicts on the root run so a resume reuses them instead of re-running the closure', async () => { + const f = fixture(` + await f.run('echo one').gate(out => out.trim() === 'one', 'echo says one'); + f.done('success');`); + const result = f.invoke(); + expect(result.status, result.stderr + result.stdout).toBe(0); + const report = JSON.parse(result.stdout) as { runId: string }; + const client = new JournalClient(socketPathFor(join(f.root, 'data')), { requestTimeoutMs: 5000 }); + await client.connect(); + await client.hello('predicate-stream-test'); + try { + const page = await client.streamRead(report.runId, 'predicate-gates', 0, 100); + const records = page.messages.map(m => ((m as { message?: unknown }).message ?? m) as Record); + expect(records).toEqual([{ gate: 'predicate', step: 'run-1', verdict: 'pass', because: 'echo says one' }]); + } finally { + client.close(); + } + }, 120_000); + +}); diff --git a/packages/sdk/tests/artifact-gates.test.ts b/packages/sdk/tests/artifact-gates.test.ts index 89da438dc..4983684ad 100644 --- a/packages/sdk/tests/artifact-gates.test.ts +++ b/packages/sdk/tests/artifact-gates.test.ts @@ -76,6 +76,25 @@ describe('worker-side artifacts', () => { }); }); +describe('worker-side artifacts: concurrency', () => { + it('serializes overlapping executions in one cwd so writes are attributed to the agent that made them', async () => { + const cwd = tempDir(); + let calls = 0; + onSpawn = (dir) => { + calls += 1; + // Only the first execution writes; the second must not see that file + // as its own artifact even though both were started together. + if (calls === 1) writeFileSync(join(dir!, 'first.md'), 'first'); + }; + const [a, b] = await Promise.all([ + runAgentCli('claude', 'one', undefined, 'claude-opus-5', undefined, undefined, 'agent', undefined, cwd), + runAgentCli('claude', 'two', undefined, 'claude-opus-5', undefined, undefined, 'agent', undefined, cwd), + ]); + expect(a.artifacts).toEqual(['first.md']); + expect(b.artifacts).toEqual([]); + }); +}); + describe('artifact_exists named gate', () => { it('validates a relative POSIX path and refuses escapes, absolute paths and NUL', () => { expect(namedGateErrors({ type: 'artifact_exists', path: 'review/security.md' }, undefined, 'g')).toEqual([]); diff --git a/packages/sdk/tests/wrapper-artifacts-cwd.test.ts b/packages/sdk/tests/wrapper-artifacts-cwd.test.ts new file mode 100644 index 000000000..189bd364f --- /dev/null +++ b/packages/sdk/tests/wrapper-artifacts-cwd.test.ts @@ -0,0 +1,37 @@ +import { chmodSync, mkdirSync, mkdtempSync, rmSync, writeFileSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import { join, resolve } from 'node:path'; +import { afterEach, describe, expect, it } from 'vitest'; +import { runAgentCli } from '../src/worker-cli.js'; + +const dirs: string[] = []; +afterEach(() => { for (const d of dirs.splice(0)) rmSync(d, { recursive: true, force: true }); }); + +/** A real wrapper CLI (no mocks) that writes `review/here.md` relative to its own cwd. */ +function wrapper(root: string): string { + const path = join(root, 'agent.mjs'); + writeFileSync(path, `#!/usr/bin/env node +import { receiveWrapperRequest } from ${JSON.stringify(resolve('../../testdata/preflight/wrapper-session.mjs'))}; +import { mkdirSync, writeFileSync } from 'node:fs'; +if (process.argv[2] === 'auth') process.exit(0); +const request = await receiveWrapperRequest(); +if (request) { mkdirSync('review', { recursive: true }); writeFileSync('review/here.md', 'x'); console.log('ok'); process.exit(0); } +`); + chmodSync(path, 0o755); + return path; +} + +describe('wrapper artifacts follow the requested cwd', () => { + it('spawns the wrapper in `cwd` and measures artifacts there, not in the worker process directory', async () => { + const root = mkdtempSync(join(tmpdir(), 'wrapper-cwd-')); + dirs.push(root); + const cli = wrapper(root); + const agentDir = join(root, 'elsewhere'); + mkdirSync(agentDir); + const result = await runAgentCli(cli, 'write', undefined, undefined, undefined, undefined, 'agent', undefined, agentDir); + expect(result.exit_code, result.stderr_tail).toBe(0); + expect(result.artifacts).toEqual(['review/here.md']); + // The file is where the agent ran, not where this test process runs. + expect(() => rmSync(join(agentDir, 'review/here.md'))).not.toThrow(); + }); +}); From 66dbf1298cbde1f562995e7b618c17480b0ae3d0 Mon Sep 17 00:00:00 2001 From: Relayflow Lead Date: Thu, 17 Sep 2026 15:37:27 -0700 Subject: [PATCH 3/4] fix(sdk): share one predicate-verdict stream load across concurrent gates on resume MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `recordedVerdict` assigned an empty map before its `stream.read` resolved, so a second `.gate(fn)` in flight at the same time (Promise.all) saw the map as empty, re-ran its closure and appended a second record — and a different verdict would have changed the lowered gate command under its existing admission key. The in-flight load promise is now what is memoized; every concurrent gate awaits the same read. Test: two parallel gates on a resumed root with slow-answered recorded verdicts — closures never called, nothing appended, both gate specs carry the recorded pass; red on the previous source, green now. Co-Authored-By: Claude Opus 5 (1M context) --- packages/sdk/src/authored-flow-executor.ts | 44 +++++++----- .../tests/authored-agent-artifacts.test.ts | 71 +++++++++++++++++++ 2 files changed, 96 insertions(+), 19 deletions(-) diff --git a/packages/sdk/src/authored-flow-executor.ts b/packages/sdk/src/authored-flow-executor.ts index ed840f756..5f1ee722c 100644 --- a/packages/sdk/src/authored-flow-executor.ts +++ b/packages/sdk/src/authored-flow-executor.ts @@ -227,27 +227,33 @@ export async function executeAuthoredFlow( */ const PREDICATE_STREAM = 'predicate-gates'; interface PredicateRecord { gate: 'predicate'; step: string; verdict: 'pass' | 'fail'; because?: string; threw?: string } - let recordedVerdicts: Map | undefined; - async function recordedVerdict(id: string): Promise { - if (options.rootRunId === undefined) return undefined; - if (recordedVerdicts === undefined) { - recordedVerdicts = new Map(); - let offset = 0; - for (;;) { - const page = await journal.streamRead(options.rootRunId, PREDICATE_STREAM, offset, 1000); - for (const message of page.messages) { - const record = (message as { message?: unknown }).message ?? message; - if (typeof record === 'object' && record !== null && (record as PredicateRecord).gate === 'predicate' - && typeof (record as PredicateRecord).step === 'string' - && ((record as PredicateRecord).verdict === 'pass' || (record as PredicateRecord).verdict === 'fail')) { - recordedVerdicts.set((record as PredicateRecord).step, record as PredicateRecord); - } + // The stream is read once per execution; concurrent gates (Promise.all) + // share the single in-flight load, so none of them can observe an empty + // map while the read is still pending and re-run a closure whose verdict + // was already recorded. + let recordedVerdicts: Promise> | undefined; + async function loadRecordedVerdicts(rootRunId: string): Promise> { + const verdicts = new Map(); + let offset = 0; + for (;;) { + const page = await journal.streamRead(rootRunId, PREDICATE_STREAM, offset, 1000); + for (const message of page.messages) { + const record = (message as { message?: unknown }).message ?? message; + if (typeof record === 'object' && record !== null && (record as PredicateRecord).gate === 'predicate' + && typeof (record as PredicateRecord).step === 'string' + && ((record as PredicateRecord).verdict === 'pass' || (record as PredicateRecord).verdict === 'fail')) { + verdicts.set((record as PredicateRecord).step, record as PredicateRecord); } - if (page.messages.length === 0 || page.next_offset <= offset) break; - offset = page.next_offset; } + if (page.messages.length === 0 || page.next_offset <= offset) break; + offset = page.next_offset; } - return recordedVerdicts.get(id); + return verdicts; + } + async function recordedVerdict(id: string): Promise { + if (options.rootRunId === undefined) return undefined; + recordedVerdicts ??= loadRecordedVerdicts(options.rootRunId); + return (await recordedVerdicts).get(id); } async function applyPredicateGate(operation: { id: string; predicateGate: unknown }, value: T): Promise { const gate = operation.predicateGate as { predicate: (value: T) => boolean; because?: string } | undefined; @@ -270,7 +276,7 @@ export async function executeAuthoredFlow( }; if (options.rootRunId !== undefined) { await journal.streamAppend(options.rootRunId, PREDICATE_STREAM, record); - recordedVerdicts?.set(id, record); + (await recordedVerdicts)?.set(id, record); } } const literal = `'${JSON.stringify(record).replaceAll("'", "'\\''")}'`; diff --git a/packages/sdk/tests/authored-agent-artifacts.test.ts b/packages/sdk/tests/authored-agent-artifacts.test.ts index 939a96ba2..e9aeb759e 100644 --- a/packages/sdk/tests/authored-agent-artifacts.test.ts +++ b/packages/sdk/tests/authored-agent-artifacts.test.ts @@ -171,3 +171,74 @@ describe('f.agent artifacts (local-agent path)', () => { } }); }); + +describe('predicate verdicts recorded on the root run', () => { + let server: Server | undefined; + let path: string | undefined; + afterEach(async () => { + if (server !== undefined) await new Promise(resolve => server!.close(() => resolve())); + if (path !== undefined) rmSync(path, { force: true }); + server = undefined; path = undefined; + }); + + it('concurrent gates on resume share one stream load, reuse the recorded verdicts and never re-run a closure', async () => { + path = sockPath(); + // A root run that already carries verdicts for both gates; the stream + // read is answered slowly so both gates are in flight before it resolves. + const recorded = [ + { gate: 'predicate', step: 'run-1', verdict: 'pass', because: 'first' }, + { gate: 'predicate', step: 'run-2', verdict: 'pass', because: 'second' }, + ]; + let streamReads = 0; + const appended: unknown[] = []; + let nextRun = 1; + const specs = new Map(); + server = startLoopback(path, { + hello: ctx => sendOk(ctx), + 'stream.read': (ctx, params) => { + streamReads += 1; + const from = params.from_offset as number; + setTimeout(() => sendResult(ctx, from === 0 + ? { messages: recorded.map(message => ({ message })), next_offset: recorded.length } + : { messages: [], next_offset: from }), 50); + }, + 'stream.append': (ctx, params) => { appended.push(params.message); sendResult(ctx, { offset: appended.length }); }, + 'run.start': (ctx, params) => { + const spec = params.spec as { steps: Array<{ id: string; command?: string }> }; + const runId = `resume-run-${nextRun++}`; + specs.set(runId, spec.steps[0]!.id + '|' + (spec.steps[0]!.command ?? '')); + sendResult(ctx, { run_id: runId, status: 'completed', completion_reason: 'success', completed_steps: 1 }); + }, + 'journal.read': (ctx, params) => { + const [id] = specs.get(params.run_id as string)!.split('|'); + sendResult(ctx, { entries: [{ entry_type: 'step.completed', step_id: id, payload: { + completionReason: 'success', disposition: 'step_done', output: { exit_code: 0, stdout_tail: 'x', stderr_tail: '' } } }] }); + }, + }); + const client = new JournalClient(path, { requestTimeoutMs: 2000 }); + await client.connect(); + await client.hello('predicate-resume-test'); + let closureCalls = 0; + try { + const handle = flow('predicate-resume', async (f) => { + await Promise.all([ + f.run('echo a').gate(() => { closureCalls += 1; return false; }, 'first'), + f.run('echo b').gate(() => { closureCalls += 1; return false; }, 'second'), + ]); + f.done('success'); + }); + const result = await executeAuthoredFlow(handle, client, undefined, { rootRunId: 'root-1' }); + expect(result.completionReason).toBe('success'); + } finally { + client.close(); + } + // The closures would have said "fail"; the recorded "pass" verdicts won, + // through a single stream load, and nothing new was appended. + expect(closureCalls).toBe(0); + expect(streamReads).toBeLessThanOrEqual(2); + expect(appended).toEqual([]); + const gateSpecs = [...specs.values()].filter(s => s.includes('.gate|')); + expect(gateSpecs).toHaveLength(2); + for (const spec of gateSpecs) expect(spec).toContain('"verdict":"pass"'); + }); +}); From e2ac9315da27552f0183b4141df8139d76de8846 Mon Sep 17 00:00:00 2001 From: Relayflow Lead Date: Thu, 17 Sep 2026 16:23:06 -0700 Subject: [PATCH 4/4] =?UTF-8?q?fix(sdk):=20standalone=20result=20verifier?= =?UTF-8?q?=20accepts=20gate=20children=20=E2=80=94=20predicate=20`.?= =?UTF-8?q?gate`=20runs=20and=20lowered=20named-gate=20steps?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The Bun/Node standalone verifier (`verifyAuthoredNodeResult`) required `journalSteps.length === N` from `complete-N` and every id to end in `-`, so a predicate-gated flow — whose gate is journaled as its own `.gate` child run — was refused as having no durable completion and the root terminalized `worker_error`. Missed locally because the Bun suite was skipped (Bun 1.3.14 vs pinned 1.4.0). Decision, per SURFACE.md §6: a gate is subordinate to the step it judges. A named gate lowers to `.gate` INSIDE the step's spec and has never consumed an ordinal; the predicate gate keeps the same `.gate` shape as a separate run only because the closure must run between the step's completion and the gate's. `complete-N` counts the operations the author wrote, so `.gate` ids are set aside from the count and contiguity check — not exempted from verification: every gate must name a claimed parent, cannot be the terminal, and is held to the same durable completion evidence as any other child (completed run, `done` state, `step.completed` success, spec named `/`). While running the suite under a real Bun 1.4.0, a second pre-existing gap in the same verifier surfaced: it assumed one step per child spec, but a NAMED gate lowers a second `.gate` step into that spec, so any named-gated authored step was refused too. The verifier now accepts `[step, step.gate]` and requires the gate step to have completed. Tests: unit (no Bun) — predicate gate accepted and verified; orphan gate, non-success gate, gate without durable completion, gate as terminal, and a count that includes gates all refused; named-gate spec shape accepted with a completed gate and refused otherwise or with a foreign second step. Runtime — new `authored-node-runtime` case runs a predicate-gated + named-gated flow through the standalone CLI and resumes it; the whole Bun suite passes under Bun 1.4.0 (`FLOWS_BUILD_BUN`). Co-Authored-By: Claude Opus 5 (1M context) --- packages/sdk/src/authored-node-runner.ts | 48 ++++++++++++++----- .../sdk/tests/authored-node-result.test.ts | 48 +++++++++++++++++++ .../sdk/tests/authored-node-runtime.test.ts | 14 ++++++ 3 files changed, 97 insertions(+), 13 deletions(-) diff --git a/packages/sdk/src/authored-node-runner.ts b/packages/sdk/src/authored-node-runner.ts index e7ac7d9fe..d646635ce 100644 --- a/packages/sdk/src/authored-node-runner.ts +++ b/packages/sdk/src/authored-node-runner.ts @@ -161,30 +161,44 @@ 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'); }; + const invalid = (why = ''): never => { throw new Error(`authored runtime result has no matching durable completion${process.env['FLOWS_VERIFIER_DEBUG'] && why ? ` (${why})` : ''}`); }; if (result.rootRunId !== rootRunId || result.name !== metadata.flowName || !isLoweredCompletion(result.completionReason) - || !Array.isArray(result.journalSteps) || result.journalSteps.length === 0) invalid(); + || !Array.isArray(result.journalSteps) || result.journalSteps.length === 0) invalid('frame'); const terminal = result.journalSteps.at(-1)!; - if (!terminal || !/^complete-[1-9][0-9]*$/.test(terminal.id)) invalid(); + if (!terminal || !/^complete-[1-9][0-9]*$/.test(terminal.id)) invalid('terminal'); const count = Number(terminal.id.slice('complete-'.length)); - if (!Number.isSafeInteger(count) || result.journalSteps.length !== count) invalid(); + // A predicate gate is journaled as `.gate`: a child run subordinate + // to the authored step it judges, in the same `.gate` shape a named + // gate lowers to inside its step's own spec. Neither consumes an ordinal — + // `complete-N` counts the operations the author wrote (SURFACE.md §6) — + // so gates are set aside from the count and the contiguity check, and + // verified separately: every gate must name a claimed parent, and is then + // held to the same durable-completion evidence as any other child run. + const isGate = (id: string): boolean => /\.gate$/.test(id); + const authored = result.journalSteps.filter(step => typeof step?.id === 'string' && !isGate(step.id)); + const gates = result.journalSteps.filter(step => typeof step?.id === 'string' && isGate(step.id)); + if (!Number.isSafeInteger(count) || authored.length !== count || isGate(terminal.id)) invalid(`count ${authored.length}!=${count}`); 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 ordered = [...authored].sort((a,b)=>ordinal(a.id)-ordinal(b.id)); + const authoredIds = new Set(ordered.map(step => step.id)); + for (const gate of gates) { + if (!authoredIds.has(gate.id.slice(0, -'.gate'.length))) invalid(`orphan ${gate.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()) { + for (const [index, claimed] of [...ordered, ...gates].entries()) { if (!claimed || typeof claimed.id !== 'string' || typeof claimed.runId !== 'string' - || ordinal(claimed.id) !== index+1 || claimed.completionReason !== 'success' - || runs.has(claimed.runId)) invalid(); + || (index < ordered.length && ordinal(claimed.id) !== index+1) || claimed.completionReason !== 'success' + || runs.has(claimed.runId)) invalid(`claim ${claimed?.id}`); 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(); + || state.steps[claimed.id]?.state !== 'done') invalid(`state ${claimed.id} ${state.status} ${state.steps[claimed.id]?.state}`); 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}> }; @@ -196,7 +210,7 @@ export async function verifyAuthoredNodeResult( 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(); + if (!Number.isSafeInteger(entry?.seq) || entry.seq < fromSeq) invalid('seq'); 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); @@ -204,12 +218,20 @@ export async function verifyAuthoredNodeResult( } const spec = entries.find(entry => entry.entry_type === 'run.spawned')?.payload?.spec; const step = spec?.steps?.[0]; + // A named gate lowers INTO the step's own spec as a second, dependent + // `.gate` step (named-gate-lowering.ts); that is the only other + // step a child spec may carry, and it must have completed too. + const lowered = spec?.steps ?? []; + const specShape = lowered.length === 1 || (lowered.length === 2 && lowered[1]?.id === `${claimed.id}.gate`); const completed = entries.filter(entry => entry.entry_type === 'step.completed' && entry.step_id === claimed.id); + const gateCompleted = lowered.length === 2 + ? entries.filter(entry => entry.entry_type === 'step.completed' && entry.step_id === `${claimed.id}.gate`) : []; const terminalFacts = entries.filter(entry => entry.entry_type === 'run.completed'); - if (spec?.name !== `${metadata.flowName}/${claimed.id}` || spec?.steps?.length !== 1 + if (spec?.name !== `${metadata.flowName}/${claimed.id}` || !specShape || step?.id !== claimed.id || completed.length !== 1 || completed[0]?.payload?.completionReason !== 'success' - || terminalFacts.length !== 1 || terminalFacts[0]?.payload?.completionReason !== 'success') invalid(); + || (lowered.length === 2 && (gateCompleted.length !== 1 || gateCompleted[0]?.payload?.completionReason !== 'success')) + || terminalFacts.length !== 1 || terminalFacts[0]?.payload?.completionReason !== 'success') invalid(`evidence ${claimed.id} spec=${spec?.name} step=${step?.id} completed=${completed.length}`); if (claimed === terminal) { // The claimed verdict must match the marker the journal actually // recorded, so an IPC frame cannot claim `success` over a run whose @@ -219,7 +241,7 @@ export async function verifyAuthoredNodeResult( // the runtime validation of untrusted IPC, and it is why nothing has // to be re-asserted here just to satisfy the type. if (step?.type !== 'deterministic' - || step.command !== completionMarker(result.completionReason)) invalid(); + || step.command !== completionMarker(result.completionReason)) invalid('marker'); } } } finally { journal.close(); } diff --git a/packages/sdk/tests/authored-node-result.test.ts b/packages/sdk/tests/authored-node-result.test.ts index 31912df07..3e4c30473 100644 --- a/packages/sdk/tests/authored-node-result.test.ts +++ b/packages/sdk/tests/authored-node-result.test.ts @@ -36,6 +36,54 @@ beforeEach(()=>{ })); }); describe('authored IPC result durable verification',()=>{ + it('accepts a predicate-gated flow: the `.gate` child is verified but does not consume an ordinal',async()=>{ + const claimed:AuthoredFlowExecutionResult={rootRunId:'root',name:'example',completionReason:'success',journalSteps:[ + {id:'run-1',runId:'child-1',completionReason:'success'}, + {id:'run-1.gate',runId:'child-gate',completionReason:'success'}, + {id:'complete-2',runId:'child-2',completionReason:'success'}, + ]}; + records.set('child-gate',entries(`printf '%s' '{"gate":"predicate","step":"run-1","verdict":"pass"}'`,'run-1.gate')); + mocks.runGet.mockImplementation(async (run_id:string)=>({run_id,status:'completed',steps:{ + [run_id==='child-1'?'run-1':run_id==='child-gate'?'run-1.gate':'complete-2']:{state:'done'}, + }})); + await verifyAuthoredNodeResult(claimed,metadata,'root','socket'); + expect(mocks.runGet).toHaveBeenCalledTimes(3); + expect(mocks.journalRead).toHaveBeenCalledWith('child-gate',1,100); + }); + it.each([ + ['an orphan gate whose parent was not claimed', [{id:'run-9.gate',runId:'child-gate',completionReason:'success'}]], + ['a gate claimed without a success completion', [{id:'run-1.gate',runId:'child-gate',completionReason:'step_failed'}]], + ['a gate whose child run has no durable completion', [{id:'run-1.gate',runId:'child-missing',completionReason:'success'}]], + ] as const)('refuses %s',async(_label,gateSteps)=>{ + const base=result(); + const claimed:AuthoredFlowExecutionResult={...base,journalSteps:[base.journalSteps[0]!,...gateSteps,base.journalSteps[1]!]}; + records.set('child-gate',entries(':','run-1.gate')); + mocks.runGet.mockImplementation(async (run_id:string)=>{ + if(run_id==='child-missing') return {run_id,status:'running',steps:{}}; + return {run_id,status:'completed',steps:{[run_id==='child-1'?'run-1':run_id==='child-gate'?'run-1.gate':'complete-2']:{state:'done'}}}; + }); + await expect(verifyAuthoredNodeResult(claimed,metadata,'root','socket')).rejects.toThrow('no matching durable completion'); + }); + it('accepts a child spec that carries its lowered named gate as a second, completed step — and refuses one where the gate did not complete',async()=>{ + const namedGated=(gateReason:string)=>[ + {entry_type:'run.spawned',payload:{spec:{name:'example/run-1',steps:[{id:'run-1',type:'deterministic',command:':'},{id:'run-1.gate',type:'deterministic',command:'node -e 1'}]}}}, + {entry_type:'step.completed',step_id:'run-1',payload:{completionReason:'success'}}, + {entry_type:'step.completed',step_id:'run-1.gate',payload:{completionReason:gateReason}}, + {entry_type:'run.completed',payload:{completionReason:'success'}}, + ]; + records.set('child-1',namedGated('success')); + await verifyAuthoredNodeResult(result(),metadata,'root','socket'); + records.set('child-1',namedGated('verification_failed')); + await expect(verifyAuthoredNodeResult(result(),metadata,'root','socket')).rejects.toThrow('no matching durable completion'); + records.set('child-1',[{entry_type:'run.spawned',payload:{spec:{name:'example/run-1',steps:[{id:'run-1',type:'deterministic',command:':'},{id:'other',type:'deterministic',command:':'}]}}}, + {entry_type:'step.completed',step_id:'run-1',payload:{completionReason:'success'}},{entry_type:'run.completed',payload:{completionReason:'success'}}]); + await expect(verifyAuthoredNodeResult(result(),metadata,'root','socket')).rejects.toThrow('no matching durable completion'); + }); + it('refuses a gate id as the terminal marker, and a count that includes gates',async()=>{ + const base=result(); + await expect(verifyAuthoredNodeResult({...base,journalSteps:[base.journalSteps[0]!,{id:'complete-2.gate',runId:'x',completionReason:'success'}]},metadata,'root','socket')).rejects.toThrow('no matching durable completion'); + await expect(verifyAuthoredNodeResult({...base,journalSteps:[base.journalSteps[0]!,{id:'run-1.gate',runId:'child-gate',completionReason:'success'},{id:'complete-3',runId:'child-2',completionReason:'success'}]},metadata,'root','socket')).rejects.toThrow('no matching durable completion'); + }); 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,100); diff --git a/packages/sdk/tests/authored-node-runtime.test.ts b/packages/sdk/tests/authored-node-runtime.test.ts index 0bf39b6d7..7efc543e8 100644 --- a/packages/sdk/tests/authored-node-runtime.test.ts +++ b/packages/sdk/tests/authored-node-runtime.test.ts @@ -109,6 +109,20 @@ 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('accepts a predicate-gated flow: the `.gate` child is journaled, verified, and not counted as an authored step', async () => { + const f = fixture(`await f.run("printf one").gate(out => out === 'one', 'echo says one'); +await f.run("printf two").gate({ type: 'word_count_bounds', min: 1, max: 1 }); +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'}); + const rootEntries = await entries(f.directory, report.runId); + const output = rootEntries.find(e=>e.entry_type==='step.completed')!.payload['output']; + // The predicate gate is its own child run (`run-1.gate`); the named gate + // lowers inside `run-2`'s spec; `complete-3` counts the two authored steps. + expect(output.journalSteps.map((s:{id:string})=>s.id)).toEqual(['run-1','run-1.gate','run-2','complete-3']); + const resumed = f.invoke(['resume', report.runId]); expect(resumed.status, resumed.stderr + resumed.stdout).toBe(0); + }, 90_000); + it.each([ ['SIGKILL', 'success'], ['SIGTERM', 'success'], ['blocked-SIGKILL', 'success'], ['SIGKILL', 'declined'],