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
63 changes: 60 additions & 3 deletions packages/sdk/src/authored-flow-loader.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
import { accessSync, constants } from 'node:fs';
import { accessSync, constants, realpathSync } from 'node:fs';
import { createRequire } from 'node:module';
import { resolve } from 'node:path';
import { dirname, resolve } from 'node:path';
import { pathToFileURL } from 'node:url';
import {
getAuthoredFlowDefinition,
Expand All @@ -9,7 +9,7 @@ import {
} from './authored-flow.js';

export class AuthoredFlowLoadError extends Error {
constructor(message: string) {
constructor(message: string, readonly kind: 'invalid_spec' | 'use_not_found' | 'use_invalid' | 'use_cycle' = 'invalid_spec') {
super(message);
this.name = 'AuthoredFlowLoadError';
}
Expand All @@ -28,10 +28,67 @@ export interface LoadedAuthoredFlow {
* `getAuthoredFlowDefinition` import.
*/
readonly getDefinition: GetFlowDefinition;
/** Dependency-first load order, each canonical absolute path appearing once. */
readonly graph: readonly LoadedAuthoredFlowNode[];
}

export interface LoadedAuthoredFlowNode {
readonly path: string;
readonly handle: FlowHandle;
readonly getDefinition: GetFlowDefinition;
readonly use: readonly string[];
}

/** Import and validate a direct-run module without executing its authored body. */
export async function loadAuthoredFlow(path: string): Promise<LoadedAuthoredFlow> {
const loaded = new Map<string, LoadedAuthoredFlowNode>();
const visiting = new Set<string>();
async function visit(sourcePath: string, isRoot = false): Promise<LoadedAuthoredFlowNode> {
let absolutePath: string;
try {
absolutePath = realpathSync(sourcePath);
accessSync(absolutePath, constants.R_OK);
} catch {
throw new AuthoredFlowLoadError(`Flow "${sourcePath}" is not readable.`, isRoot ? 'invalid_spec' : 'use_not_found');
}
if (visiting.has(absolutePath)) {
throw new AuthoredFlowLoadError(`Flow use cycle: ${[...visiting, absolutePath].join(' -> ')}`, 'use_cycle');
}
const cached = loaded.get(absolutePath);
if (cached !== undefined) return cached;
visiting.add(absolutePath);
try {
const { handle, getDefinition } = await importAuthoredFlow(absolutePath);
const dependencies: string[] = [];
for (const entry of getDefinition(handle).header.use ?? []) {
// Validate again at the SDK boundary: the author's surface package may
// be a different version from the runtime loading the graph.
if (!/^(?:\.\/|\.\.\/).+\.flow\.ts$/.test(entry) || /[?#\\\\]/.test(entry)) {
throw new AuthoredFlowLoadError(`Flow "${absolutePath}" use must contain relative .flow.ts paths.`, 'use_invalid');
}
const child = await visit(resolve(dirname(absolutePath), entry));
if (dependencies.includes(child.path)) {
throw new AuthoredFlowLoadError(`Flow "${absolutePath}" declares the same use path more than once.`, 'use_invalid');
}
dependencies.push(child.path);
}
const node = Object.freeze({ path: absolutePath, handle, getDefinition, use: Object.freeze(dependencies) });
loaded.set(absolutePath, node);
return node;
} catch (error) {
if (error instanceof AuthoredFlowLoadError && error.kind === 'invalid_spec' && !isRoot) {
throw new AuthoredFlowLoadError(error.message, 'use_invalid');
}
throw error;
} finally {
visiting.delete(absolutePath);
}
}
const root = await visit(resolve(path), true);
return Object.freeze({ handle: root.handle, getDefinition: root.getDefinition, graph: Object.freeze([...loaded.values()]) });
}

async function importAuthoredFlow(path: string): Promise<Pick<LoadedAuthoredFlow, 'handle' | 'getDefinition'>> {
const absolutePath = resolve(path);
try {
accessSync(absolutePath, constants.R_OK);
Expand Down
59 changes: 59 additions & 0 deletions packages/sdk/tests/authored-use-loader.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,59 @@
import { mkdtempSync, mkdirSync, 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 { loadAuthoredFlow } from '../src/authored-flow-loader.js';

const directories: string[] = [];
afterEach(() => {
for (const directory of directories.splice(0)) rmSync(directory, { recursive: true, force: true });
});

function fixture(files: Record<string, string[]>): string {
const directory = mkdtempSync(join(tmpdir(), 'flows-use-'));
directories.push(directory);
mkdirSync(join(directory, 'node_modules', '@relayflows'), { recursive: true });
symlinkSync(resolve('node_modules/@relayflows/surface'), join(directory, 'node_modules/@relayflows/surface'));
writeFileSync(join(directory, 'package.json'), '{"type":"module"}');
for (const [name, use] of Object.entries(files)) {
writeFileSync(join(directory, `${name}.flow.ts`), `
import { flow } from '@relayflows/surface';
export default flow(${JSON.stringify(name)}, { use: ${JSON.stringify(use)} }, async f => {
throw new Error('loader must never execute authored bodies');
});
`);
}
return directory;
}

describe('authored use graph loader', () => {
it('loads a diamond in dependency order with one node per canonical path', async () => {
const directory = fixture({
root: ['./mid1.flow.ts', './mid2.flow.ts'],
mid1: ['./leaf.flow.ts'], mid2: ['./leaf.flow.ts'], leaf: [],
});
const loaded = await loadAuthoredFlow(join(directory, 'root.flow.ts'));
expect(loaded.graph.map(node => node.handle.name)).toEqual(['leaf', 'mid1', 'mid2', 'root']);
expect(loaded.graph[1]?.use).toEqual(loaded.graph[2]?.use);
expect(Object.isFrozen(loaded.graph)).toBe(true);
});

it('refuses a transitive use cycle before executing any body', async () => {
const directory = fixture({ a: ['./b.flow.ts'], b: ['./a.flow.ts'] });
await expect(loadAuthoredFlow(join(directory, 'a.flow.ts'))).rejects.toMatchObject({ kind: 'use_cycle' });
});

it('refuses a missing declared flow', async () => {
const directory = fixture({ root: ['./missing.flow.ts'] });
await expect(loadAuthoredFlow(join(directory, 'root.flow.ts'))).rejects.toMatchObject({ kind: 'use_not_found' });
});

it.each([
'throw new Error("import failed");',
'export default { name: "forged" };',
])('refuses a declared module that cannot supply a flow: %s', async source => {
const directory = fixture({ root: ['./invalid.flow.ts'] });
writeFileSync(join(directory, 'invalid.flow.ts'), source);
await expect(loadAuthoredFlow(join(directory, 'root.flow.ts'))).rejects.toMatchObject({ kind: 'use_invalid' });
});
});
22 changes: 20 additions & 2 deletions packages/surface/src/flow.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,8 @@ import type { Ctx } from "./context.js";

/** Optional escalation header; the empty header is the common case. */
export interface FlowHeader {
/** Relative paths to reusable authored flows composed by this body. */
use?: string[];
identity?: string;
memory?: { script?: boolean; agent?: boolean };
budget?: string;
Expand All @@ -12,6 +14,7 @@ export interface FlowHeader {
export type FlowBody<Input = unknown> = (f: Ctx, input: Input) => Promise<void>;

export interface ReadonlyFlowHeader {
readonly use?: readonly string[];
readonly identity?: string;
readonly memory?: Readonly<{ script?: boolean; agent?: boolean }>;
readonly budget?: string;
Expand Down Expand Up @@ -114,6 +117,7 @@ function isStoredDefinition(

function freezeHeader(header: FlowHeader): ReadonlyFlowHeader {
const unknownFields = Object.keys(header).filter((field) => ![
"use",
"identity",
"memory",
"budget",
Expand All @@ -137,6 +141,7 @@ function freezeHeader(header: FlowHeader): ReadonlyFlowHeader {
: { mcp: Object.freeze([...header.tools.mcp]) }),
});
return Object.freeze({
...(header.use === undefined ? {} : { use: Object.freeze([...header.use]) }),
...(header.identity === undefined ? {} : { identity: header.identity }),
...(memory === undefined ? {} : { memory }),
...(header.budget === undefined ? {} : { budget: header.budget }),
Expand All @@ -150,12 +155,20 @@ function assertFlowHeader(value: unknown, flowName: string): asserts value is Fl
assertHeaderObject(value, at);
assertKnownKeys(
value,
["identity", "memory", "budget", "tools", "workspace"],
["use", "identity", "memory", "budget", "tools", "workspace"],
at,
);
assertOptionalString(value, "identity", at);
assertOptionalString(value, "budget", at);
assertOptionalString(value, "workspace", at);
assertOptionalStringArray(value, "use", at);
if (value.use !== undefined) {
for (const path of value.use as string[]) {
if (!/^(?:\.\/|\.\.\/).+\.flow\.ts$/.test(path) || /[?#\\\\]/.test(path)) {
throw new TypeError(`${at}.use: expected relative .flow.ts paths`);
}
}
}

if (value.memory !== undefined) {
assertHeaderObject(value.memory, `${at}.memory`);
Expand Down Expand Up @@ -238,8 +251,13 @@ function assertOptionalStringArray(
const candidate = value[key];
if (
candidate !== undefined
&& (!Array.isArray(candidate) || candidate.some((item) => typeof item !== "string"))
&& (!Array.isArray(candidate) || Array.from(candidate).some((item) => typeof item !== "string"))
) {
throw new TypeError(`${at}.${key}: expected an array of strings`);
}
if (key === "use" && Array.isArray(candidate)
&& (candidate.some((item) => item.trim().length === 0)
|| new Set(candidate).size !== candidate.length)) {
throw new TypeError(`${at}.${key}: expected unique nonempty paths`);
}
}
21 changes: 21 additions & 0 deletions packages/surface/tests/flow.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,27 @@ import {
import { getFlowDefinition } from "@relayflows/surface/runtime";

describe("flow", () => {
it("retains an immutable copy of declared relative flow imports", () => {
const use = ["./reviewer.flow.ts", "../shared/helper.flow.ts"];
const handle = flow("chief", { use }, async () => undefined);
use.push("./later.flow.ts");
expect(getFlowDefinition(handle).header.use).toEqual([
"./reviewer.flow.ts", "../shared/helper.flow.ts",
]);
expect(Object.isFrozen(getFlowDefinition(handle).header.use)).toBe(true);
});

it.each([
null, "./reviewer.flow.ts", [42], [""], [" "],
["./reviewer.flow.ts", "./reviewer.flow.ts"],
["reviewer.flow.ts"], ["@team/reviewer"], ["/reviewer.flow.ts"],
["./reviewer.ts"], ["./reviewer.flow.ts?other.flow.ts"],
Array(1),
])("refuses malformed use declarations: %j", (use) => {
expect(() => flow("chief", { use } as FlowHeader, async () => undefined))
.toThrow("header.use");
});

it("defines a flow with the empty header as the default", () => {
const body = async (_f: Ctx): Promise<void> => undefined;
const definition = flow("release-note", body);
Expand Down
Loading