diff --git a/kernel/relayflowd-journal/src/lib.rs b/kernel/relayflowd-journal/src/lib.rs index 5385cedcf..62251e05f 100644 --- a/kernel/relayflowd-journal/src/lib.rs +++ b/kernel/relayflowd-journal/src/lib.rs @@ -79,28 +79,33 @@ impl SqliteJournal { std::fs::create_dir_all(parent)?; } let run_id = run_id.into(); - let connection = Connection::open_with_flags( + let mut connection = Connection::open_with_flags( &path, OpenFlags::SQLITE_OPEN_READ_WRITE | OpenFlags::SQLITE_OPEN_CREATE, )?; configure(&connection)?; - connection.execute_batch(SCHEMA)?; - let mut journal = Self { - connection, - run_id, - path, - }; - let transaction = journal.connection.transaction()?; + // Schema DDL + meta/segment rows land in one transaction so SIGKILL + // between create() and the first append() never leaves the file with + // schema but no meta row. Without this, open()'s + // "SELECT value FROM meta WHERE key = 'run_id'" surfaces as + // "Query returned no rows" on resume (kernel/journal SIGKILL race: + // reproduced by webhook-live's SIGKILL-after-spawn-before-ack test). + let transaction = connection.transaction()?; + transaction.execute_batch(SCHEMA)?; transaction.execute( "INSERT INTO meta(key, value) VALUES ('run_id', ?1), ('created_at_ms', ?2), ('journal_version', ?3)", - params![journal.run_id, created_at_ms.to_string(), relayflowd_core::JOURNAL_VERSION.to_string()], + params![run_id, created_at_ms.to_string(), relayflowd_core::JOURNAL_VERSION.to_string()], )?; transaction.execute( "INSERT INTO segments(segment_id, journal_version, opened_seq) VALUES (1, ?1, 1)", [i64::from(relayflowd_core::JOURNAL_VERSION)], )?; transaction.commit()?; - Ok(journal) + Ok(Self { + connection, + run_id, + path, + }) } pub fn open(path: impl AsRef) -> Result { diff --git a/kernel/relayflowd/src/engine/wake.rs b/kernel/relayflowd/src/engine/wake.rs index a3ab8b325..4668adf09 100644 --- a/kernel/relayflowd/src/engine/wake.rs +++ b/kernel/relayflowd/src/engine/wake.rs @@ -11,6 +11,14 @@ use ulid::Ulid; use super::{DriveOptions, Engine, RunOutcome, canonical_hash}; +/// Additive event-run payload: normal run.spawned readers retain compatibility. +#[derive(Serialize)] +struct EventRunSpawnedPayload { + #[serde(flatten)] + run: RunSpawnedPayload, + event: serde_json::Value, +} + /// Default silence budget when a trigger does not declare its own /// `stale_after_ms`. 5 minutes is long enough not to trip a sluggish /// external stream during a normal quiet stretch, short enough that a @@ -32,8 +40,31 @@ impl Engine { spec: RunSpec, event: Event, created_by: &str, + ) -> Result { + self.submit_event_inner(spec, event, created_by, None) + } + + /// Inbox retries resume the claimed run before acknowledging its file. + pub fn submit_webhook_event( + &self, + spec: RunSpec, + event: Event, + resume: &dyn Fn(&str) -> Result, + ) -> Result { + self.submit_event_inner(spec, event, "webhook", Some(resume)) + } + + fn submit_event_inner( + &self, + spec: RunSpec, + event: Event, + created_by: &str, + inbox_resume: Option<&dyn Fn(&str) -> Result>, ) -> Result { spec.validate().context("invalid run spec")?; + if inbox_resume.is_some() { + self.preflight_placement(&spec)?; + } let Some(trigger) = spec .triggers .iter() @@ -107,16 +138,20 @@ impl Engine { stale_after_ms, self.clock.now_ms(), )?; - if self - .registry()? - .claim_event(&flow_key, &trigger.id, &event_key, &run_id, self.boot_id())? - .is_some() - { + if let Some(existing_run) = self.registry()?.claim_event( + &flow_key, + &trigger.id, + &event_key, + &run_id, + self.boot_id(), + )? { return Ok(EventSubmitOutcome { matched: true, deduped: true, subscription_id: Some(trigger.id), - run: None, + run: inbox_resume + .map(|resume| resume(&existing_run)) + .transpose()?, }); } // The claim is now held by this boot, and `claim_event` will tell any @@ -216,12 +251,15 @@ impl Engine { None, None, now_ms, - RunSpawnedPayload { - spec: spec_value.clone(), - spec_hash: canonical_hash(&spec_value), - parent_run_id: None, - journal_version: relayflowd_core::JOURNAL_VERSION, - created_by: created_by.to_owned(), + EventRunSpawnedPayload { + run: RunSpawnedPayload { + spec: spec_value.clone(), + spec_hash: canonical_hash(&spec_value), + parent_run_id: None, + journal_version: relayflowd_core::JOURNAL_VERSION, + created_by: created_by.to_owned(), + }, + event: event.payload.clone(), }, ), )?; @@ -292,16 +330,15 @@ impl Drop for ClaimGuard { // during an unwind, when holding a borrow of the engine would constrain // the guard's lifetime to it for no benefit. `Registry::open` is what // every other caller here does per operation anyway. - let released = Registry::open(self.data_dir.join("relayflowd.sqlite3")).and_then( - |registry| { + let released = + Registry::open(self.data_dir.join("relayflowd.sqlite3")).and_then(|registry| { registry.release_claim( &self.flow_key, &self.subscription_id, &self.event_key, &self.run_id, ) - }, - ); + }); if let Err(error) = released { // Report, do not panic. Panicking in `Drop` during an unwind aborts // the process, which would turn a stranded event into a dead daemon. @@ -344,7 +381,9 @@ mod claim_guard_tests { /// Take a claim the way `submit_event` does, without registering a run. fn claim(data_dir: &std::path::Path, run_id: &str) { assert_eq!( - registry(data_dir).claim_event(FLOW, SUB, KEY, run_id, BOOT).unwrap(), + registry(data_dir) + .claim_event(FLOW, SUB, KEY, run_id, BOOT) + .unwrap(), None, "the first claim must be granted" ); diff --git a/kernel/relayflowd/src/lib.rs b/kernel/relayflowd/src/lib.rs index 85e3fe674..004cdcb63 100644 --- a/kernel/relayflowd/src/lib.rs +++ b/kernel/relayflowd/src/lib.rs @@ -4,6 +4,7 @@ pub mod exec_det; pub mod memory; pub mod server; pub mod socket_path; +pub mod trigger_watcher; pub mod worker; pub use engine::{ diff --git a/kernel/relayflowd/src/server.rs b/kernel/relayflowd/src/server.rs index 8f1f5865e..7bcbbb784 100644 --- a/kernel/relayflowd/src/server.rs +++ b/kernel/relayflowd/src/server.rs @@ -74,6 +74,16 @@ pub fn serve(data_dir: &Path) -> Result<()> { // because the failure mode it catches is "minutes without an event", // not "seconds without a heartbeat". liveness::spawn_liveness_sweep(data_dir.to_path_buf()); + let trigger_hub = hub.clone(); + let trigger_dir = data_dir.to_path_buf(); + crate::trigger_watcher::spawn_watcher(trigger_dir.clone(), move |spec, event| { + let engine = Engine::with_runtime(&trigger_dir, trigger_hub.clone(), trigger_hub.clone()); + engine.submit_webhook_event(spec, event, &|run_id| { + let lock = trigger_hub.run_lock(run_id); + let _guard = lock.lock().expect("run lock"); + engine.resume_live(run_id, trigger_hub.as_ref()) + }) + }); let next_connection = Arc::new(AtomicU64::new(1)); for connection in listener.incoming() { let connection = connection?; diff --git a/kernel/relayflowd/src/trigger_watcher.rs b/kernel/relayflowd/src/trigger_watcher.rs new file mode 100644 index 000000000..bdf3a43ea --- /dev/null +++ b/kernel/relayflowd/src/trigger_watcher.rs @@ -0,0 +1,137 @@ +//! Local inbox ingress. Each `triggers/.json` is a compiled RunSpec whose +//! event_type and executor equal . Provision these files before ingress; +//! do not change a binding while its inbox is pending (dedupe is spec-scoped). +//! TODO https://github.com/AgentWorkforce/flows/issues/301: bind sealed bundles +//! through flows deploy; TS handler deployment and Cloud mounts are separate. +//! No author process stays alive between events. The daemon polls at 1 Hz. + +use crate::engine::{EventSubmitOutcome, read_spec}; +use anyhow::{Context, Result, bail}; +use relayflowd_core::{Event, RunSpec}; +use std::{ + fs, + path::{Path, PathBuf}, + thread, + time::Duration, +}; + +const MAX_EVENT_BYTES: u64 = 1024 * 1024; + +pub fn valid_name(name: &str) -> bool { + !name.is_empty() + && name.len() <= 128 + && name.as_bytes()[0].is_ascii_alphanumeric() + && name + .bytes() + .all(|c| c.is_ascii_alphanumeric() || c == b'_' || c == b'-') +} + +/// The callback is the engine boundary; filesystem traversal never executes code. +/// Errors retain the offending file and do not starve other inboxes. +pub fn poll_once( + data_dir: &Path, + submit: &mut impl FnMut(RunSpec, Event) -> Result, +) -> Result> { + let inbox = data_dir.join("inbox"); + fs::create_dir_all(&inbox)?; + let mut errors = Vec::new(); + for directory in fs::read_dir(&inbox)? { + let directory = directory?; + if !directory.file_type()?.is_dir() { + continue; + } + let Some(name) = directory.file_name().to_str().map(str::to_owned) else { + continue; + }; + if !valid_name(&name) { + continue; + } + for file in fs::read_dir(directory.path())? { + let file = file?; + if !file.file_type()?.is_file() { + continue; + } + let filename = file.file_name(); + let Some(filename) = filename.to_str() else { + continue; + }; + let Some(id) = filename.strip_suffix(".json") else { + continue; + }; + if !valid_name(id) { + continue; + } + if let Err(error) = process_file(data_dir, &name, filename, &file.path(), submit) { + errors.push(format!("{}: {error:#}", file.path().display())); + } + } + } + Ok(errors) +} + +fn process_file( + data_dir: &Path, + name: &str, + filename: &str, + path: &Path, + submit: &mut impl FnMut(RunSpec, Event) -> Result, +) -> Result<()> { + if fs::metadata(path)?.len() > MAX_EVENT_BYTES { + bail!("event exceeds 1 MiB"); + } + let event = Event { + event_type: name.to_owned(), + payload: serde_json::from_slice(&fs::read(path)?).context("parse inbox event")?, + key: Some(filename.to_owned()), + }; + let binding = data_dir.join("triggers").join(format!("{name}.json")); + if !fs::symlink_metadata(&binding)?.is_file() { + bail!("trigger binding must be a regular file"); + } + let spec = read_spec(&binding)?; + spec.validate().context("invalid trigger spec")?; + if spec.triggers.is_empty() + || spec + .triggers + .iter() + .any(|trigger| trigger.executor != name || trigger.event_type.as_deref() != Some(name)) + { + bail!("trigger binding must declare executor and event_type {name:?}"); + } + let outcome = submit(spec, event)?; + // A nonmatch is a consumed filter rejection. A matched submission must be + // durable and driven (or resumed) before moving: crash-before-move retries. + if outcome.matched && outcome.run.is_none() { + bail!("matched event has no durable run receipt"); + } + let processed = data_dir.join("inbox-processed").join(name); + fs::create_dir_all(&processed)?; + if !fs::symlink_metadata(&processed)?.is_dir() { + bail!("processed inbox must be a directory"); + } + fs::rename(path, processed.join(filename)).context("archive inbox event")?; + fs::File::open(&processed)?.sync_all()?; + fs::File::open(path.parent().context("inbox parent")?)?.sync_all()?; + Ok(()) +} + +pub(crate) fn spawn_watcher( + data_dir: PathBuf, + mut submit: impl FnMut(RunSpec, Event) -> Result + Send + 'static, +) { + thread::spawn(move || { + loop { + match poll_once(&data_dir, &mut submit) { + Ok(errors) => { + for error in errors { + eprintln!("relayflowd: inbox retained: {error}"); + } + } + Err(error) => eprintln!("relayflowd: inbox poll failed: {error:#}"), + } + // TODO https://github.com/AgentWorkforce/flows/issues/301: native watch + // notifications are a follow-up; this slice deliberately polls at 1 Hz. + thread::sleep(Duration::from_secs(1)); + } + }); +} diff --git a/kernel/relayflowd/tests/trigger_watcher.rs b/kernel/relayflowd/tests/trigger_watcher.rs new file mode 100644 index 000000000..5d17be25e --- /dev/null +++ b/kernel/relayflowd/tests/trigger_watcher.rs @@ -0,0 +1,113 @@ +use relayflowd::{Engine, trigger_watcher::poll_once}; +use relayflowd_journal::SqliteJournal; +use serde_json::{Value, json}; +use std::fs; + +fn provision(path: &std::path::Path) { + fs::create_dir_all(path.join("triggers")).unwrap(); + fs::create_dir_all(path.join("inbox/release")).unwrap(); + fs::write( + path.join("triggers/release.json"), + json!({ + "version": "0.1.0", "name": "release", + "triggers": [{"id":"release", "executor":"release", "event_type":"release", + "pattern":{"action":"released"}, "dedupe_key_template":"{{event.type}}"}], + "steps": [{"id":"log", "type":"deterministic", "command":"printf accepted"}] + }) + .to_string(), + ) + .unwrap(); +} + +#[test] +fn journals_payload_and_filename_key_then_archives_and_dedupes_replay() { + let dir = tempfile::tempdir().unwrap(); + let root = dir.path(); + provision(root); + let payload = json!({"action":"released", "nested":{"ok":true}}); + let file = root.join("inbox/release/event-1.json"); + fs::write(&file, payload.to_string()).unwrap(); + let engine = Engine::new(root); + let mut ids = Vec::new(); + let mut submit = |spec, event| { + let receipt = engine.submit_webhook_event(spec, event, &|id| engine.resume(id, None))?; + ids.push(receipt.run.as_ref().unwrap().run_id.clone()); + Ok(receipt) + }; + assert!(poll_once(root, &mut submit).unwrap().is_empty()); + assert!(!file.exists()); + fs::rename(root.join("inbox-processed/release/event-1.json"), &file).unwrap(); + assert!(poll_once(root, &mut submit).unwrap().is_empty()); + assert_eq!(ids.len(), 2); + assert_eq!(ids[0], ids[1]); + let entries = SqliteJournal::open(root.join(format!("runs/{}.sqlite3", ids[0]))) + .unwrap() + .scan_all() + .unwrap(); + assert_eq!( + entries + .iter() + .filter(|e| e.entry_type.as_str() == "run.spawned") + .count(), + 1 + ); + assert_eq!(entries[0].payload["event"], payload); + let received = entries + .iter() + .find(|e| e.entry_type.as_str() == "event.received") + .unwrap(); + assert_eq!(received.payload["event_key"], "event-1.json"); +} + +#[test] +fn retains_bad_and_unregistered_events_while_consuming_filter_nonmatches() { + let dir = tempfile::tempdir().unwrap(); + let root = dir.path(); + provision(root); + fs::write(root.join("inbox/release/bad.json"), "{").unwrap(); + fs::write( + root.join("inbox/release/ignored.json"), + "{\"action\":\"ignored\"}", + ) + .unwrap(); + fs::write(root.join("inbox/release/.pending.tmp"), "{").unwrap(); + fs::create_dir_all(root.join("inbox/unregistered")).unwrap(); + fs::write(root.join("inbox/unregistered/event.json"), "{}").unwrap(); + let engine = Engine::new(root); + let errors = poll_once(root, &mut |spec, event| { + engine.submit_webhook_event(spec, event, &|id| engine.resume(id, None)) + }) + .unwrap(); + assert_eq!(errors.len(), 2); + assert!(root.join("inbox/release/bad.json").exists()); + assert!(root.join("inbox/release/.pending.tmp").exists()); + assert!(root.join("inbox/unregistered/event.json").exists()); + assert!(root.join("inbox-processed/release/ignored.json").exists()); + assert!(!root.join("runs").exists()); +} + +#[test] +fn failed_archive_retries_the_same_durable_run() { + let dir = tempfile::tempdir().unwrap(); + let root = dir.path(); + provision(root); + fs::write( + root.join("inbox/release/event.json"), + "{\"action\":\"released\"}", + ) + .unwrap(); + // Force a failure after spawn, precisely at the acknowledgement boundary. + fs::write(root.join("inbox-processed"), "blocked").unwrap(); + let engine = Engine::new(root); + let mut receipts = Vec::::new(); + let mut submit = |spec, event| { + let result = engine.submit_webhook_event(spec, event, &|id| engine.resume(id, None))?; + receipts.push(serde_json::to_value(&result)?); + Ok(result) + }; + assert_eq!(poll_once(root, &mut submit).unwrap().len(), 1); + fs::remove_file(root.join("inbox-processed")).unwrap(); + assert!(poll_once(root, &mut submit).unwrap().is_empty()); + assert_eq!(receipts[0]["run"]["run_id"], receipts[1]["run"]["run_id"]); + assert_eq!(receipts[1]["deduped"], true); +} diff --git a/packages/sdk/src/cli.ts b/packages/sdk/src/cli.ts index 1e0dd6a33..a20bd2947 100644 --- a/packages/sdk/src/cli.ts +++ b/packages/sdk/src/cli.ts @@ -17,6 +17,8 @@ import { type RunProgress, type RunReport, } from './cli/run.js'; +import { checkAuthoredTriggers } from './cli/check-triggers.js'; +import { parseWebhookArgs, runServeWebhook } from './cli/serve-webhook.js'; import { runDirectFlow } from './cli/direct-run.js'; import { parseReplayArgs, replayJournal, type ReplayArgs } from './cli/replay.js'; import { checkTypeScriptFlow } from './cli/check-typescript.js'; @@ -42,6 +44,7 @@ type CliExitCode = 0 | 1 | 2 | 3; type ParsedArgs = | ReplayArgs | BuildArgs + | { command: 'serve-webhook'; dataDir: string; port: number } | { command: 'cloud-run'; value: string; json: boolean; wait: boolean } | { command: 'check'; json: boolean; watch: boolean; value: string } | { command: 'run'; reuseFromRunId: string | undefined; localAgent: boolean; dataDir: string; input: string | undefined; json: boolean; spawn: boolean; noObserverLink: boolean; value: string } @@ -58,6 +61,7 @@ const USAGE = [ 'flows build [--out ] ', 'flows build --verify ', 'flows check [--watch] [--json] ', + 'flows serve-webhook --data-dir --port

', 'flows run [--json] [--no-spawn] [--no-observer-link] [--data-dir

] [--local-agent] [--reuse-from ] ', 'flows run --cloud [--json] [--wait] ', 'flows run [--json] [--no-spawn] [--no-observer-link] [--data-dir ] [--local-agent] --input ', @@ -99,6 +103,8 @@ export async function runCli( return 2; } + if (parsed.command === 'serve-webhook') return runServeWebhook(parsed, io); + if (parsed.command === 'cloud-run') return runCloudCli(parsed, io); if (parsed.command === 'replay') return replayJournal(parsed, io); if (parsed.command === 'build') return runBuild(parsed, io); @@ -223,11 +229,16 @@ async function checkAuthoredFlowComposed(path: string): Promise<{ report: CheckR const helper = await checkHelperBody(path); if (!helper.report.ok) return helper; const mcp = await checkTypeScriptFlow(path); + const triggers = isAuthoredFlowPath(path) + ? await checkAuthoredTriggers(path) + : undefined; + const triggerDiagnostics = triggers?.report.diagnostics ?? []; + const triggerOk = triggers?.report.ok ?? true; return { report: { ...mcp.report, - diagnostics: [...helper.report.diagnostics, ...mcp.report.diagnostics], - ok: helper.report.ok && mcp.report.ok, + diagnostics: [...helper.report.diagnostics, ...mcp.report.diagnostics, ...triggerDiagnostics], + ok: helper.report.ok && mcp.report.ok && triggerOk, }, }; } @@ -393,6 +404,7 @@ function parseArgs(args: readonly string[]): ParsedArgs | undefined { const command = args[0]; if (command === 'replay') return parseReplayArgs(args.slice(1)); if (command === 'build') return parseBuildArgs(args.slice(1)); + if (command === 'serve-webhook') return parseWebhookArgs(args.slice(1)); if (command === 'hn-monitor') return parseHnMonitorArgs(args.slice(1)); if (command === 'tick') return parseTickArgs(args.slice(1)); if (command === 'observer') return parseObserverArgs(args.slice(1)); diff --git a/packages/sdk/src/cli/check-triggers.ts b/packages/sdk/src/cli/check-triggers.ts new file mode 100644 index 000000000..aa76d5ef1 --- /dev/null +++ b/packages/sdk/src/cli/check-triggers.ts @@ -0,0 +1,41 @@ +import { dirname, resolve } from 'node:path'; +import { loadAuthoredFlow, type LoadedAuthoredFlow } from '../authored-flow-loader.js'; +import { preflightWebhookTriggers } from '../preflight.js'; +import { checkSlackHelpers } from '../slack-preflight.js'; +import { inputFailureReport, readProjectConfig, type CheckReport } from './check.js'; + +/** + * Inspect an authored flow without running a handler or contacting the daemon. + * Combines the two authored-only preflight paths that share a loaded flow: the + * webhook-trigger check (E) and the f.slack helper check (B). Skipping either + * turned this into a silent trapdoor -- a `.flow.ts` using `f.slack.post` + * without a token would pass `flows check` and only crash at run. + */ +export async function checkAuthoredTriggers(path: string): Promise<{ + report: CheckReport; + loaded?: LoadedAuthoredFlow; +}> { + try { + const loaded = await loadAuthoredFlow(path); + const definition = loaded.getDefinition(loaded.handle); + const config = readProjectConfig(dirname(resolve(path))); + const triggerDiagnostics = preflightWebhookTriggers( + (definition.handlers ?? []).map(handler => handler.trigger), config.executors, + ); + const helperReport = checkSlackHelpers(definition); + const diagnostics = [...triggerDiagnostics, ...helperReport.diagnostics]; + return { + loaded, + report: { + ok: diagnostics.length === 0, path, gates: [], resolutions: [], diagnostics, + ...(config.path === undefined ? {} : { projectConfigPath: config.path }), + }, + }; + } catch (error) { + return { report: inputFailureReport({ + kind: typeof error === 'object' && error !== null && 'kind' in error && error.kind === 'config_invalid' + ? 'config_invalid' : 'invalid_spec', + message: error instanceof Error ? error.message : 'Cannot inspect authored triggers', + }, path) }; + } +} diff --git a/packages/sdk/src/cli/check-typescript.ts b/packages/sdk/src/cli/check-typescript.ts index 82fc5c66f..65d30deb8 100644 --- a/packages/sdk/src/cli/check-typescript.ts +++ b/packages/sdk/src/cli/check-typescript.ts @@ -27,8 +27,9 @@ export async function checkMcpHeader( ): Promise { const empty = { servers: Object.freeze({}), inventory: Object.freeze({}) }; const KNOWN_HEADER_FIELDS = new Set(['tools', 'budget', 'identity', 'memory', 'workspace', 'use']); - const unsupported = Object.keys(definition.header).filter(key => !KNOWN_HEADER_FIELDS.has(key)); - if (definition.header.tools?.relayfile !== undefined) unsupported.push('tools.relayfile'); + const header = definition.header ?? {}; + const unsupported = Object.keys(header).filter(key => !KNOWN_HEADER_FIELDS.has(key)); + if (header.tools?.relayfile !== undefined) unsupported.push('tools.relayfile'); if (unsupported.length) { return { ...empty, report: inputFailureReport({ kind: 'invalid_spec', message: `flow "${definition.name}" uses unsupported header fields: ${unsupported.join(', ')}` }, path) }; @@ -37,7 +38,7 @@ export async function checkMcpHeader( const config = readProjectConfig(dirname(resolve(path))); const result = await preflight({ version: SPEC_SCHEMA_VERSION, name: definition.name, steps: [{ id: 'header', type: 'deterministic', command: ':' }] }, { - mcpServers: definition.header.tools?.mcp ?? [], mcp: config.mcp, + mcpServers: header.tools?.mcp ?? [], mcp: config.mcp, probes: { command: () => true, cli: () => { throw new Error('no CLI declared'); }, executor: () => false }, }); return { diff --git a/packages/sdk/src/cli/direct-run.ts b/packages/sdk/src/cli/direct-run.ts index d9bc82652..80fbf7199 100644 --- a/packages/sdk/src/cli/direct-run.ts +++ b/packages/sdk/src/cli/direct-run.ts @@ -6,10 +6,11 @@ import { AuthoredFlowExecutionError, executeAuthoredFlow, } from '../authored-flow-executor.js'; -import { AuthoredFlowLoadError, loadAuthoredFlow } from '../authored-flow-loader.js'; +import { AuthoredFlowLoadError } from '../authored-flow-loader.js'; import { DirectInputError, parseDirectInput } from '../direct-input.js'; import { JournalClient } from '../journal-client.js'; import { inputFailureReport } from './check.js'; +import { checkAuthoredTriggers } from './check-triggers.js'; import { connect, emptyReport, @@ -41,6 +42,15 @@ export async function runDirectFlow( }; } + // Declared triggers are knowable before any daemon or step is started. + // Importing the authored module is unavoidable here — trigger sources + // are only observable after `flow(...).on(webhook(...))` has run — but + // the authored body is not called, so a side-effect-in-body flow still + // has its body deferred until after daemon-attach below. + const checked = await checkAuthoredTriggers(path); + if (!checked.report.ok || checked.loaded === undefined) { + return { exitCode: 2, report: fromCheckReport('run', checked.report) }; + } const socketPath = socketFor(dataDir); const base: RunReport = { ...emptyReport('run'), path }; const client = new JournalClient(socketPath); @@ -52,7 +62,7 @@ export async function runDirectFlow( let llmClient: JournalClient | undefined; let llmFailure: unknown; try { - const { handle, getDefinition } = await loadAuthoredFlow(path); + const { handle, getDefinition } = checked.loaded; if (options.localAgent) { localAgent = await attachLocalAgent(client); // A session owns one worker registration. Keep the workspace-free LLM diff --git a/packages/sdk/src/cli/serve-webhook.ts b/packages/sdk/src/cli/serve-webhook.ts new file mode 100644 index 000000000..1c216de5c --- /dev/null +++ b/packages/sdk/src/cli/serve-webhook.ts @@ -0,0 +1,141 @@ +import { randomUUID } from 'node:crypto'; +import { mkdir, open, rename, unlink, lstat } from 'node:fs/promises'; +import { createServer, type Server, type ServerResponse } from 'node:http'; +import { join, resolve } from 'node:path'; +import { TextDecoder } from 'node:util'; +import type { CliIo } from '../cli.js'; + +const MAX_BODY_BYTES = 1024 * 1024; +const NAME = /^[A-Za-z0-9][A-Za-z0-9_-]{0,127}$/; + +export function parseWebhookArgs(args: readonly string[]): { + command: 'serve-webhook'; dataDir: string; port: number; +} | undefined { + const values = new Map(); + for (let i = 0; i < args.length; i += 2) { + const flag = args[i]!; + const value = args[i + 1]; + if (!['--data-dir', '--port'].includes(flag) || values.has(flag) + || !value || value.startsWith('-')) return undefined; + values.set(flag, value); + } + const dataDir = values.get('--data-dir'); + const portText = values.get('--port'); + if (!dataDir || !portText || !/^\d+$/.test(portText)) return undefined; + const port = Number(portText); + if (!Number.isInteger(port) || port < 0 || port > 65535) return undefined; + return { command: 'serve-webhook', dataDir, port }; +} + +function reply(response: ServerResponse, status: number, body: object): void { + response.writeHead(status, { 'content-type': 'application/json' }); + response.end(JSON.stringify(body)); +} + +/** POST / accepts JSON, including scalar values. No daemon connection. */ +export async function startWebhookServer(dataDir: string, port: number): Promise { + const inbox = join(resolve(dataDir), 'inbox'); + await directory(inbox); + // TODO https://github.com/AgentWorkforce/flows/issues/301: provider signatures, + // public ingress and Cloud mount provisioning belong to the deployment slice. + const server = createServer(async (request, response) => { + if (request.method !== 'POST') { + request.resume(); + response.setHeader('allow', 'POST'); + reply(response, 405, { error: 'method_not_allowed' }); + return; + } + const name = request.url?.slice(1); + if (!name || !NAME.test(name)) { + request.resume(); + reply(response, 404, { error: 'invalid_webhook_name' }); + return; + } + let temporary: string | undefined; + try { + let bytes = 0; + const chunks: Buffer[] = []; + for await (const chunk of request) { + const buffer = Buffer.from(chunk as Uint8Array); + bytes += buffer.length; + if (bytes > MAX_BODY_BYTES) { + reply(response, 413, { error: 'payload_too_large' }); + return; + } + chunks.push(buffer); + } + let payload: unknown; + try { + payload = JSON.parse(new TextDecoder('utf-8', { fatal: true }).decode(Buffer.concat(chunks)), (_key, value: unknown) => { + if (typeof value === 'number' && !Number.isFinite(value)) throw new Error('non-finite JSON number'); + return value; + }); + } catch { + reply(response, 400, { error: 'invalid_json' }); + return; + } + const target = join(inbox, name); + await directory(target); + const id = randomUUID(); + temporary = join(target, `.${id}.tmp`); + const file = await open(temporary, 'wx', 0o600); + try { + await file.writeFile(JSON.stringify(payload)); + await file.sync(); + } finally { + await file.close(); + } + await rename(temporary, join(target, `${id}.json`)); + temporary = undefined; + const dir = await open(target, 'r'); + try { await dir.sync(); } finally { await dir.close(); } + reply(response, 202, { accepted: true, id }); + } catch { + if (temporary) await unlink(temporary).catch(() => undefined); + if (!response.headersSent) reply(response, 500, { error: 'inbox_write_failed' }); + else response.destroy(); + } + }); + server.requestTimeout = 30_000; + server.headersTimeout = 10_000; + await new Promise((accept, reject) => { + server.once('error', reject); + server.listen(port, '127.0.0.1', () => { + server.off('error', reject); + accept(); + }); + }); + return server; +} + +async function directory(path: string): Promise { + await mkdir(path, { recursive: true }); + if (!(await lstat(path)).isDirectory()) throw new Error('inbox directory must not be a symlink'); +} + +export async function runServeWebhook( + options: { dataDir: string; port: number }, io: CliIo, +): Promise<0 | 1> { + try { + const server = await startWebhookServer(options.dataDir, options.port); + const address = server.address(); + io.stdout(`WEBHOOK http://127.0.0.1:${typeof address === 'object' && address ? address.port : options.port}`); + await new Promise((accept, reject) => { + const stop = (): void => { + server.close(error => error ? reject(error) : accept()); + server.closeAllConnections(); + }; + server.once('error', reject); + process.once('SIGINT', stop); + process.once('SIGTERM', stop); + server.once('close', () => { + process.off('SIGINT', stop); + process.off('SIGTERM', stop); + }); + }); + return 0; + } catch (error) { + io.stderr(`FAILED [webhook_server] ${error instanceof Error ? error.message : 'receiver failed'}`); + return 1; + } +} diff --git a/packages/sdk/src/preflight.ts b/packages/sdk/src/preflight.ts index e44f6e4b4..77b9564a3 100644 --- a/packages/sdk/src/preflight.ts +++ b/packages/sdk/src/preflight.ts @@ -2,6 +2,7 @@ import type { FlowSpec, StepSpec, TriggerSpec, McpServerConfig } from './spec.js import { McpError, openMcpSession, type McpDiagnostic } from './mcp-client.js'; import { BudgetSyntaxError } from './budget.js'; import { budgetDiagnostics } from './budget-preflight.js'; +import type { TriggerSource } from '@relayflows/surface'; import { acceptsAnyOutput, inspectStepGate, type StepGateInspection } from './gate-contract.js'; import { compileSpec, CompileError } from './compile.js'; import type { @@ -109,6 +110,21 @@ export interface PreflightResult { diagnostics: PreflightDiagnostic[]; } +/** Pure declared-surface check, before opening a receiver or journal. */ +export function preflightWebhookTriggers( + triggers: readonly TriggerSource[], + executors: readonly string[], +): PreflightRefusal[] { + return [...new Set(triggers.map(trigger => trigger.name))] + .filter(name => !executors.includes(name)) + .map(name => ({ + severity: 'refusal', + kind: 'no_executor', + executor: name, + message: `webhook trigger "${name}" is not registered in flows.json`, + })); +} + // Preserve the synchronous declarative API; an authored MCP declaration opts // into asynchronous connection probes after the same pure refusal pass. export function preflight(flow: FlowSpec, options: PreflightOptions & { mcpServers: readonly string[] }): Promise; @@ -544,11 +560,12 @@ function firstCommandWord(command: string): string | undefined { /** Helper preflight never evaluates the authored body. Dynamic uses are checked at call time. */ export function preflightHelpers( - definition: { header: { tools?: { slack?: boolean } }; body: Function }, + definition: { header?: { tools?: { slack?: boolean } }; body?: Function }, facts: { slackToken?: string; slackMount: boolean; slackMock: boolean }, ): PreflightResult { - const usesSlack = definition.header.tools?.slack === true - || /(?:\.\s*slack\b|\[\s*['"]slack['"]\s*\])/.test(Function.prototype.toString.call(definition.body)); + const usesSlack = definition.header?.tools?.slack === true + || (typeof definition.body === 'function' + && /(?:\.\s*slack\b|\[\s*['"]slack['"]\s*\])/.test(Function.prototype.toString.call(definition.body))); const diagnostics: PreflightDiagnostic[] = usesSlack && !facts.slackMock && !facts.slackToken?.trim() && !facts.slackMount ? [{ severity: 'refusal', kind: 'helper_slack.credential_missing', diff --git a/packages/sdk/tests/authored-flow-slack.test.ts b/packages/sdk/tests/authored-flow-slack.test.ts index d4250cf34..1330bbe5b 100644 --- a/packages/sdk/tests/authored-flow-slack.test.ts +++ b/packages/sdk/tests/authored-flow-slack.test.ts @@ -112,6 +112,23 @@ describe('authored Slack helper effects', () => { }).ok).toBe(true); }); + it('surfaces helper_slack.credential_missing on `.flow.ts` paths too, not just helper modules', async () => { + // Regression: `flows check` on `.flow.ts` routes through `checkAuthoredTriggers` + // (trigger-aware) while `.mjs` helper modules routed through `checkHelperBody`. + // Before this fix, the `.flow.ts` path silently skipped the helper preflight, + // so a flow using `f.slack.post` without SLACK_BOT_TOKEN passed check and only + // crashed at run. + for (const key of ['SLACK_BOT_TOKEN', 'RELAYFLOWS_SLACK_MOCK', 'RELAYFILE_MOUNT_PATH', 'WORKSPACE_ROOT', 'WORKFORCE_SANDBOX_ROOT', 'RELAYFILE_MOUNT_ROOT', 'RELAYFILE_ROOT']) vi.stubEnv(key, ''); + const dataDir = temporary(); + mkdirSync(join(dataDir, 'node_modules/@relayflows'), { recursive: true }); + symlinkSync(join(root, 'packages/sdk/node_modules/@relayflows/surface'), join(dataDir, 'node_modules/@relayflows/surface')); + const fixture = join(dataDir, 'slack.flow.ts'); + writeFileSync(fixture, `import { flow } from ${JSON.stringify(join(root, 'packages/sdk/node_modules/@relayflows/surface/dist/index.js'))};\nexport default flow('slack-authored', async f => { await f.slack.post('#test', 'hi'); f.done('success'); });`); + const output: string[] = []; + expect(await runCli(['check', '--json', fixture], { stdout: line => output.push(line), stderr: line => output.push(line) })).toBe(2); + expect(output.join('')).toContain('helper_slack.credential_missing'); + }); + it.each(['confirm', 'complete'] as const)('replays after SIGKILL before %s with the same token and one successful completion', async boundary => { vi.stubEnv('RELAYFLOWS_SLACK_MOCK', '1'); const dataDir = temporary(); diff --git a/packages/sdk/tests/direct-input.test.ts b/packages/sdk/tests/direct-input.test.ts index ade726805..a69039b50 100644 --- a/packages/sdk/tests/direct-input.test.ts +++ b/packages/sdk/tests/direct-input.test.ts @@ -107,18 +107,26 @@ describe('direct .flow.ts input through the built CLI and live runtime', () => { // `--no-spawn` keeps this case about the property it names. `flows run` now // starts a daemon when none is serving (kernel/DAEMON-LIFECYCLE.md §3), so // without the flag the refusal under test would be about the spawn rather - // than about the authored module never being imported. - it('does not import or execute authored code before daemon availability', () => { + // than about the authored BODY never being run. + // + // Trigger preflight (checkAuthoredTriggers) DOES import the authored + // module before daemon-attach — trigger sources are only observable + // after `flow(...).on(webhook(...))` has run — so the import-time + // marker is expected to exist. The load-bearing invariant is that the + // authored BODY does not run before the daemon is confirmed available; + // the body-run marker path is what this test checks. + it('does not run the authored body before daemon availability', () => { const directory = temporaryDirectory(); - const marker = join(directory, 'marker.txt'); + const bodyMarker = join(directory, 'body-marker.txt'); + const importMarker = join(directory, 'import-marker.txt'); const result = invokeCli([ - 'run', '--no-spawn', SIDE_EFFECT_FLOW, '--input', JSON.stringify({ marker }), + 'run', '--no-spawn', SIDE_EFFECT_FLOW, '--input', JSON.stringify({ marker: bodyMarker }), '--data-dir', join(directory, 'absent-daemon'), - ], { RELAYFLOWS_TEST_IMPORT_MARKER: marker }); + ], { RELAYFLOWS_TEST_IMPORT_MARKER: importMarker }); expect(result.status, result.stderr).toBe(2); expect(result.stderr).toContain('REFUSED [daemon_unreachable]'); - expect(existsSync(marker)).toBe(false); + expect(existsSync(bodyMarker)).toBe(false); }); it('refuses oversized file input before contacting relayflowd', () => { diff --git a/packages/sdk/tests/fixtures/pre-journal-side-effect.flow.ts b/packages/sdk/tests/fixtures/pre-journal-side-effect.flow.ts index dc20d2d82..c85ebde3e 100644 --- a/packages/sdk/tests/fixtures/pre-journal-side-effect.flow.ts +++ b/packages/sdk/tests/fixtures/pre-journal-side-effect.flow.ts @@ -1,6 +1,12 @@ import { writeFileSync } from 'node:fs'; import { flow } from '@relayflows/surface'; +// The test asserts the AUTHORED BODY does not run before daemon +// availability. Trigger preflight requires importing the module (see +// checkAuthoredTriggers in cli/check-triggers.ts) so an import-time +// side effect is allowed and, if the caller supplies +// RELAYFLOWS_TEST_IMPORT_MARKER, recorded — the body-run marker is +// distinct so the daemon-not-imported contract still has teeth. const importMarker = process.env['RELAYFLOWS_TEST_IMPORT_MARKER']; if (importMarker !== undefined) writeFileSync(importMarker, 'authored module imported'); diff --git a/packages/sdk/tests/webhook-live.test.ts b/packages/sdk/tests/webhook-live.test.ts new file mode 100644 index 000000000..3108a8e0c --- /dev/null +++ b/packages/sdk/tests/webhook-live.test.ts @@ -0,0 +1,129 @@ +import { afterEach, expect, it } from 'vitest'; +import { spawn, type ChildProcess } from 'node:child_process'; +import { existsSync } from 'node:fs'; +import { mkdtemp, mkdir, readFile, readdir, rename, rm, writeFile } from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import { join, resolve } from 'node:path'; +import { setTimeout as delay } from 'node:timers/promises'; +import { socketPathFor } from '../src/daemon-connection.js'; +import { JournalClient } from '../src/journal-client.js'; + +const binary = process.env['RELAYFLOWD_BIN'] ?? resolve('../../kernel/target/debug/relayflowd'); +const children: ChildProcess[] = []; +const directories: string[] = []; +const clients: JournalClient[] = []; +afterEach(async () => { + clients.splice(0).forEach(client => client.close()); + await Promise.all(children.splice(0).map(child => stop(child))); + await Promise.all(directories.splice(0).map(dir => rm(dir, { recursive: true, force: true }))); +}); +async function stop(child: ChildProcess): Promise { + if (child.exitCode !== null || child.signalCode !== null) return; + await new Promise(done => { child.once('exit', () => done()); child.kill('SIGKILL'); }); +} +function start(command: string, args: string[]): { child: ChildProcess; output: () => string } { + const child = spawn(command, args, { stdio: ['ignore', 'pipe', 'pipe'] }); + children.push(child); + let output = ''; + child.stdout!.on('data', chunk => { output += String(chunk); }); + child.stderr!.on('data', chunk => { output += String(chunk); }); + child.on('error', error => { output += error.message; }); + return { child, output: () => output }; +} +async function until(predicate: () => boolean | Promise, detail: () => string = () => ''): Promise { + const deadline = Date.now() + 10_000; + while (Date.now() < deadline) { if (await predicate()) return; await delay(20); } + throw new Error(`webhook integration timed out: ${detail()}`); +} +async function daemon(dir: string): Promise { + const process = start(binary, ['--data-dir', dir, 'serve']); + await until(async () => { + const client = new JournalClient(socketPathFor(dir), { requestTimeoutMs: 200 }); + try { await client.connect(); await client.hello('webhook-test'); return true; } + catch { return false; } finally { client.close(); } + }, process.output); + return process.child; +} +async function setup(command = 'printf accepted'): Promise<{ dir: string; base: string }> { + const dir = await mkdtemp(join(tmpdir(), 'flows-inbox-live-')); + directories.push(dir); + await mkdir(join(dir, 'triggers')); + await writeFile(join(dir, 'triggers', 'release.json'), JSON.stringify({ + version: '0.1.0', name: 'release', + triggers: [{ id: 'release', executor: 'release', event_type: 'release', dedupe_key_template: '{{event.type}}' }], + steps: [{ id: 'log', type: 'deterministic', command }], + })); + const receiver = start(process.execPath, [resolve('dist/cli.js'), 'serve-webhook', '--data-dir', dir, '--port', '0']); + await until(() => /WEBHOOK http:\/\/127.0.0.1:\d+/.test(receiver.output()), receiver.output); + return { dir, base: receiver.output().match(/http:\/\/127.0.0.1:\d+/)![0] }; +} +async function post(base: string): Promise { + const response = await fetch(`${base}/release`, { method: 'POST', body: JSON.stringify({ action: 'released', nested: [1, null] }) }); + expect(response.status).toBe(202); + return `${(await response.json() as { id: string }).id}.json`; +} +async function journals(dir: string): Promise { + if (!existsSync(join(dir, 'runs'))) return []; + return (await readdir(join(dir, 'runs'))).filter(file => file.endsWith('.sqlite3')); +} +async function readJournal(dir: string, runFile: string) { + const client = new JournalClient(socketPathFor(dir)); + clients.push(client); + await client.connect(); + return client.journalRead(runFile.slice(0, -8), 1); +} + +it('flows serve-webhook writes JSON before the daemon starts, then journals and archives exactly once', async () => { + const { dir, base } = await setup(); + const filename = await post(base); + const inbox = join(dir, 'inbox/release', filename); + expect(JSON.parse(await readFile(inbox, 'utf8'))).toEqual({ action: 'released', nested: [1, null] }); + expect(await journals(dir)).toEqual([]); + await daemon(dir); + const processed = join(dir, 'inbox-processed/release', filename); + await until(() => existsSync(processed)); + const files = await journals(dir); + expect(files).toHaveLength(1); + const journal = await readJournal(dir, files[0]!); + expect(journal.entries.find(entry => entry.entry_type === 'run.spawned')?.payload['event']).toEqual({ action: 'released', nested: [1, null] }); + expect(journal.entries.filter(entry => entry.entry_type === 'step.completed')).toHaveLength(1); + await rename(processed, inbox); + await until(() => existsSync(processed)); + expect(await journals(dir)).toEqual(files); + expect((await readJournal(dir, files[0]!)).entries.filter(entry => entry.entry_type === 'step.completed')).toHaveLength(1); +}, 20_000); + +it('replays a dropped file after SIGKILL before spawn', async () => { + const { dir, base } = await setup(); + const first = await daemon(dir); + first.kill('SIGSTOP'); + const filename = await post(base); + expect(existsSync(join(dir, 'inbox/release', filename))).toBe(true); + await stop(first); + expect(await journals(dir)).toEqual([]); + await daemon(dir); + await until(() => existsSync(join(dir, 'inbox-processed/release', filename))); + expect(await journals(dir)).toHaveLength(1); +}, 20_000); + +it('resumes the same journal after SIGKILL after spawn and before acknowledgement', async () => { + const { dir, base } = await setup('sleep 2; printf accepted'); + const first = await daemon(dir); + const filename = await post(base); + await until(async () => (await journals(dir)).length === 1); + // Wait until a real attempt exists, proving the run is registered and driven. + const files = await journals(dir); + await until(async () => { + const journal = await readJournal(dir, files[0]!); + return journal.entries.some(entry => entry.entry_type === 'step.attempt.started'); + }); + clients.splice(0).forEach(client => client.close()); + await stop(first); + expect(existsSync(join(dir, 'inbox/release', filename))).toBe(true); + await daemon(dir); + await until(() => existsSync(join(dir, 'inbox-processed/release', filename))); + expect(await journals(dir)).toEqual(files); + const journal = await readJournal(dir, files[0]!); + expect(journal.entries.filter(entry => entry.entry_type === 'run.spawned')).toHaveLength(1); + expect(journal.entries.some(entry => entry.entry_type === 'run.completed')).toBe(true); +}, 20_000); diff --git a/packages/sdk/tests/webhook.test.ts b/packages/sdk/tests/webhook.test.ts new file mode 100644 index 000000000..593394646 --- /dev/null +++ b/packages/sdk/tests/webhook.test.ts @@ -0,0 +1,96 @@ +import { afterEach, describe, expect, it } from 'vitest'; +import { mkdtemp, readFile, readdir, rm, symlink, writeFile } from 'node:fs/promises'; +import { join, resolve } from 'node:path'; +import { tmpdir } from 'node:os'; +import type { Server } from 'node:http'; +import { startWebhookServer, parseWebhookArgs } from '../src/cli/serve-webhook.js'; +import { preflightWebhookTriggers } from '../src/preflight.js'; +import { webhook } from '@relayflows/surface'; +import { runCli } from '../src/cli.js'; +import { runDirectFlow } from '../src/cli/direct-run.js'; + +const dirs: string[] = []; +const servers: Server[] = []; +afterEach(async () => { + for (const server of servers.splice(0)) await new Promise(done => { + server.close(() => done()); server.closeAllConnections(); + }); + await Promise.all(dirs.splice(0).map(dir => rm(dir, { recursive: true, force: true }))); +}); +async function temporary(): Promise { + const dir = await mkdtemp(join(tmpdir(), 'flows-webhook-')); + dirs.push(dir); return dir; +} +async function receiver(dir: string): Promise { + const server = await startWebhookServer(dir, 0); + servers.push(server); + const address = server.address(); + if (!address || typeof address === 'string') throw new Error('missing HTTP address'); + return `http://127.0.0.1:${address.port}`; +} + +describe('webhook ingress', () => { + it('requires registered executor names with the exact refusal', () => { + expect(preflightWebhookTriggers([webhook('unregistered')], [])).toEqual([{ + severity: 'refusal', kind: 'no_executor', executor: 'unregistered', + message: 'webhook trigger "unregistered" is not registered in flows.json', + }]); + expect(preflightWebhookTriggers([webhook('registered')], ['registered'])).toEqual([]); + }); + it('checks TS declarations against flows.json without invoking handlers', async () => { + const dir = await temporary(); + await symlink(resolve('node_modules'), join(dir, 'node_modules'), 'dir'); + await writeFile(join(dir, 'package.json'), '{"type":"module"}'); + const path = join(dir, 'test.flow.ts'); + await writeFile(path, `import { flow, webhook } from '@relayflows/surface';\nexport default flow('test').on(webhook('unregistered'), async () => { throw new Error('handler ran'); });`); + const reports: string[] = []; + const io = { stdout: (line: string) => reports.push(line), stderr: (line: string) => reports.push(line) }; + await writeFile(join(dir, 'flows.json'), '{"executors":[]}'); + expect(await runCli(['check', '--json', path], io)).toBe(2); + expect(reports.join('\n')).toContain('no_executor'); + const refusedRun = await runDirectFlow(path, '{}', join(dir, 'daemon'), { daemon: { spawn: false } }); + expect(refusedRun.exitCode).toBe(2); + expect(refusedRun.report.diagnostics[0]?.kind).toBe('no_executor'); + expect(await readdir(dir)).not.toContain('daemon'); + reports.length = 0; + await writeFile(join(dir, 'flows.json'), '{"executors":["unregistered"]}'); + expect(await runCli(['check', '--json', path], io)).toBe(0); + }); + it('accepts JSON through atomic files without creating a daemon', async () => { + const dir = await temporary(); + const base = await receiver(dir); + const payload = { nested: [1, null, 'é'], action: 'release' }; + const response = await fetch(`${base}/release`, { method: 'POST', body: JSON.stringify(payload) }); + expect(response.status).toBe(202); + const receipt = await response.json() as { id: string }; + expect(receipt.id).toMatch(/^[0-9a-f-]{36}$/); + expect(await readdir(join(dir, 'inbox', 'release'))).toEqual([`${receipt.id}.json`]); + expect(JSON.parse(await readFile(join(dir, 'inbox', 'release', `${receipt.id}.json`), 'utf8'))).toEqual(payload); + expect(await readdir(dir)).toEqual(['inbox']); + }); + it('rejects malformed, oversized, traversal, and non-POST requests', async () => { + const dir = await temporary(); + const base = await receiver(dir); + expect((await fetch(`${base}/release`)).status).toBe(405); + expect((await fetch(`${base}/bad%2Fpath`, { method: 'POST', body: '{}' })).status).toBe(404); + expect((await fetch(`${base}/release`, { method: 'POST', body: '{' })).status).toBe(400); + expect((await fetch(`${base}/release`, { method: 'POST', body: 'x'.repeat(1024 * 1024 + 1) })).status).toBe(413); + expect((await fetch(`${base}/release`, { method: 'POST', body: '{"value":1e400}' })).status).toBe(400); + expect(await readdir(join(dir, 'inbox'))).toEqual([]); + }); + it('fails closed on a symlink inbox target', async () => { + const dir = await temporary(); + const outside = await temporary(); + const base = await receiver(dir); + await symlink(outside, join(dir, 'inbox', 'release'), 'dir'); + expect((await fetch(`${base}/release`, { method: 'POST', body: '{}' })).status).toBe(500); + expect(await readdir(outside)).toEqual([]); + }); + it('parses CLI options strictly', () => { + expect(parseWebhookArgs(['--data-dir', 'data', '--port', '0'])).toEqual({ command: 'serve-webhook', dataDir: 'data', port: 0 }); + for (const args of [[], ['--port', '80'], ['--data-dir', 'x', '--port', '65536'], + ['--data-dir', 'x', '--port', '1', '--port', '2'], ['--data-dir', 'x', '--port', '3.1']]) { + expect(parseWebhookArgs(args)).toBeUndefined(); + } + }); +}); diff --git a/packages/surface/src/flow.ts b/packages/surface/src/flow.ts index 4013f7358..4eec2c712 100644 --- a/packages/surface/src/flow.ts +++ b/packages/surface/src/flow.ts @@ -1,4 +1,5 @@ import type { Ctx } from "./context.js"; +import { webhook, type TriggerSource } from "./triggers.js"; /** Optional escalation header; the empty header is the common case. */ export interface FlowHeader { @@ -31,15 +32,25 @@ export interface AuthoredFlowDefinition { readonly name: string; readonly header: ReadonlyFlowHeader; readonly body: FlowBody; + readonly handlers: readonly TriggerHandler[]; +} + +export interface TriggerHandler { + readonly trigger: TriggerSource; + readonly body: FlowBody; } /** Opaque authored-flow handle. Execution stays behind the journal runtime. */ export interface FlowHandle { readonly name: string; + on>(trigger: TriggerSource, body: FlowBody): TriggeredFlowHandle; } +export interface TriggeredFlowHandle extends FlowHandle {} + const definitions = new WeakMap(); +export function flow(name: string, header?: FlowHeader): FlowHandle; export function flow(name: string, body: FlowBody): FlowHandle; export function flow( name: string, @@ -48,7 +59,7 @@ export function flow( ): FlowHandle; export function flow( name: string, - headerOrBody: FlowHeader | FlowBody, + headerOrBody: FlowHeader | FlowBody = {}, body?: FlowBody, ): FlowHandle { const flowBody = typeof headerOrBody === "function" ? headerOrBody : body; @@ -57,7 +68,7 @@ export function flow( if (name.trim().length === 0) { throw new TypeError("flow name must not be empty"); } - if (typeof flowBody !== "function") { + if (flowBody !== undefined && typeof flowBody !== "function") { throw new TypeError(`flow "${name}" requires a body`); } assertFlowHeader(header, name); @@ -65,9 +76,28 @@ export function flow( const definition: AuthoredFlowDefinition = Object.freeze({ name, header: freezeHeader(header), - body: flowBody, + body: flowBody ?? (async () => { throw new TypeError(`flow "${name}" has no direct-run body`); }), + handlers: Object.freeze([]), + }); + return makeHandle(definition as AuthoredFlowDefinition); +} + +function makeHandle(definition: AuthoredFlowDefinition): FlowHandle { + const handle = { name: definition.name } as FlowHandle; + Object.defineProperty(handle, "on", { + value: (trigger: TriggerSource, body: FlowBody): TriggeredFlowHandle => { + if (typeof body !== "function") throw new TypeError("trigger handler requires a body"); + assertHeaderObject(trigger, "trigger"); + assertKnownKeys(trigger, ["kind", "name", "filter"], "trigger"); + if (trigger.kind !== "webhook") throw new TypeError("unsupported trigger kind"); + const source = webhook(trigger.name, trigger.filter); + return makeHandle(Object.freeze({ + ...definition, + handlers: Object.freeze([...definition.handlers, Object.freeze({ trigger: source, body: body as FlowBody })]), + })); + }, }); - const handle: FlowHandle = Object.freeze({ name }); + Object.freeze(handle); // One map holds definitions of many input types, so it is stored at the // default parameterisation and `getFlowDefinition` re-parameterises on // the way out. The cast is needed because `body` puts `Input` in a parameter @@ -75,7 +105,7 @@ export function flow( // assignable to `AuthoredFlowDefinition` even though every read // recovers the author's own type. Sound here because the handle-to-definition // pairing is 1:1 and both sides are keyed by the same authored flow. - definitions.set(handle, definition as AuthoredFlowDefinition); + definitions.set(handle, definition); return handle; } @@ -83,7 +113,7 @@ export function flow( * Runtime bridge used by the SDK after it imports an authored `.flow.ts`. * The root package deliberately does not re-export this accessor. */ -export function getFlowDefinition(handle: FlowHandle): AuthoredFlowDefinition { +export function getFlowDefinition(handle: Pick): AuthoredFlowDefinition { if ((typeof handle !== "object" && typeof handle !== "function") || handle === null) { throw new TypeError("expected an @relayflows/surface flow handle"); } diff --git a/packages/surface/src/index.ts b/packages/surface/src/index.ts index dbb2d6962..7af220819 100644 --- a/packages/surface/src/index.ts +++ b/packages/surface/src/index.ts @@ -18,8 +18,10 @@ export type { Step } from "./step.js"; export { flow, type FlowHandle, + type TriggeredFlowHandle, type FlowHeader, } from "./flow.js"; export { flowRunWritebackIdempotency, type SlackHelper, type SlackReceipt } from "./slack.js"; export type { Helpers } from "./helpers/index.js"; export type { MemoryHelper, MemoryFinding, MemoryRecallOptions, HistoryEntry, TrajectoryEntry } from "./memory.js"; +export { webhook, type TriggerSource, type WebhookFilter, type WebhookValue } from "./triggers.js"; diff --git a/packages/surface/src/runtime.ts b/packages/surface/src/runtime.ts index 880ef322c..08663667f 100644 --- a/packages/surface/src/runtime.ts +++ b/packages/surface/src/runtime.ts @@ -3,5 +3,7 @@ export { type AuthoredFlowDefinition, type FlowBody, type FlowHandle, + type TriggerHandler, + type TriggeredFlowHandle, type ReadonlyFlowHeader, } from "./flow.js"; diff --git a/packages/surface/src/triggers.ts b/packages/surface/src/triggers.ts new file mode 100644 index 000000000..60b304b56 --- /dev/null +++ b/packages/surface/src/triggers.ts @@ -0,0 +1,49 @@ +export type WebhookValue = null | boolean | number | string + | readonly WebhookValue[] | { readonly [key: string]: WebhookValue }; + +/** Recursive object subset; array and scalar leaves match exactly. */ +export type WebhookFilter = { readonly [key: string]: WebhookValue }; + +export interface TriggerSource { + readonly kind: "webhook"; + readonly name: string; + readonly filter?: WebhookFilter; +} + +/** Plain, immutable data. Constructing a source opens no receiver. */ +export function webhook(name: string, filter?: WebhookFilter): TriggerSource { + if (typeof name !== "string" || !/^[A-Za-z0-9][A-Za-z0-9_-]{0,127}$/.test(name)) { + throw new TypeError("webhook name must be 1-128 letters, digits, underscores or hyphens, starting with a letter or digit"); + } + if (filter === undefined) return Object.freeze({ kind: "webhook", name }); + if (filter === null || typeof filter !== "object" || Array.isArray(filter)) { + throw new TypeError("webhook filter must be a JSON object"); + } + return Object.freeze({ kind: "webhook", name, filter: snapshot(filter) as WebhookFilter }); +} + +function snapshot(value: unknown, ancestors = new Set()): WebhookValue { + if (value === null || typeof value === "string" || typeof value === "boolean") return value; + if (typeof value === "number" && Number.isFinite(value)) return value; + if (typeof value !== "object" || value === null || ancestors.has(value)) { + throw new TypeError("webhook filter must contain finite, acyclic JSON data"); + } + const array = Array.isArray(value); + if (!array && Object.getPrototypeOf(value) !== Object.prototype && Object.getPrototypeOf(value) !== null) { + throw new TypeError("webhook filter must contain plain JSON objects"); + } + ancestors.add(value); + const entries: [string, WebhookValue][] = []; + for (const key of Reflect.ownKeys(value)) { + if (array && key === "length") continue; + const descriptor = Object.getOwnPropertyDescriptor(value, key)!; + if (typeof key !== "string" || !descriptor.enumerable || !("value" in descriptor)) { + throw new TypeError("webhook filter must contain JSON data properties"); + } + if (array && !/^(0|[1-9][0-9]*)$/.test(key)) throw new TypeError("invalid JSON array property"); + entries.push([key, snapshot(descriptor.value, ancestors)]); + } + ancestors.delete(value); + if (array && entries.length !== value.length) throw new TypeError("webhook filter arrays must not be sparse"); + return Object.freeze(array ? entries.map(([, item]) => item) : Object.fromEntries(entries)); +} diff --git a/packages/surface/tests/triggers.test.ts b/packages/surface/tests/triggers.test.ts new file mode 100644 index 000000000..df1337c48 --- /dev/null +++ b/packages/surface/tests/triggers.test.ts @@ -0,0 +1,50 @@ +import { describe, expect, it } from "vitest"; +import { flow, webhook, type Ctx, type WebhookFilter } from "@relayflows/surface"; +import { getFlowDefinition } from "@relayflows/surface/runtime"; + +describe("webhook declarations", () => { + it("records handlers without executing them and preserves immutable chains", async () => { + const events: unknown[] = []; + const body = async (_f: Ctx, event: unknown): Promise => { events.push(event); }; + const base = flow("chief", { identity: "chief" }); + const first = base.on(webhook("release"), body); + const second = first.on(webhook("deploy"), body); + expect(getFlowDefinition(base).handlers).toEqual([]); + expect(getFlowDefinition(first).handlers).toEqual([{ trigger: { kind: "webhook", name: "release" }, body }]); + expect(getFlowDefinition(second).handlers).toHaveLength(2); + expect(Object.isFrozen(getFlowDefinition(second).handlers)).toBe(true); + expect(events).toEqual([]); + await getFlowDefinition(first).handlers[0]!.body({} as Ctx, { tag: "v1" }); + expect(events).toEqual([{ tag: "v1" }]); + }); + + it("keeps direct bodies distinct from event handlers", () => { + const direct = async (): Promise => undefined; + const handler = async (): Promise => undefined; + const declared = flow("both", direct).on(webhook("release"), handler); + expect(getFlowDefinition(declared).body).toBe(direct); + expect(getFlowDefinition(declared).handlers[0]!.body).toBe(handler); + }); + + it("snapshots and deeply freezes filter data", () => { + const filter = { action: "released", nested: { tags: ["a"] } }; + const source = webhook("release", filter); + filter.nested.tags.push("b"); + expect(source.filter).toEqual({ action: "released", nested: { tags: ["a"] } }); + expect(Object.isFrozen(source.filter?.nested)).toBe(true); + expect(JSON.parse(JSON.stringify(source))).toEqual(source); + }); + + it("refuses paths, non-data filters, and invalid handlers", () => { + for (const name of ["", ".", "..", "a/b", "a%2fb", "a\\b", "a b"]) { + expect(() => webhook(name)).toThrow(); + } + const circular: Record = {}; + circular.self = circular; + for (const filter of [null, [], { x: undefined }, { x: NaN }, { x: new Date() }, circular, + Object.defineProperty({}, "x", { get: () => { throw new Error("getter executed"); } })]) { + expect(() => webhook("release", filter as WebhookFilter)).toThrow(/webhook filter/); + } + expect(() => flow("empty").on(webhook("release"), undefined as never)).toThrow("requires a body"); + }); +});