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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
30 changes: 20 additions & 10 deletions packages/sdk/src/authored-flow-executor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -181,22 +182,27 @@ export async function executeAuthoredFlow<Input = undefined>(
function llmOperation(prompt: string | TemplateStringsArray, ...values: unknown[]): Step<unknown> {
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<unknown>;
llmOp = new AuthoredFlowOperation(
id, 'llm',
() => assertOperationAllowed('llm', definition.name, requestedCompletion),
() => 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.');
}
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();
Expand Down Expand Up @@ -254,26 +260,30 @@ export async function executeAuthoredFlow<Input = undefined>(
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<string>;
runOp = new AuthoredFlowOperation<string>(
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<AgentResult>;
agentOp = new AuthoredFlowOperation<AgentResult>(
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);
Expand Down
43 changes: 38 additions & 5 deletions packages/sdk/src/authored-flow-operation.ts
Original file line number Diff line number Diff line change
@@ -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<T> {
readonly step: Step<T>;
/**
* 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;
Expand Down Expand Up @@ -40,11 +59,21 @@ export class AuthoredFlowOperation<T> {

observeRejection(this.promise, (error) => this.recordRootFailure(error));
const operation = this;
this.step = Object.freeze({
gate(): never {
const step: Step<T> = {
gate(configOrPredicate: NamedGate | ((value: T) => boolean), _because?: string): Step<T> {
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;

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Config gates skip helper steps

Medium Severity

.gate(config) now accepts a named gate on every Step, but only f.run, f.llm, and f.agent read namedGate. Slack, helper, MCP, and plugin starts ignore it, so the call succeeds and the step runs with no verification.

Additional Locations (2)
Fix in Cursor Fix in Web

Reviewed by Cursor Bugbot for commit 0bcd7f5. Configure here.

}
// 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<TResult1 = T, TResult2 = never>(
Expand All @@ -64,7 +93,11 @@ export class AuthoredFlowOperation<T> {
(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);
}

Expand Down
20 changes: 14 additions & 6 deletions packages/sdk/src/authored-worker-step.ts
Original file line number Diff line number Diff line change
@@ -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';
Expand Down Expand Up @@ -83,7 +83,7 @@ export function authoredWorkerRunner(
}

return {
async agent(id: string, options: AgentOptions): Promise<AgentResult> {
async agent(id: string, options: AgentOptions, verification?: NamedGate): Promise<AgentResult> {
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.');
Expand Down Expand Up @@ -125,14 +125,15 @@ 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`);
}
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<unknown> {
async llm(id: string, prompt: string, options?: LlmOptions, verification?: NamedGate): Promise<unknown> {
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.');
}
Expand All @@ -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`);
}
Expand All @@ -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<string> => {
return async (id: string, command: string, terminal = false, leaseMs?: number, verification?: NamedGate): Promise<string> => {
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 }),
}],

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Run gates ignore verification failure

High Severity

f.run().gate(config) lowers a second .gate step into the kernel spec, but the deterministic runner still only reads the producer step's step.completed. A failed named gate therefore never fails the authored f.run step, so postfix verification on f.run does not enforce.

Additional Locations (1)
Fix in Cursor Fix in Web

Reviewed by Cursor Bugbot for commit 0bcd7f5. Configure here.

}));
if (terminal) return readSuccessfulOutput(journal, await journal.runStart(spec), id, journalSteps);
return budget.execute(journal, spec, outcome => readSuccessfulOutput(journal, outcome, id, journalSteps));
Expand Down
12 changes: 11 additions & 1 deletion packages/sdk/tests/authored-flow.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 () => {
Expand Down
9 changes: 8 additions & 1 deletion packages/surface/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
53 changes: 52 additions & 1 deletion packages/surface/src/step.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,56 @@
/** A journal-backed step result with its postfix verification gate. */
export interface Step<T> extends PromiseLike<T> {
/** 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<T>;
/**
* 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<T>;
}

/**
* 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<string | number>;
}

export interface SubprocessNamedGate {
type: 'subprocess_gate';
command: string;
from_output?: Array<string | number>;
}

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<string | number>;
}
Loading