diff --git a/packages/sdk/src/authored-flow-loader.ts b/packages/sdk/src/authored-flow-loader.ts index dde721aee..23193dcec 100644 --- a/packages/sdk/src/authored-flow-loader.ts +++ b/packages/sdk/src/authored-flow-loader.ts @@ -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, @@ -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'; } @@ -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 { + const loaded = new Map(); + const visiting = new Set(); + async function visit(sourcePath: string, isRoot = false): Promise { + 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> { const absolutePath = resolve(path); try { accessSync(absolutePath, constants.R_OK); diff --git a/packages/sdk/tests/authored-use-loader.test.ts b/packages/sdk/tests/authored-use-loader.test.ts new file mode 100644 index 000000000..4f083417b --- /dev/null +++ b/packages/sdk/tests/authored-use-loader.test.ts @@ -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 { + 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' }); + }); +}); diff --git a/packages/surface/src/flow.ts b/packages/surface/src/flow.ts index fa1f358d8..1323f004d 100644 --- a/packages/surface/src/flow.ts +++ b/packages/surface/src/flow.ts @@ -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; @@ -12,6 +14,7 @@ export interface FlowHeader { export type FlowBody = (f: Ctx, input: Input) => Promise; export interface ReadonlyFlowHeader { + readonly use?: readonly string[]; readonly identity?: string; readonly memory?: Readonly<{ script?: boolean; agent?: boolean }>; readonly budget?: string; @@ -114,6 +117,7 @@ function isStoredDefinition( function freezeHeader(header: FlowHeader): ReadonlyFlowHeader { const unknownFields = Object.keys(header).filter((field) => ![ + "use", "identity", "memory", "budget", @@ -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 }), @@ -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`); @@ -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`); + } } diff --git a/packages/surface/tests/flow.test.ts b/packages/surface/tests/flow.test.ts index 5c74aa81d..3d9af98da 100644 --- a/packages/surface/tests/flow.test.ts +++ b/packages/surface/tests/flow.test.ts @@ -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 => undefined; const definition = flow("release-note", body);