diff --git a/packages/sdk/src/authored-flow-executor.ts b/packages/sdk/src/authored-flow-executor.ts index f04b868b8..3e52d2e37 100644 --- a/packages/sdk/src/authored-flow-executor.ts +++ b/packages/sdk/src/authored-flow-executor.ts @@ -13,6 +13,7 @@ import { assertMemoryReachable, authoredMemory, scriptMemoryScope } from './auth import { authoredDeterministicRunner, authoredWorkerRunner } from './authored-worker-step.js'; import { isSurfaceRunCompletionReason } from './authored-step-output.js'; import { + type AgentResult, type LlmOptions, type CloudHelper, type CompletionReason as SurfaceCompletionReason, @@ -181,7 +182,11 @@ export async function executeAuthoredFlow( function llmOperation(prompt: string | TemplateStringsArray, ...values: unknown[]): Step { assertOperationAllowed('llm', definition.name, requestedCompletion); const id = `llm-${nextStep++}`; - return trackStep(authoredSteps, new AuthoredFlowOperation( + // Hoist the operation reference so the start closure can read the + // caller's `.gate(config)` at spec-build time. AuthoredFlowOperation + // begins after construction returns, so `llmOp` is defined by then. + let llmOp!: AuthoredFlowOperation; + llmOp = new AuthoredFlowOperation( id, 'llm', () => assertOperationAllowed('llm', definition.name, requestedCompletion), () => observeStep(id, 'llm', () => { @@ -189,14 +194,15 @@ export async function executeAuthoredFlow( if (values.length !== 1 || values[0] === undefined) { throw new AuthoredFlowExecutionError('llm_cli_unresolved', 'f.llm(prompt, options) requires an output JSON Schema.'); } - return worker.llm(id, prompt, values[0] as LlmOptions); + return worker.llm(id, prompt, values[0] as LlmOptions, llmOp.namedGate); } const text = prompt.reduce((result, part, index) => result + part + (index < values.length ? String(values[index]) : ''), ''); - return worker.llm(id, text); + return worker.llm(id, text, undefined, llmOp.namedGate); }, onProgress), lifecycle, - )); + ); + return trackStep(authoredSteps, llmOp); } const slackRun = randomUUID(); @@ -254,26 +260,30 @@ export async function executeAuthoredFlow( assertOperationAllowed('run', definition.name, requestedCompletion); const leaseMs = runOptions?.timeout === undefined ? undefined : parseStepTimeout(runOptions.timeout); const id = `run-${nextStep++}`; - return trackStep(authoredSteps, new AuthoredFlowOperation( + let runOp!: AuthoredFlowOperation; + runOp = new AuthoredFlowOperation( id, 'run', () => assertOperationAllowed('run', definition.name, requestedCompletion), - () => observeStep(id, 'deterministic', () => lowerDeterministic(id, command, false, leaseMs), options.onProgress), + () => observeStep(id, 'deterministic', () => lowerDeterministic(id, command, false, leaseMs, runOp.namedGate), options.onProgress), lifecycle, - )); + ); + return trackStep(authoredSteps, runOp); }, llm: llmOperation, agent(name, options) { assertOperationAllowed('agent', definition.name, requestedCompletion); void name; // Authored headers do not yet declare reusable named agents. const id = `agent-${nextStep++}`; - return trackStep(authoredSteps, new AuthoredFlowOperation( + let agentOp!: AuthoredFlowOperation; + agentOp = new AuthoredFlowOperation( id, 'agent', () => assertOperationAllowed('agent', definition.name, requestedCompletion), - () => observeStep(id, 'agent', () => worker.agent(id, options), onProgress), + () => observeStep(id, 'agent', () => worker.agent(id, options, agentOp.namedGate), onProgress), lifecycle, - )); + ); + return trackStep(authoredSteps, agentOp); }, human() { assertOperationAllowed('human', definition.name, requestedCompletion); diff --git a/packages/sdk/src/authored-flow-operation.ts b/packages/sdk/src/authored-flow-operation.ts index 92d662114..faee0776f 100644 --- a/packages/sdk/src/authored-flow-operation.ts +++ b/packages/sdk/src/authored-flow-operation.ts @@ -1,17 +1,36 @@ import { executionAsyncId } from 'node:async_hooks'; -import type { Step } from '@relayflows/surface'; +import type { NamedGate, Step } from '@relayflows/surface'; import { AuthoredFlowExecutionError } from './authored-flow-error.js'; import { AuthoredFlowLifecycle, type AuthoredOperationInvocation, } from './authored-flow-lifecycle.js'; +/** 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', +]); + +function isNamedGateConfig(candidate: unknown): candidate is NamedGate { + return typeof candidate === 'object' && candidate !== null && !Array.isArray(candidate) + && typeof (candidate as { type?: unknown }).type === 'string' + && NAMED_GATE_KINDS.has((candidate as { type: string }).type); +} + type OperationState = 'created' | 'running' | 'fulfilled' | 'rejected'; const nativePromiseThen = Promise.prototype.then; /** A root authored operation whose outcome cannot be hidden by promise handlers. */ export class AuthoredFlowOperation { readonly step: Step; + /** + * Named-gate config attached via `.gate({type: '…', …})`. The start closure + * reads this at spec-build time and injects it into the compiled StepSpec's + * `verification:` field — journal-honest via slice P's named-gate lowering. + * Undefined means no postfix gate (default behavior). Predicate `.gate(fn)` + * still throws `unsupported_gate` and never sets this field. + */ + namedGate: NamedGate | undefined = undefined; private state: OperationState = 'created'; private thenInvoked = false; private rootFailureRecorded = false; @@ -40,11 +59,21 @@ export class AuthoredFlowOperation { observeRejection(this.promise, (error) => this.recordRootFailure(error)); const operation = this; - this.step = Object.freeze({ - gate(): never { + const step: Step = { + gate(configOrPredicate: NamedGate | ((value: T) => boolean), _because?: string): Step { + 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 gates are not lowered by the initial authored executor', + 'postfix .gate(predicate) closures cannot be journaled; ' + + 'use .gate({type: "…", …}) with a slice-P named gate instead.', ); }, then( @@ -64,7 +93,11 @@ export class AuthoredFlowOperation { (error) => operation.recordCallbackFailure(error), ); }, - }); + }; + // Freeze the `then`/`gate` method table so the surface's Step contract + // stays intact, but keep `namedGate` reachable via the enclosing + // operation reference (not through the surface Step handle). + this.step = Object.freeze(step); this.scope.registerStep(this.step, this); } diff --git a/packages/sdk/src/authored-worker-step.ts b/packages/sdk/src/authored-worker-step.ts index 897f4ba95..3352e8274 100644 --- a/packages/sdk/src/authored-worker-step.ts +++ b/packages/sdk/src/authored-worker-step.ts @@ -1,6 +1,6 @@ import type { AuthoredBudget } from './authored-budget.js'; import { parseBudget } from './budget.js'; -import type { AgentOptions, AgentResult, LlmOptions } from '@relayflows/surface'; +import type { AgentOptions, AgentResult, LlmOptions, NamedGate } from '@relayflows/surface'; import { compileSpec, toKernelSpec } from './compile.js'; import { checkAuthoredFlow } from './cli/check.js'; import { classifyOutcome, type RunLifecycleOptions } from './cli/run.js'; @@ -83,7 +83,7 @@ export function authoredWorkerRunner( } return { - async agent(id: string, options: AgentOptions): Promise { + async agent(id: string, options: AgentOptions, verification?: NamedGate): Promise { if (options.workspace !== undefined && localAgentStream !== undefined) { throw new AuthoredFlowExecutionError('unsupported_workspace_permission', 'The local agent worker accepts stream-only steps. Remove workspace or attach a worker that holds its revision pins.'); @@ -125,6 +125,7 @@ export function authoredWorkerRunner( ...(options.cli === undefined ? {} : { cli: options.cli }), ...(options.model === undefined ? {} : { model: options.model }), ...(options.cwd === undefined ? {} : { cwd: options.cwd }), + ...(verification === undefined ? {} : { verification }), }); if (typeof output !== 'object' || output === null || Array.isArray(output)) { throw new AuthoredFlowExecutionError('journal_protocol_violation', `step "${id}" produced a non-object output`); @@ -132,7 +133,7 @@ export function authoredWorkerRunner( const stdout = 'stdout_tail' in output ? output.stdout_tail : undefined; return { summary: typeof stdout === 'string' ? stdout : JSON.stringify(output), artifacts: [] }; }, - async llm(id: string, prompt: string, options?: LlmOptions): Promise { + async llm(id: string, prompt: string, options?: LlmOptions, verification?: NamedGate): Promise { if (options !== undefined && (typeof options !== 'object' || options === null || options.output === undefined)) { throw new AuthoredFlowExecutionError('llm_cli_unresolved', 'f.llm(prompt, options) requires an output JSON Schema.'); } @@ -141,7 +142,10 @@ export function authoredWorkerRunner( || Object.keys(snapshot).some(key => !['output', 'cli', 'model'].includes(key)))) { throw new AuthoredFlowExecutionError('llm_cli_unresolved', 'f.llm options accepts only output, cli, and model.'); } - const output = await run({ id, type: 'llm', prompt, ...snapshot as unknown as LlmOptions }); + const output = await run({ + id, type: 'llm', prompt, ...snapshot as unknown as LlmOptions, + ...(verification === undefined ? {} : { verification }), + }); if (options === undefined && typeof output !== 'string') { throw new AuthoredFlowExecutionError('journal_protocol_violation', `text step "${id}" produced a non-string output`); } @@ -154,11 +158,15 @@ export function authoredWorkerRunner( export function authoredDeterministicRunner( name: string, journal: JournalClient, journalSteps: AuthoredFlowJournalStep[], budget: AuthoredBudget, ) { - return async (id: string, command: string, terminal = false, leaseMs?: number): Promise => { + return async (id: string, command: string, terminal = false, leaseMs?: number, verification?: NamedGate): Promise => { const spec = toKernelSpec(compileSpec({ version: SPEC_SCHEMA_VERSION, name: `${name}/${id}`, - steps: [{ id, type: 'deterministic', command, ...(leaseMs === undefined ? {} : { lease_ms: leaseMs }) }], + steps: [{ + id, type: 'deterministic', command, + ...(leaseMs === undefined ? {} : { lease_ms: leaseMs }), + ...(verification === undefined ? {} : { verification }), + }], })); if (terminal) return readSuccessfulOutput(journal, await journal.runStart(spec), id, journalSteps); return budget.execute(journal, spec, outcome => readSuccessfulOutput(journal, outcome, id, journalSteps)); diff --git a/packages/sdk/tests/authored-flow.test.ts b/packages/sdk/tests/authored-flow.test.ts index 40ce4111a..4e4ea80eb 100644 --- a/packages/sdk/tests/authored-flow.test.ts +++ b/packages/sdk/tests/authored-flow.test.ts @@ -140,10 +140,20 @@ describe('authored flow journal executor', () => { async (f) => f.done('success'), ), disconnectedJournal)).rejects.toMatchObject({ code: 'unsupported_header' }); - await expect(executeAuthoredFlow(flow('gate-not-lowered', async (f) => { + // Predicate .gate(fn) still refuses — JS closures can't be journaled. + await expect(executeAuthoredFlow(flow('gate-predicate-not-lowered', async (f) => { await f.run('true').gate(Boolean); f.done('success'); }), disconnectedJournal)).rejects.toMatchObject({ code: 'unsupported_gate' }); + + // Config-object .gate({...}) MUST NOT throw unsupported_gate — it lowers + // into the compiled step's `verification:` field via slice-P named gates. + // The disconnected journal still refuses the run, but on a different code + // path than 'unsupported_gate' (the step spec reaches journal.runStart). + await expect(executeAuthoredFlow(flow('gate-config-is-lowered', async (f) => { + await f.run('echo ok').gate({ type: 'regex_match', pattern: 'ok' }); + f.done('success'); + }), disconnectedJournal)).rejects.not.toMatchObject({ code: 'unsupported_gate' }); }); it('refuses a workspace permission annotation f.agent cannot enforce, before contacting the journal', async () => { diff --git a/packages/surface/src/index.ts b/packages/surface/src/index.ts index e36dfa65f..619664e22 100644 --- a/packages/surface/src/index.ts +++ b/packages/surface/src/index.ts @@ -15,7 +15,14 @@ export { type RunCompletionReason, type FlowCompletionReason, } from "./completion.js"; -export type { Step } from "./step.js"; +export type { + Step, + NamedGate, + ReferencesInputNamedGate, + SubprocessNamedGate, + WordCountBoundsNamedGate, + RegexMatchNamedGate, +} from "./step.js"; export { flow, type FlowHandle, diff --git a/packages/surface/src/step.ts b/packages/surface/src/step.ts index 715ea8c2b..e5bd4def4 100644 --- a/packages/surface/src/step.ts +++ b/packages/surface/src/step.ts @@ -1,5 +1,56 @@ /** A journal-backed step result with its postfix verification gate. */ export interface Step extends PromiseLike { - /** Fail the step with `verification_failed` when the predicate is false. */ + /** + * Attach a named data gate to this step's `verification:` field. The gate is + * lowered by the SDK into a slice-P journal-honest check that survives + * replay — its shape is data, not code, so covenant 1 (journal-as-truth) + * holds. See `packages/sdk/src/named-gates.ts` for the enforcement. + */ + 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. + */ gate(predicate: (value: T) => boolean, because?: string): Step; } + +/** + * Slice-P named data gates the authored surface can attach postfix via + * `.gate(config)`. The SDK's named-gate lowering enforces these exactly (see + * `packages/sdk/src/named-gates.ts`); the shapes here are the surface-visible + * subset kept in sync with `NamedDataGate` in the SDK spec. + */ +export type NamedGate = + | ReferencesInputNamedGate + | SubprocessNamedGate + | WordCountBoundsNamedGate + | RegexMatchNamedGate; + +export interface ReferencesInputNamedGate { + type: 'references_input'; + input_key: string; + in_output_at?: Array; +} + +export interface SubprocessNamedGate { + type: 'subprocess_gate'; + command: string; + from_output?: Array; +} + +export interface WordCountBoundsNamedGate { + type: 'word_count_bounds'; + min?: number; + max?: number; +} + +export interface RegexMatchNamedGate { + type: 'regex_match'; + pattern: string; + /** Only i, m, and s; evaluated by a non-backtracking RE2 engine. */ + flags?: string; + in_output_at?: Array; +}