From 099cbe5150e238a6deaa18e6722cd07ba42f677a Mon Sep 17 00:00:00 2001 From: Miya Date: Fri, 11 Sep 2026 10:16:57 +0200 Subject: [PATCH 1/5] feat(surface,sdk,kernel): event triggers via webhook and inbox watcher (#301) Session-Id: 01a08f7e-5920-7321-9c9b-ed10eb8ac4b8 Session-Id: efeda5df-9b7c-48d4-b2ce-957f5bef0a82 Session-Id: efeda5df-9b7c-48d4-b2ce-957f5bef0a82 Session-Id: efeda5df-9b7c-48d4-b2ce-957f5bef0a82 Session-Id: efeda5df-9b7c-48d4-b2ce-957f5bef0a82 Session-Id: efeda5df-9b7c-48d4-b2ce-957f5bef0a82 Session-Id: efeda5df-9b7c-48d4-b2ce-957f5bef0a82 --- kernel/relayflowd/src/engine/wake.rs | 73 ++++++++--- kernel/relayflowd/src/lib.rs | 1 + kernel/relayflowd/src/server.rs | 10 ++ kernel/relayflowd/src/trigger_watcher.rs | 137 ++++++++++++++++++++ kernel/relayflowd/tests/trigger_watcher.rs | 113 +++++++++++++++++ packages/sdk/src/cli.ts | 16 ++- packages/sdk/src/cli/check-triggers.ts | 32 +++++ packages/sdk/src/cli/direct-run.ts | 10 +- packages/sdk/src/cli/serve-webhook.ts | 141 +++++++++++++++++++++ packages/sdk/src/preflight.ts | 16 +++ packages/sdk/tests/webhook-live.test.ts | 129 +++++++++++++++++++ packages/sdk/tests/webhook.test.ts | 96 ++++++++++++++ packages/surface/src/flow.ts | 42 +++++- packages/surface/src/index.ts | 2 + packages/surface/src/runtime.ts | 2 + packages/surface/src/triggers.ts | 49 +++++++ packages/surface/tests/triggers.test.ts | 50 ++++++++ 17 files changed, 892 insertions(+), 27 deletions(-) create mode 100644 kernel/relayflowd/src/trigger_watcher.rs create mode 100644 kernel/relayflowd/tests/trigger_watcher.rs create mode 100644 packages/sdk/src/cli/check-triggers.ts create mode 100644 packages/sdk/src/cli/serve-webhook.ts create mode 100644 packages/sdk/tests/webhook-live.test.ts create mode 100644 packages/sdk/tests/webhook.test.ts create mode 100644 packages/surface/src/triggers.ts create mode 100644 packages/surface/tests/triggers.test.ts 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..1fa3a1a1c --- /dev/null +++ b/packages/sdk/src/cli/check-triggers.ts @@ -0,0 +1,32 @@ +import { dirname, resolve } from 'node:path'; +import { loadAuthoredFlow, type LoadedAuthoredFlow } from '../authored-flow-loader.js'; +import { preflightWebhookTriggers } from '../preflight.js'; +import { inputFailureReport, readProjectConfig, type CheckReport } from './check.js'; + +/** Inspect declarations without running a handler or contacting the daemon. */ +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 diagnostics = preflightWebhookTriggers( + (definition.handlers ?? []).map(handler => handler.trigger), config.executors, + ); + 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/direct-run.ts b/packages/sdk/src/cli/direct-run.ts index d9bc82652..4adbb2ff7 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,11 @@ export async function runDirectFlow( }; } + // Declared triggers are knowable before any daemon or step is started. + 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 +58,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..8311e2505 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; 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"); + }); +}); From 653c09f826465b716b7bad98c702c5ab33630306 Mon Sep 17 00:00:00 2001 From: kjgbot Date: Fri, 11 Sep 2026 12:28:27 +0200 Subject: [PATCH 2/5] fix(direct-run): distinguish import from body-run in pre-journal invariant MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit E's trigger check imports the authored module before daemon-attach because trigger sources are only observable after `flow(...).on(webhook(...))` has run — the CI failure on direct-input.test.ts:121 was the old assertion catching this legitimate import. The load-bearing invariant is that the authored BODY is not called before daemon availability; the fixture now records import vs body separately and the test targets the body-run marker. Session-Id: efeda5df-9b7c-48d4-b2ce-957f5bef0a82 Session-Id: efeda5df-9b7c-48d4-b2ce-957f5bef0a82 Session-Id: efeda5df-9b7c-48d4-b2ce-957f5bef0a82 Session-Id: efeda5df-9b7c-48d4-b2ce-957f5bef0a82 Session-Id: efeda5df-9b7c-48d4-b2ce-957f5bef0a82 Session-Id: efeda5df-9b7c-48d4-b2ce-957f5bef0a82 --- packages/sdk/src/cli/direct-run.ts | 4 ++++ packages/sdk/tests/direct-input.test.ts | 20 +++++++++++++------ .../fixtures/pre-journal-side-effect.flow.ts | 6 ++++++ 3 files changed, 24 insertions(+), 6 deletions(-) diff --git a/packages/sdk/src/cli/direct-run.ts b/packages/sdk/src/cli/direct-run.ts index 4adbb2ff7..80fbf7199 100644 --- a/packages/sdk/src/cli/direct-run.ts +++ b/packages/sdk/src/cli/direct-run.ts @@ -43,6 +43,10 @@ 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) }; 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'); From bff827f4d98588172acbf8b402fed30be1070e18 Mon Sep 17 00:00:00 2001 From: kjgbot Date: Fri, 11 Sep 2026 12:43:24 +0200 Subject: [PATCH 3/5] fix(check-triggers): compose f.slack helper preflight for .flow.ts paths Cursor Bugbot MED at packages/sdk/src/cli.ts:115: flows check on a .flow.ts routed through checkAuthoredTriggers (webhook-trigger preflight) but silently skipped checkHelperBody (added for f.slack in PR#314). A .flow.ts using f.slack.post without SLACK_BOT_TOKEN passed check and only crashed at run. Fix by folding the slack helper preflight into checkAuthoredTriggers on the same loaded definition, merging diagnostics. Direct-run also inherits the helper check because it calls the same function. Regression test: a .flow.ts using f.slack.post with no SLACK_BOT_TOKEN must produce helper_slack.credential_missing via flows check --json. Session-Id: efeda5df-9b7c-48d4-b2ce-957f5bef0a82 Session-Id: efeda5df-9b7c-48d4-b2ce-957f5bef0a82 Session-Id: efeda5df-9b7c-48d4-b2ce-957f5bef0a82 Session-Id: efeda5df-9b7c-48d4-b2ce-957f5bef0a82 Session-Id: efeda5df-9b7c-48d4-b2ce-957f5bef0a82 Session-Id: efeda5df-9b7c-48d4-b2ce-957f5bef0a82 --- packages/sdk/src/cli/check-triggers.ts | 13 +++++++++++-- packages/sdk/tests/authored-flow-slack.test.ts | 17 +++++++++++++++++ 2 files changed, 28 insertions(+), 2 deletions(-) diff --git a/packages/sdk/src/cli/check-triggers.ts b/packages/sdk/src/cli/check-triggers.ts index 1fa3a1a1c..aa76d5ef1 100644 --- a/packages/sdk/src/cli/check-triggers.ts +++ b/packages/sdk/src/cli/check-triggers.ts @@ -1,9 +1,16 @@ 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 declarations without running a handler or contacting the daemon. */ +/** + * 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; @@ -12,9 +19,11 @@ export async function checkAuthoredTriggers(path: string): Promise<{ const loaded = await loadAuthoredFlow(path); const definition = loaded.getDefinition(loaded.handle); const config = readProjectConfig(dirname(resolve(path))); - const diagnostics = preflightWebhookTriggers( + const triggerDiagnostics = preflightWebhookTriggers( (definition.handlers ?? []).map(handler => handler.trigger), config.executors, ); + const helperReport = checkSlackHelpers(definition); + const diagnostics = [...triggerDiagnostics, ...helperReport.diagnostics]; return { loaded, report: { 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(); From ebfadb792e7a93615f0b85717cbb1c288289bb7c Mon Sep 17 00:00:00 2001 From: kjgbot Date: Fri, 11 Sep 2026 16:02:39 +0200 Subject: [PATCH 4/5] fix(sdk): guard against missing header in preflight paths reachable via mocked authored flows direct-run-failure.test.ts mocks getDefinition to return {} to simulate loader failures. Both preflightHelpers (slack) and checkMcpHeader unguarded-read definition.header, causing the tests to trap on TypeError instead of surfacing the mocked worker-error classification. Session-Id: efeda5df-9b7c-48d4-b2ce-957f5bef0a82 Session-Id: efeda5df-9b7c-48d4-b2ce-957f5bef0a82 Session-Id: efeda5df-9b7c-48d4-b2ce-957f5bef0a82 --- packages/sdk/src/cli/check-typescript.ts | 7 ++++--- packages/sdk/src/preflight.ts | 7 ++++--- 2 files changed, 8 insertions(+), 6 deletions(-) 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/preflight.ts b/packages/sdk/src/preflight.ts index 8311e2505..77b9564a3 100644 --- a/packages/sdk/src/preflight.ts +++ b/packages/sdk/src/preflight.ts @@ -560,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', From 4e0a02a01709f3f4b2d90ca5e83a174aceb0e742 Mon Sep 17 00:00:00 2001 From: kjgbot Date: Fri, 11 Sep 2026 17:14:10 +0200 Subject: [PATCH 5/5] fix(kernel/journal): atomic schema+meta init to survive SIGKILL race MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit SqliteJournal::create ran execute_batch(SCHEMA) in autocommit before the meta INSERT transaction. SIGKILL in that microsecond window left the file with schema but no meta row, so open()'s SELECT ... FROM meta returned NoRows (surfaced as 'Query returned no rows' on resume). The webhook-live SIGKILL-after-spawn-before-ack test reproduced this deterministically on GH runners. Fold DDL + meta/segment inserts into a single WAL transaction. Either SIGKILL leaves an empty file (no tables) or a fully initialized journal — nothing in between. Session-Id: efeda5df-9b7c-48d4-b2ce-957f5bef0a82 --- kernel/relayflowd-journal/src/lib.rs | 25 +++++++++++++++---------- 1 file changed, 15 insertions(+), 10 deletions(-) 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 {