From 0b71713a5b01ea87188bbabbc706b7e81134df42 Mon Sep 17 00:00:00 2001 From: Cyrus Gray <144336577+graycyrus@users.noreply.github.com> Date: Mon, 6 Jul 2026 22:57:49 +0530 Subject: [PATCH] fix(flows): builder produces resolvable graphs (dry-run catches null args + correct envelope binding) (#4586) --- .../flows/agents/workflow_builder/prompt.md | 59 ++- src/openhuman/flows/builder_tools.rs | 257 ++++++++++++- src/openhuman/flows/builder_tools_tests.rs | 359 ++++++++++++++++++ src/openhuman/flows/ops.rs | 194 +++++++++- src/openhuman/flows/ops_tests.rs | 166 ++++++++ src/openhuman/flows/tools.rs | 21 +- src/openhuman/flows/tools_tests.rs | 43 +++ src/openhuman/tinyflows/caps.rs | 266 ++++++++++++- src/openhuman/tinyflows/tests.rs | 2 +- 9 files changed, 1340 insertions(+), 27 deletions(-) diff --git a/src/openhuman/flows/agents/workflow_builder/prompt.md b/src/openhuman/flows/agents/workflow_builder/prompt.md index cc98c53ae..49b5610b3 100644 --- a/src/openhuman/flows/agents/workflow_builder/prompt.md +++ b/src/openhuman/flows/agents/workflow_builder/prompt.md @@ -108,6 +108,18 @@ connected, `list_flow_connections` → build the `tool_call` node with the real 3. **Build the graph** (see the model below). 4. **Self-check with `dry_run_workflow`** on the draft — catch missing edges, wrong ports, unreachable nodes. Fix and re-run. + + **Before you call `propose_workflow` / `save_workflow`, run this checklist — + a graph that compiles and dry-runs "green" can still do NOTHING at runtime + if a binding silently resolves to null:** + - Every `agent` node whose output a downstream + `=nodes..item.json.` binding reads MUST declare + `config.output_parser.schema` naming that field under `properties`. No + schema ⇒ the agent's item is `{text: "..."}` and the binding is null. + - If `dry_run_workflow` reports `"ok": false` with a `null_resolutions` + list, **fix every one** before proposing — add the missing schema, or + rewire the expression to a real upstream field. Don't propose/save a + graph `dry_run_workflow` flagged. 5. **`propose_workflow`** (first draft) or **`revise_workflow`** (iterating on a prior draft — apply the change to the existing graph, don't regenerate from scratch). If validation fails, read the error, fix the graph, call again. @@ -131,15 +143,18 @@ A `WorkflowGraph` is `{ name?, nodes: [...], edges: [...] }`. **If the agent's output feeds a `tool_call`, it MUST declare an output schema** — set `config.output_parser.schema` (a JSON Schema object) — so its emitted item is a structured object whose fields downstream nodes can - address (`=nodes..item.`). Without a schema the agent - emits `{text: "..."}` and `=item.to`-style bindings resolve to null. + address (`=nodes..item.json.` — see "the envelope" below). + Without a schema the agent emits `{text: "..."}` (no other fields) and any + `.item.json.`-style binding to it resolves to null. 3. **`tool_call`** — an action. Two flavours by `config.slug`: - **Composio app action** — `config.slug` = a real action slug (from `search_tool_catalog`, e.g. `GMAIL_SEND_EMAIL`) + `config.connection_ref` for the account. **Wire every REQUIRED arg in `config.args` from a named upstream node** — e.g. an email send needs `to`/`recipient_email`, usually - `"to": "=nodes..item.email"`. A required arg left unwired (or - whose expression misses) now fails BEFORE the provider call — both in + `"to": "=nodes..item.json.email"` (drop `.json` only if + `` is a `code`/`transform`/`split_out`/`merge`/`trigger` node + — see "the envelope" below). A required arg left unwired (or whose + expression misses) now fails BEFORE the provider call — both in `dry_run_workflow` and in real runs — with an error naming the field. - **Native OpenHuman tool** — `config.slug` = `oh:` (e.g. `oh:web_search`) to call one of the assistant's own built-in tools (search, @@ -177,15 +192,34 @@ The scope exposes: - `nodes` — **every completed node's output, keyed by node id**: `nodes..item` (first item) and `nodes..items` (all items). Use this to reference ANY upstream node — not just the immediate predecessor — and to - disambiguate a fan-in node's inputs. Dotted form: `"=nodes.fetch.item.email"`; - jq form: `"=.nodes[\"fetch\"].items[0].email"`. Ids (not names) are the key. + disambiguate a fan-in node's inputs. Ids (not names) are the key. Use expressions to thread data between steps (a `transform`'s `set`, an `agent`'s `prompt`, a `tool_call`'s `args`). Prefer `=nodes..…` for `tool_call` args so the binding survives graph re-wiring. +**The envelope — `.item` vs. `.item.json`.** `agent`, `tool_call`, and +`http_request` nodes wrap their result in a stable +`{ json, text, raw }` envelope, so `nodes..item` for one of THOSE node +kinds is that envelope, NOT the structured value itself: + +- Structured fields live under **`.json`** — `"=nodes..item.json."` + (jq: `"=.nodes[\"\"].items[0].json."`). +- Prose lives under **`.text`** — `"=nodes..item.text"`. +- `code`, `transform`, `split_out`, `merge`, `output_parser`, `sub_workflow`, + and `trigger` nodes do **NOT** envelope — their output is addressed directly, + `"=nodes..item."`, same as the ungrouped `item`/`items` scope + entries above (which are always the raw predecessor value, envelope + included when the predecessor is one of the three enveloping kinds). + +**Getting this wrong is the single most common way a graph "builds" (compiles, +dry-runs against echo mocks) but does nothing at runtime** — the expression +resolves to `null` silently rather than erroring. `dry_run_workflow` catches a +null-resolved `tool_call` arg and fails with `null_resolutions`; if you see +one, check first whether the upstream node needs `.json.` inserted. + **Worked example — agent → Gmail send.** The agent must declare a schema, and -the tool_call wires each required arg from the agent BY ID: +the tool_call wires each required arg from the agent BY ID, through `.json.`: ```json { "id": "extract", "kind": "agent", "config": { @@ -195,13 +229,14 @@ the tool_call wires each required arg from the agent BY ID: "properties": { "email": {"type": "string"}, "subject": {"type": "string"}, "body": {"type": "string"} } } } } } { "id": "send", "kind": "tool_call", "config": { "slug": "GMAIL_SEND_EMAIL", "connection_ref": "composio:gmail:", - "args": { "to": "=nodes.extract.item.email", - "subject": "=nodes.extract.item.subject", - "body": "=nodes.extract.item.body" } } } + "args": { "to": "=nodes.extract.item.json.email", + "subject": "=nodes.extract.item.json.subject", + "body": "=nodes.extract.item.json.body" } } } ``` -Without the schema, `=nodes.extract.item.email` would be null (the agent's -item would be `{text: ...}`) and the send would fail preflight naming `to`. +Without the schema, `=nodes.extract.item.json.email` would be null (the +agent's `.json` has no `email` key — it's just `{text: "...", ...}`) and +`dry_run_workflow` would report it as a `null_resolutions` entry naming `to`. ### Trigger kinds — which ones actually fire diff --git a/src/openhuman/flows/builder_tools.rs b/src/openhuman/flows/builder_tools.rs index 92cdecd28..0386af050 100644 --- a/src/openhuman/flows/builder_tools.rs +++ b/src/openhuman/flows/builder_tools.rs @@ -167,6 +167,25 @@ impl Tool for ReviseWorkflowTool { } }; + // Enforcing binding-resolvability gate (see + // `ops::validate_binding_resolvability`): reject outright — rather + // than merely warn — a `tool_call` binding that is guaranteed to + // resolve null (or the wrong value) at runtime, so the builder must + // fix the graph before the revision can even be proposed. + let binding_errors = ops::validate_binding_resolvability(&graph); + if !binding_errors.is_empty() { + tracing::debug!( + target: "flows", + %name, + error_count = binding_errors.len(), + "[flows] revise_workflow: binding-resolvability check rejected the revised graph" + ); + return Ok(ToolResult::error(format!( + "{}\n\nFix these bindings and call revise_workflow again.", + binding_errors.join("\n\n") + ))); + } + let summary = super::tools::build_summary(&graph); let mut warnings = ops::graph_trigger_warnings(&graph); // Author-time wiring check: unwired REQUIRED Composio args come back @@ -718,6 +737,34 @@ impl Tool for ListAgentProfilesTool { /// so a Composio `tool_call` whose required arg is missing or `=`-resolved to /// null fails the dry run with the same actionable, field-naming error a real /// run would produce — the echo mocks alone would happily accept a null `to`. +/// +/// **Null-resolution check (the "produces functionally-broken workflows" fix):** +/// a required arg can be present *and non-Composio* (a native `oh:` tool, or a +/// Composio arg the catalog has no cached schema for) and still be wired to a +/// `=`-expression that silently resolves to `null` — the preflight above only +/// catches a *missing/null Composio-required* arg, so a graph like that used to +/// dry-run green and then do nothing at runtime. The run is driven through +/// [`tinyflows::engine::run_with_observer`] with a [`CapturingObserver`] that +/// records every node's [`ExecutionStep::diagnostics`](tinyflows::observability::ExecutionStep) +/// — the `=`-expressions the vendored engine itself traced as null-resolved +/// (see `tinyflows::expr::resolve_traced`). After the run settles, every +/// diagnostic on a **`tool_call` node's `args.*` location** is collected; any +/// hit fails the dry run with `ok: false` and the offending +/// `{ node_id, location, expression }` list, rather than reporting `ok: true` +/// for a graph that would silently no-op. Diagnostics on `agent`-node prompt +/// expressions (or any non-`tool_call` node) are NOT fatal here — a null in +/// prose is not execution-breaking the way a null tool arg is. +/// +/// **`on_error: continue`/`route` does not mask a `tool_call` failure either.** +/// Those policies convert an executor error (e.g. the required-arg preflight +/// rejecting a null arg) into a routed error ITEM so the *run* still completes +/// (`Ok(outcome)`) — the failing node's `ExecutionStep` carries an EMPTY +/// `diagnostics` (the null check above would miss it) but its `status` is +/// [`StepStatus::Error`](tinyflows::observability::StepStatus::Error). Every +/// such `tool_call` step is collected into `node_errors` +/// (`{ node_id, error }`, the error text read back out of the run's `output` +/// state — see [`tool_call_error_message`]) and fails the dry run the same as +/// a null resolution. pub struct DryRunWorkflowTool { security: Arc, config: Arc, @@ -825,14 +872,16 @@ impl Tool for DryRunWorkflowTool { } }; - // Wire the mock `AgentRunner` (echoes `agent_ref`/request/conn) so a - // draft with `agent` nodes exercises the agent-node path during the - // dry run instead of erroring on a missing capability — the plain - // `mock_capabilities()` leaves `agent: None`. No real agent turn fires; - // the mock runner is a deterministic echo, same contract as the other - // sandbox mocks. + // Wire the schema-aware mock `AgentRunner` so a draft with `agent` + // nodes exercises the agent-node path during the dry run instead of + // erroring on a missing capability — the plain `mock_capabilities()` + // leaves `agent: None`. No real agent turn fires; the mock runner is a + // deterministic echo, same contract as the other sandbox mocks, except + // it additionally honors `config.output_parser.schema` (see its doc) + // so the null-resolution check below doesn't false-positive on an + // agent node that correctly declared a schema. let mut caps = tinyflows::caps::mock::mock_capabilities_with_agent( - tinyflows::caps::mock::MockAgentRunner, + crate::openhuman::tinyflows::caps::SchemaAwareMockAgentRunner, ); // Wiring preflight over the echo mocks (see the struct doc): required // Composio args must be present and non-null even in the sandbox. @@ -840,7 +889,25 @@ impl Tool for DryRunWorkflowTool { config: self.config.clone(), inner: caps.tools.clone(), }); - let run = tinyflows::engine::run(&compiled, input, &caps); + + // Which node ids are `tool_call` nodes — the null-resolution check + // below is scoped to just these (see the struct doc: a null in an + // `agent`'s prompt is not execution-breaking the way a null tool arg + // is, so only `tool_call` diagnostics fail the dry run). + let tool_call_node_ids: std::collections::HashSet<&str> = graph + .nodes + .iter() + .filter(|node| node.kind == tinyflows::model::NodeKind::ToolCall) + .map(|node| node.id.as_str()) + .collect(); + + // Capture every node's execution diagnostics (null-resolved + // `=`-expressions the engine itself traced — see + // `tinyflows::expr::resolve_traced`) as the sandbox run executes, so + // they can be inspected once the run settles. + let observer = Arc::new(CapturingObserver::default()); + let observer_dyn: Arc = observer.clone(); + let run = tinyflows::engine::run_with_observer(&compiled, input, &caps, &observer_dyn); let outcome = match tokio::time::timeout( std::time::Duration::from_secs(DRY_RUN_TIMEOUT_SECS), run, @@ -864,23 +931,161 @@ impl Tool for DryRunWorkflowTool { } }; + // Collect every null-resolved `=`-expression that landed on a + // `tool_call` node's `args.*` config path — the class of binding + // mistake that "builds" (compiles, dry-runs against echo mocks) but + // does nothing at runtime because the wired field never had a value. + let null_resolutions: Vec = observer + .steps() + .iter() + .filter(|step| tool_call_node_ids.contains(step.node_id.as_str())) + .flat_map(|step| { + step.diagnostics.iter().filter_map(|diag| { + (diag.location == "args" || diag.location.starts_with("args.")).then(|| { + json!({ + "node_id": step.node_id, + "location": diag.location, + "expression": diag.expression, + }) + }) + }) + }) + .collect(); + + // Collect every `tool_call` node whose EXECUTOR errored (e.g. the + // Composio required-arg preflight rejecting a missing/null arg) — + // regardless of that node's `on_error` policy. A `"continue"`/`"route"` + // policy converts the failure into a routed error ITEM and the run + // still completes successfully (`Ok(outcome)`), so the naive + // `null_resolutions` check above misses it entirely: the failing + // node's `ExecutionStep` carries an EMPTY `diagnostics` (the engine + // never got far enough to trace an `=`-expression — see + // `tinyflows::engine`'s error-item path) even though the node + // genuinely failed. Only `"stop"` (the default) fails the whole run — + // and that's already caught above via `Ok(Err(e))` before this point, + // so every `StepStatus::Error` step reachable here is exactly the + // continue/route case. The error text itself isn't on the step (the + // engine only attaches it to the routed error item), so it's read + // back out of `outcome.output`. + let node_errors: Vec = observer + .steps() + .iter() + .filter(|step| { + tool_call_node_ids.contains(step.node_id.as_str()) + && matches!(step.status, tinyflows::observability::StepStatus::Error) + }) + .map(|step| { + let error = + tool_call_error_message(&outcome.output, &step.node_id).unwrap_or_else(|| { + format!( + "tool_call node '{}' failed during the sandbox run — its `on_error` \ + policy turned the failure into routed/continued data instead of \ + failing the whole dry run, but the underlying error still means the \ + node is broken.", + step.node_id + ) + }); + json!({ "node_id": step.node_id, "error": error }) + }) + .collect(); + tracing::info!( target: "flows", node_count = graph.nodes.len(), pending_approvals = outcome.pending_approvals.len(), + null_resolution_count = null_resolutions.len(), + node_error_count = node_errors.len(), "[flows] dry_run_workflow: sandbox run finished" ); + if !null_resolutions.is_empty() || !node_errors.is_empty() { + tracing::debug!( + target: "flows", + ?null_resolutions, + ?node_errors, + "[flows] dry_run_workflow: tool_call issue(s) found — failing the dry run" + ); + return Ok(ToolResult::success(serde_json::to_string_pretty(&json!({ + "sandbox": true, + "ok": false, + "null_resolutions": null_resolutions, + "node_errors": node_errors, + "message": "These tool_call args resolved to null, or a tool_call node failed \ + during the sandbox run (even one recovered via on_error: continue/route) — \ + wire null-resolved args from an upstream node's real output (give any agent \ + node an output_parser.schema so its fields are addressable), and fix or \ + rewire whatever tool_call node_errors names.", + }))?)); + } + Ok(ToolResult::success(serde_json::to_string_pretty(&json!({ "sandbox": true, "ok": true, "output": outcome.output, "pending_approvals": outcome.pending_approvals, + "null_resolutions": null_resolutions, + "node_errors": node_errors, "note": "SANDBOX (mock) output — LLM/tool/HTTP/code nodes returned deterministic echoes; NO real side effects occurred. This checks wiring/routing only, not whether real integrations work.", }))?)) } } +/// Best-effort extraction of the human-readable error message the engine +/// recorded for a `tool_call` node whose `on_error` policy is `"continue"` or +/// `"route"`. Such a node's failure is converted into an error ITEM on its +/// output (`{ "error": { "message", "node" } }` — see `tinyflows::engine`'s +/// `error_item`) rather than failing the whole run, so the message lives in +/// the run's `output` state, not on the [`tinyflows::observability::ExecutionStep`] +/// itself (whose `diagnostics` stays empty for an error step — see +/// [`DryRunWorkflowTool::execute`]'s `node_errors` collection). +fn tool_call_error_message(output: &Value, node_id: &str) -> Option { + output + .get("nodes")? + .get(node_id)? + .get("items")? + .as_array()? + .iter() + .find_map(|item| { + item.get("json")? + .get("error")? + .get("message")? + .as_str() + .map(str::to_string) + }) +} + +/// A [`tinyflows::observability::RunObserver`] that captures every finished +/// node's [`ExecutionStep`](tinyflows::observability::ExecutionStep) — in +/// particular its `diagnostics` (null-resolved `=`-expressions the engine +/// traced during that node's config resolution) — so [`DryRunWorkflowTool`] +/// can inspect them once the sandbox run settles. See the struct's "Null- +/// resolution check" doc for why this exists. +#[derive(Default)] +struct CapturingObserver { + steps: std::sync::Mutex>, +} + +impl tinyflows::observability::RunObserver for CapturingObserver { + fn on_step_finish(&self, step: &tinyflows::observability::ExecutionStep) { + self.steps + .lock() + .expect("CapturingObserver steps mutex poisoned") + .push(step.clone()); + } +} + +impl CapturingObserver { + /// A snapshot of every step recorded so far (steps are pushed + /// synchronously from `on_step_finish`, so once the run's future resolves + /// every step it will ever record is already present). + fn steps(&self) -> Vec { + self.steps + .lock() + .expect("CapturingObserver steps mutex poisoned") + .clone() + } +} + // ───────────────────────────────────────────────────────────────────────────── // save_workflow — persist a built graph onto an EXISTING saved flow // ───────────────────────────────────────────────────────────────────────────── @@ -990,6 +1195,41 @@ impl Tool for SaveWorkflowTool { .filter(|s| !s.is_empty()) .map(str::to_string); + // Same migrate/validate + enforcing binding-resolvability gate as + // propose_workflow/revise_workflow, run HERE at the tool level (not + // inside `ops::flows_update`, which the UI/RPC also call for a + // human's own edits and which must stay permissive) — so an agent + // can never persist a graph with an unresolvable `tool_call` binding + // either. See `ops::validate_binding_resolvability`. + let graph = match validate_and_migrate_graph(graph_json.clone()) { + Ok(graph) => graph, + Err(e) => { + tracing::debug!(target: "flows", %flow_id, error = %e, "[flows] save_workflow: validation failed"); + return Ok(ToolResult::error(format!( + "Workflow graph is invalid: {e}. Fix the graph and call save_workflow again." + ))); + } + }; + let binding_errors = ops::validate_binding_resolvability(&graph); + if !binding_errors.is_empty() { + tracing::debug!( + target: "flows", + %flow_id, + error_count = binding_errors.len(), + "[flows] save_workflow: binding-resolvability check rejected the graph" + ); + return Ok(ToolResult::error(format!( + "{}\n\nFix these bindings and call save_workflow again.", + binding_errors.join("\n\n") + ))); + } + // Author-time warnings (unfired trigger kinds + unwired REQUIRED + // Composio args) were previously computed by propose/revise but never + // surfaced again at save time — add them here so the agent sees any + // non-fatal wiring gaps that remain in the final persisted graph. + let mut warnings = ops::graph_trigger_warnings(&graph); + warnings.extend(ops::graph_wiring_warnings(&self.config, &graph).await); + tracing::info!( target: "flows", %flow_id, @@ -1014,6 +1254,7 @@ impl Tool for SaveWorkflowTool { "enabled": flow.enabled, "require_approval": flow.require_approval, "node_count": flow.graph.nodes.len(), + "warnings": warnings, }))?)) } Err(e) => { diff --git a/src/openhuman/flows/builder_tools_tests.rs b/src/openhuman/flows/builder_tools_tests.rs index 0178103b4..7814b9a49 100644 --- a/src/openhuman/flows/builder_tools_tests.rs +++ b/src/openhuman/flows/builder_tools_tests.rs @@ -351,6 +351,273 @@ async fn dry_run_catches_unwired_required_composio_arg() { ); } +// ── dry_run_workflow: null-resolution check ───────────────────────────────── + +#[tokio::test] +async fn dry_run_flags_tool_call_arg_null_resolved_from_unschemad_agent() { + // The `summarize` agent has no `output_parser.schema`, so (via the + // schema-aware mock agent) its structured output has no `channel` field — + // the exact "builds but does nothing" shape this check exists to catch. + let tool = DryRunWorkflowTool::new( + policy(AutonomyLevel::Supervised), + test_config(&TempDir::new().unwrap()), + ); + let graph = json!({ + "nodes": [ + { "id": "t", "kind": "trigger", "name": "Manual" }, + { "id": "summarize", "kind": "agent", "name": "Summarize", + "config": { "agent_ref": "researcher", "prompt": "summarize" } }, + { "id": "post", "kind": "tool_call", "name": "Post", + "config": { "slug": "oh:noop", + "args": { "channel": "=nodes.summarize.item.json.channel" } } } + ], + "edges": [ + { "from_node": "t", "to_node": "summarize" }, + { "from_node": "summarize", "to_node": "post" } + ] + }); + + let result = tool.execute(json!({ "graph": graph })).await.unwrap(); + assert!(!result.is_error, "{}", result.output()); + let parsed: Value = serde_json::from_str(&result.output()).unwrap(); + assert_eq!( + parsed["sandbox"], true, + "still labeled a sandbox result: {parsed}" + ); + assert_eq!( + parsed["ok"], false, + "a null-resolved tool_call arg must fail the dry run: {parsed}" + ); + let null_resolutions = parsed["null_resolutions"] + .as_array() + .expect("null_resolutions array"); + assert_eq!(null_resolutions.len(), 1, "{parsed}"); + assert_eq!(null_resolutions[0]["node_id"], "post"); + assert_eq!(null_resolutions[0]["location"], "args.channel"); + assert_eq!( + null_resolutions[0]["expression"], + "=nodes.summarize.item.json.channel" + ); + assert!( + parsed["message"] + .as_str() + .unwrap() + .to_lowercase() + .contains("output_parser"), + "{parsed}" + ); +} + +#[tokio::test] +async fn dry_run_passes_when_agent_schema_matches_tool_call_binding() { + // The FALSE-POSITIVE-PREVENTION case: `summarize` DOES declare a schema + // covering `channel`, and `post` binds exactly that field. Without the + // schema-aware mock agent (i.e. with the vendored `MockAgentRunner`, which + // always echoes `{ agent, request, connection }` regardless of schema) + // this would incorrectly fail — proving the mock is what makes the check + // accurate rather than perpetually red for correctly-built graphs. + let tool = DryRunWorkflowTool::new( + policy(AutonomyLevel::Supervised), + test_config(&TempDir::new().unwrap()), + ); + let graph = json!({ + "nodes": [ + { "id": "t", "kind": "trigger", "name": "Manual" }, + { "id": "summarize", "kind": "agent", "name": "Summarize", + "config": { "agent_ref": "researcher", "prompt": "summarize", + "output_parser": { "schema": { "type": "object", + "required": ["channel"], + "properties": { "channel": { "type": "string" } } } } } }, + { "id": "post", "kind": "tool_call", "name": "Post", + "config": { "slug": "oh:noop", + "args": { "channel": "=nodes.summarize.item.json.channel" } } } + ], + "edges": [ + { "from_node": "t", "to_node": "summarize" }, + { "from_node": "summarize", "to_node": "post" } + ] + }); + + let result = tool.execute(json!({ "graph": graph })).await.unwrap(); + assert!(!result.is_error, "{}", result.output()); + let parsed: Value = serde_json::from_str(&result.output()).unwrap(); + assert_eq!( + parsed["ok"], true, + "schema-aware mock must satisfy the declared schema: {parsed}" + ); + assert!( + parsed["null_resolutions"].as_array().unwrap().is_empty(), + "{parsed}" + ); +} + +#[tokio::test] +async fn dry_run_passes_when_tool_call_binds_to_upstream_tool_output() { + // A `tool_call` binding to another `tool_call`'s real output (not an + // agent at all) must not be affected by the agent-schema machinery above. + let tool = DryRunWorkflowTool::new( + policy(AutonomyLevel::Supervised), + test_config(&TempDir::new().unwrap()), + ); + let graph = json!({ + "nodes": [ + { "id": "t", "kind": "trigger", "name": "Manual" }, + { "id": "lookup", "kind": "tool_call", "name": "Lookup", + "config": { "slug": "oh:lookup", "args": {} } }, + { "id": "post", "kind": "tool_call", "name": "Post", + "config": { "slug": "oh:noop", + "args": { "channel": "=nodes.lookup.item.json.tool" } } } + ], + "edges": [ + { "from_node": "t", "to_node": "lookup" }, + { "from_node": "lookup", "to_node": "post" } + ] + }); + + let result = tool.execute(json!({ "graph": graph })).await.unwrap(); + assert!(!result.is_error, "{}", result.output()); + let parsed: Value = serde_json::from_str(&result.output()).unwrap(); + assert_eq!(parsed["ok"], true, "{parsed}"); + assert!( + parsed["null_resolutions"].as_array().unwrap().is_empty(), + "{parsed}" + ); +} + +#[tokio::test] +async fn dry_run_flags_tool_call_error_when_on_error_is_route() { + // `on_error: "route"` converts the preflight failure into a routed error + // ITEM so the SANDBOX RUN as a whole still completes (`Ok(outcome)`) — + // exactly the case the naive `null_resolutions`-only check would miss, + // because the failing node's diagnostics stay empty (the engine never + // got far enough to trace an `=`-expression before the preflight error). + // Seed the same schema as `dry_run_catches_unwired_required_composio_arg` + // (process-global cache; keep the arg list identical across tests). + let mut entries = std::collections::HashMap::new(); + entries.insert( + "GMAIL_SEND_EMAIL".to_string(), + vec!["to".to_string(), "body".to_string()], + ); + crate::openhuman::tinyflows::caps::seed_required_args_cache("gmail", entries); + + let tool = DryRunWorkflowTool::new( + policy(AutonomyLevel::Supervised), + test_config(&TempDir::new().unwrap()), + ); + let graph = json!({ + "nodes": [ + { "id": "t", "kind": "trigger", "name": "Manual" }, + { "id": "post", "kind": "tool_call", "name": "Send email", + "config": { "slug": "GMAIL_SEND_EMAIL", "on_error": "route", + "args": { "to": "=item.email", "body": "hello" } } } + ], + "edges": [ { "from_node": "t", "to_node": "post" } ] + }); + + // `to` misses (trigger input has no `email`) — a real run would fail the + // preflight; `on_error: "route"` must not let that slip through as `ok: true`. + let result = tool + .execute(json!({ "graph": graph, "input": {} })) + .await + .unwrap(); + assert!(!result.is_error, "{}", result.output()); + let parsed: Value = serde_json::from_str(&result.output()).unwrap(); + assert_eq!( + parsed["ok"], false, + "on_error: route must not mask a real tool_call failure: {parsed}" + ); + let node_errors = parsed["node_errors"].as_array().expect("node_errors array"); + assert_eq!(node_errors.len(), 1, "{parsed}"); + assert_eq!(node_errors[0]["node_id"], "post"); + assert!( + node_errors[0]["error"].as_str().unwrap().contains("to"), + "error must name the missing field: {parsed}" + ); +} + +#[tokio::test] +async fn dry_run_flags_tool_call_error_when_on_error_is_continue() { + // Same case as above, but `on_error: "continue"` — the other policy that + // converts a node failure into routed data instead of failing the run. + let mut entries = std::collections::HashMap::new(); + entries.insert( + "GMAIL_SEND_EMAIL".to_string(), + vec!["to".to_string(), "body".to_string()], + ); + crate::openhuman::tinyflows::caps::seed_required_args_cache("gmail", entries); + + let tool = DryRunWorkflowTool::new( + policy(AutonomyLevel::Supervised), + test_config(&TempDir::new().unwrap()), + ); + let graph = json!({ + "nodes": [ + { "id": "t", "kind": "trigger", "name": "Manual" }, + { "id": "post", "kind": "tool_call", "name": "Send email", + "config": { "slug": "GMAIL_SEND_EMAIL", "on_error": "continue", + "args": { "to": "=item.email", "body": "hello" } } } + ], + "edges": [ { "from_node": "t", "to_node": "post" } ] + }); + + let result = tool + .execute(json!({ "graph": graph, "input": {} })) + .await + .unwrap(); + assert!(!result.is_error, "{}", result.output()); + let parsed: Value = serde_json::from_str(&result.output()).unwrap(); + assert_eq!( + parsed["ok"], false, + "on_error: continue must not mask a real tool_call failure: {parsed}" + ); + assert_eq!( + parsed["node_errors"].as_array().unwrap().len(), + 1, + "{parsed}" + ); +} + +#[tokio::test] +async fn dry_run_passes_when_agent_enum_schema_binds_to_tool_call() { + // The agent declares an `enum`-constrained field; the schema-aware mock + // must synthesize an ALLOWED value (not a generic `""` placeholder, which + // would fail the vendored validator's `enum` check) so a correctly-built + // graph using an enum schema dry-runs green instead of false-positiving. + let tool = DryRunWorkflowTool::new( + policy(AutonomyLevel::Supervised), + test_config(&TempDir::new().unwrap()), + ); + let graph = json!({ + "nodes": [ + { "id": "t", "kind": "trigger", "name": "Manual" }, + { "id": "triage", "kind": "agent", "name": "Triage", + "config": { "agent_ref": "researcher", "prompt": "triage this", + "output_parser": { "schema": { "type": "object", + "required": ["priority"], + "properties": { + "priority": { "type": "string", "enum": ["urgent", "normal"] } + } } } } }, + { "id": "post", "kind": "tool_call", "name": "Post", + "config": { "slug": "oh:noop", + "args": { "priority": "=nodes.triage.item.json.priority" } } } + ], + "edges": [ + { "from_node": "t", "to_node": "triage" }, + { "from_node": "triage", "to_node": "post" } + ] + }); + + let result = tool.execute(json!({ "graph": graph })).await.unwrap(); + assert!(!result.is_error, "{}", result.output()); + let parsed: Value = serde_json::from_str(&result.output()).unwrap(); + assert_eq!( + parsed["ok"], true, + "enum-schema agent must dry-run green: {parsed}" + ); + assert!(parsed["null_resolutions"].as_array().unwrap().is_empty()); + assert!(parsed["node_errors"].as_array().unwrap().is_empty()); +} + #[tokio::test] async fn revise_workflow_warns_on_unwired_required_composio_arg() { let mut entries = std::collections::HashMap::new(); @@ -503,3 +770,95 @@ async fn save_workflow_rejects_invalid_graph_and_leaves_flow_intact() { "original graph must be untouched" ); } + +// ── save_workflow: enforcing binding-resolvability gate ───────────────────── + +/// The proven live-failure shape (same as +/// `tools_tests::propose_workflow_rejects_unschemad_agent_binding`): a +/// `summarize` agent with no `output_parser.schema`, and a `notify` tool_call +/// binding `args.channel` to its (unschemad, therefore unresolvable) output. +fn unresolvable_binding_graph() -> Value { + json!({ + "nodes": [ + { "id": "t", "kind": "trigger", "name": "Manual" }, + { "id": "summarize", "kind": "agent", "name": "Summarize", + "config": { "agent_ref": "researcher", "prompt": "summarize" } }, + { "id": "notify", "kind": "tool_call", "name": "Notify", + "config": { "slug": "SLACK_SEND_MESSAGE", + "args": { "channel": "=nodes.summarize.item.json.channel" } } } + ], + "edges": [ + { "from_node": "t", "to_node": "summarize" }, + { "from_node": "summarize", "to_node": "notify" } + ] + }) +} + +#[tokio::test] +async fn save_workflow_rejects_unschemad_agent_binding() { + let tmp = TempDir::new().unwrap(); + let config = test_config(&tmp); + let flow_id = seed_flow(&config, "Blank flow").await; + let tool = SaveWorkflowTool::new(config.clone()); + + let result = tool + .execute(json!({ "flow_id": flow_id, "graph": unresolvable_binding_graph() })) + .await + .unwrap(); + + assert!(result.is_error, "must be rejected: {}", result.output()); + let output = result.output(); + assert!(output.contains("notify"), "{output}"); + assert!(output.contains("channel"), "{output}"); + assert!(output.contains("summarize"), "{output}"); + assert!(output.contains("output_parser.schema"), "{output}"); + + // The flow it tried to save onto must be untouched. + let saved = ops::flows_get(&config, &flow_id).await.unwrap().value; + assert_eq!(saved.name, "Blank flow"); + assert_eq!( + saved.graph.nodes.len(), + 1, + "original graph must be untouched" + ); +} + +#[tokio::test] +async fn save_workflow_accepts_correctly_schemad_graph() { + let tmp = TempDir::new().unwrap(); + let config = test_config(&tmp); + let flow_id = seed_flow(&config, "Blank flow").await; + let tool = SaveWorkflowTool::new(config.clone()); + + let graph = json!({ + "nodes": [ + { "id": "t", "kind": "trigger", "name": "Manual" }, + { "id": "summarize", "kind": "agent", "name": "Summarize", + "config": { "agent_ref": "researcher", "prompt": "summarize", + "output_parser": { "schema": { "type": "object", + "required": ["channel"], + "properties": { "channel": { "type": "string" } } } } } }, + { "id": "notify", "kind": "tool_call", "name": "Notify", + "config": { "slug": "SLACK_SEND_MESSAGE", + "args": { "channel": "=nodes.summarize.item.json.channel" } } } + ], + "edges": [ + { "from_node": "t", "to_node": "summarize" }, + { "from_node": "summarize", "to_node": "notify" } + ] + }); + + let result = tool + .execute(json!({ "flow_id": flow_id, "graph": graph, "name": "Summarize and notify" })) + .await + .unwrap(); + + assert!(!result.is_error, "{}", result.output()); + let parsed: Value = serde_json::from_str(&result.output()).unwrap(); + assert_eq!(parsed["type"], "workflow_saved"); + assert_eq!(parsed["node_count"], 3); + + let saved = ops::flows_get(&config, &flow_id).await.unwrap().value; + assert_eq!(saved.name, "Summarize and notify"); + assert_eq!(saved.graph.nodes.len(), 3); +} diff --git a/src/openhuman/flows/ops.rs b/src/openhuman/flows/ops.rs index 489d8b84e..75d4882c8 100644 --- a/src/openhuman/flows/ops.rs +++ b/src/openhuman/flows/ops.rs @@ -7,7 +7,7 @@ use std::sync::Arc; use chrono::Utc; use serde_json::{json, Value}; -use tinyflows::model::{TriggerKind, WorkflowGraph}; +use tinyflows::model::{NodeKind, TriggerKind, WorkflowGraph}; use crate::openhuman::agent::turn_origin::{with_origin, AgentTurnOrigin, TrustedAutomationSource}; use crate::openhuman::config::Config; @@ -203,8 +203,9 @@ pub(crate) async fn graph_wiring_warnings(config: &Config, graph: &WorkflowGraph ); warnings.push(format!( "Node '{}': required arg `{missing}` of `{slug}` is not wired — set \ - args.{missing}, e.g. \"=nodes..item.\" (an agent feeding \ - this value needs an output schema so its fields are addressable).", + args.{missing}, e.g. \"=nodes..item.json.\" (an agent \ + feeding this value needs an output schema — `output_parser.schema` — so its \ + fields are addressable).", node.id )); } @@ -212,6 +213,193 @@ pub(crate) async fn graph_wiring_warnings(config: &Config, graph: &WorkflowGraph warnings } +// ───────────────────────────────────────────────────────────────────────────── +// Enforcing binding-resolvability gate +// ───────────────────────────────────────────────────────────────────────────── +// +// `graph_wiring_warnings` (above) is advisory — it, and `dry_run_workflow`'s +// null-resolution check (issue #4586), only WARN the author that a binding +// resolves null. Neither is consulted by the builder before it proposes or +// saves a graph, so a warned-about-but-ignored binding still ships. The +// functions below are the HARD counterpart: `validate_binding_resolvability` +// statically proves a `tool_call` node's `args` bindings are resolvable +// *before* `propose_workflow`/`revise_workflow`/`save_workflow` accept the +// graph at all (see their call sites), so the LLM builder is forced to fix +// the wiring rather than merely being told about it. + +/// Node kinds whose real capability adapter wraps its structured output in +/// the stable `{ json, text, raw }` envelope (`src/openhuman/tinyflows/caps.rs`): +/// a binding into one of these must dereference `.item.json.`, never +/// `.item.` directly — the latter reads the envelope wrapper itself +/// (an object with `json`/`text`/`raw` keys), not the field inside it, and +/// resolves `null` at runtime. Every other node kind (`code`, `transform`, +/// `split_out`, `merge`, `output_parser`, `sub_workflow`, `trigger`, +/// `condition`, `switch`) emits its item directly with no envelope, so no +/// convention applies to a binding that targets one of them. +const ENVELOPING_KINDS: &[NodeKind] = &[NodeKind::Agent, NodeKind::ToolCall, NodeKind::HttpRequest]; + +/// Recursively collects every `=`-prefixed expression leaf in a config +/// `Value` tree, paired with its dotted location (array elements as numeric +/// segments, e.g. `"args.cc.0"`) — the same location convention as +/// `tinyflows::expr::resolve_traced`. Unlike that function this never +/// evaluates an expression against a scope; it only locates the leaves so +/// [`validate_binding_resolvability`] can statically pattern-match them. +fn collect_expressions(value: &Value) -> Vec<(String, String)> { + fn walk(value: &Value, location: &str, out: &mut Vec<(String, String)>) { + match value { + Value::Object(map) => { + for (k, v) in map { + let child = if location.is_empty() { + k.clone() + } else { + format!("{location}.{k}") + }; + walk(v, &child, out); + } + } + Value::Array(items) => { + for (i, v) in items.iter().enumerate() { + let child = if location.is_empty() { + i.to_string() + } else { + format!("{location}.{i}") + }; + walk(v, &child, out); + } + } + Value::String(s) if tinyflows::expr::is_expression(s) => { + out.push((location.to_string(), s.clone())); + } + _ => {} + } + } + let mut out = Vec::new(); + walk(value, "", &mut out); + out +} + +/// Matches the dotted-path form of a node-output binding — +/// `=nodes..item[.json].` — returning `(ref_id, has_json, +/// field)`. `has_json` is `true` when the expression dereferenced the +/// `{json,text,raw}` envelope wrapper (`.item.json.`) rather than the +/// item directly (`.item.`). +/// +/// Only the dotted-path form is recognized here — the equivalent jq form +/// (e.g. `=.nodes["ref"].items[0].field`) is an arbitrary jq program, not a +/// fixed grammar, so it is not statically pattern-matched; that form is still +/// covered dynamically by `dry_run_workflow`'s null-resolution check (#4586), +/// which actually evaluates the expression at run time. +fn parse_node_binding(expr: &str) -> Option<(String, bool, String)> { + fn node_binding_regex() -> &'static regex::Regex { + static RE: std::sync::OnceLock = std::sync::OnceLock::new(); + RE.get_or_init(|| { + regex::Regex::new( + r"^=nodes\.([A-Za-z_][A-Za-z0-9_]*)\.item(?:\.(json))?\.([A-Za-z_][A-Za-z0-9_]*)", + ) + .expect("static regex is valid") + }) + } + let caps = node_binding_regex().captures(expr)?; + let ref_id = caps.get(1)?.as_str().to_string(); + let has_json = caps.get(2).is_some(); + let field = caps.get(3)?.as_str().to_string(); + Some((ref_id, has_json, field)) +} + +/// Human-readable label for a [`NodeKind`], for +/// [`validate_binding_resolvability`]'s envelope-violation message. +fn node_kind_label(kind: &NodeKind) -> &'static str { + match kind { + NodeKind::Agent => "an agent", + NodeKind::ToolCall => "a tool_call", + NodeKind::HttpRequest => "an http_request", + _ => "a node", + } +} + +/// Statically proves every `tool_call` node's `config.args` bindings are +/// resolvable, rejecting the graph (a non-empty `Vec` = reject; empty = +/// pass) when one is GUARANTEED to resolve `null` (or the wrong value) at +/// runtime. See the [module section](self) header for why this exists +/// alongside the advisory `graph_wiring_warnings`/`dry_run_workflow` checks. +/// +/// Scoped to `tool_call` `args` only — deliberately NOT `agent` prompts: a +/// prose string that merely mentions a bad `nodes.` path degrades output +/// quality but does not break execution the way a `null` tool argument does, +/// so enforcing there would false-positive on legitimate free-text prompts. +/// +/// For every `=nodes..item[.json].` binding found in a +/// `tool_call`'s `args` (via [`collect_expressions`] + [`parse_node_binding`]): +/// - a `` that doesn't resolve to a node in the graph is skipped — a +/// dangling reference is already a `tinyflows::validate::validate` +/// structural error, caught upstream of this pass. +/// - a `` that IS an [`ENVELOPING_KINDS`] node and the expression used +/// `.item.` (no `.json`) is REJECTED: it dereferences the envelope +/// wrapper, not the field inside it. +/// - a `` that is an `agent` node is REJECTED unless it declares +/// `config.output_parser.schema` with an object `properties` map +/// containing `` — the exact shape a real run's output-parser +/// sub-port enforces; without it the agent's structured output has no +/// addressable ``. +/// - a `` that is `tool_call`/`http_request` only gets the envelope +/// check above — neither has a static output schema to check field +/// membership against ahead of a real run. +/// - any other referenced kind (`code`, `transform`, `split_out`, `merge`, +/// `output_parser`, `sub_workflow`, `trigger`, `condition`, `switch`) has no +/// schema or envelope convention to enforce and is accepted. +pub(crate) fn validate_binding_resolvability(graph: &WorkflowGraph) -> Vec { + let mut errors = Vec::new(); + for node in &graph.nodes { + if node.kind != NodeKind::ToolCall { + continue; + } + let Some(args) = node.config.get("args") else { + continue; + }; + for (location, expr) in collect_expressions(args) { + let Some((ref_id, has_json, field)) = parse_node_binding(&expr) else { + continue; + }; + let Some(ref_node) = graph.node(&ref_id) else { + continue; + }; + + if ENVELOPING_KINDS.contains(&ref_node.kind) && !has_json { + errors.push(format!( + "Node '{}': arg `{location}` (`{expr}`) uses `.item.{field}` on {} node \ + `{ref_id}`, but agent/tool_call/http_request nodes wrap output in {{json, \ + text, raw}} — use `=nodes.{ref_id}.item.json.{field}` instead.", + node.id, + node_kind_label(&ref_node.kind), + )); + continue; + } + + if ref_node.kind == NodeKind::Agent { + let has_field = ref_node + .config + .get("output_parser") + .and_then(|p| p.get("schema")) + .filter(|s| !s.is_null()) + .and_then(|s| s.get("properties")) + .and_then(Value::as_object) + .is_some_and(|props| props.contains_key(&field)); + if !has_field { + errors.push(format!( + "Node '{}': arg `{location}` (`{expr}`) binds to agent node `{ref_id}`, \ + which has no `output_parser.schema` declaring `{field}` — its \ + structured output has no addressable `{field}`, so this binding \ + resolves null at runtime. Fix: add `{field}` to node `{ref_id}`'s \ + output_parser.schema and bind via `=nodes.{ref_id}.item.json.{field}`.", + node.id + )); + } + } + } + } + errors +} + /// Validates a candidate graph without persisting it — the same /// migrate/validate path `flows_create` and `ProposeWorkflowTool` use — and /// reports structural errors alongside non-fatal trigger warnings diff --git a/src/openhuman/flows/ops_tests.rs b/src/openhuman/flows/ops_tests.rs index 2d54dd122..82e050130 100644 --- a/src/openhuman/flows/ops_tests.rs +++ b/src/openhuman/flows/ops_tests.rs @@ -1739,6 +1739,90 @@ fn flow_stream_target_generates_request_id_when_absent_or_blank() { assert_ne!(a.request_id, b.request_id); } +// ── validate_binding_resolvability ────────────────────────────────────────── + +/// Runs a candidate graph `Value` through the exact same migrate/validate +/// path the builder tools use, for a [`WorkflowGraph`] test fixture. +fn graph(value: Value) -> WorkflowGraph { + validate_and_migrate_graph(value).expect("structurally valid test graph") +} + +#[test] +fn binding_to_agent_without_schema_is_rejected() { + // The exact live-failure shape: `summarize` has no `output_parser.schema` + // at all, so its structured output has no addressable `channel` field. + let g = graph(json!({ + "nodes": [ + { "id": "t", "kind": "trigger", "name": "Manual" }, + { "id": "summarize", "kind": "agent", "name": "Summarize", + "config": { "agent_ref": "researcher", "prompt": "summarize" } }, + { "id": "post", "kind": "tool_call", "name": "Post", + "config": { "slug": "SLACK_SEND_MESSAGE", + "args": { "channel": "=nodes.summarize.item.json.channel" } } } + ], + "edges": [ + { "from_node": "t", "to_node": "summarize" }, + { "from_node": "summarize", "to_node": "post" } + ] + })); + let errors = validate_binding_resolvability(&g); + assert_eq!(errors.len(), 1, "{errors:?}"); + assert!(errors[0].contains("post"), "{}", errors[0]); + assert!(errors[0].contains("channel"), "{}", errors[0]); + assert!(errors[0].contains("summarize"), "{}", errors[0]); + assert!(errors[0].contains("output_parser.schema"), "{}", errors[0]); +} + +#[test] +fn binding_to_agent_with_schema_missing_field_is_rejected() { + // A schema IS declared, but it doesn't cover the field the binding reads. + let g = graph(json!({ + "nodes": [ + { "id": "t", "kind": "trigger", "name": "Manual" }, + { "id": "summarize", "kind": "agent", "name": "Summarize", + "config": { "prompt": "summarize", + "output_parser": { "schema": { "type": "object", + "properties": { "summary": { "type": "string" } } } } } }, + { "id": "post", "kind": "tool_call", "name": "Post", + "config": { "slug": "SLACK_SEND_MESSAGE", + "args": { "channel": "=nodes.summarize.item.json.channel" } } } + ], + "edges": [ + { "from_node": "t", "to_node": "summarize" }, + { "from_node": "summarize", "to_node": "post" } + ] + })); + let errors = validate_binding_resolvability(&g); + assert_eq!(errors.len(), 1, "{errors:?}"); + assert!(errors[0].contains("channel"), "{}", errors[0]); +} + +#[test] +fn binding_to_agent_with_matching_schema_is_accepted() { + let g = graph(json!({ + "nodes": [ + { "id": "t", "kind": "trigger", "name": "Manual" }, + { "id": "summarize", "kind": "agent", "name": "Summarize", + "config": { "prompt": "summarize", + "output_parser": { "schema": { "type": "object", + "required": ["channel"], + "properties": { "channel": { "type": "string" } } } } } }, + { "id": "post", "kind": "tool_call", "name": "Post", + "config": { "slug": "SLACK_SEND_MESSAGE", + "args": { "channel": "=nodes.summarize.item.json.channel" } } } + ], + "edges": [ + { "from_node": "t", "to_node": "summarize" }, + { "from_node": "summarize", "to_node": "post" } + ] + })); + assert!( + validate_binding_resolvability(&g).is_empty(), + "{:?}", + validate_binding_resolvability(&g) + ); +} + // ───────────────────────────────────────────────────────────────────────────── // degrade_completed_status (PR2 — run honesty) // ───────────────────────────────────────────────────────────────────────────── @@ -1817,6 +1901,88 @@ fn failed_step_error_summary_names_every_errored_node() { ); } +#[test] +fn envelope_violation_detected() { + // `summarize` DOES declare a matching schema, but the binding reaches + // into `.item.channel` (skipping `.json`) — that dereferences the + // `{json,text,raw}` envelope wrapper itself, not the field inside it. + let g = graph(json!({ + "nodes": [ + { "id": "t", "kind": "trigger", "name": "Manual" }, + { "id": "summarize", "kind": "agent", "name": "Summarize", + "config": { "prompt": "summarize", + "output_parser": { "schema": { "type": "object", + "properties": { "channel": { "type": "string" } } } } } }, + { "id": "post", "kind": "tool_call", "name": "Post", + "config": { "slug": "SLACK_SEND_MESSAGE", + "args": { "channel": "=nodes.summarize.item.channel" } } } + ], + "edges": [ + { "from_node": "t", "to_node": "summarize" }, + { "from_node": "summarize", "to_node": "post" } + ] + })); + let errors = validate_binding_resolvability(&g); + assert_eq!(errors.len(), 1, "{errors:?}"); + assert!(errors[0].contains("json"), "{}", errors[0]); + assert!(errors[0].contains("summarize"), "{}", errors[0]); +} + +#[test] +fn non_enveloping_node_binding_is_accepted() { + // `code` nodes emit their item directly (no envelope) — `.item.` + // is the correct, and only, form. + let g = graph(json!({ + "nodes": [ + { "id": "t", "kind": "trigger", "name": "Manual" }, + { "id": "compute", "kind": "code", "name": "Compute", + "config": { "language": "javascript", "source": "return {channel:'general'};" } }, + { "id": "post", "kind": "tool_call", "name": "Post", + "config": { "slug": "SLACK_SEND_MESSAGE", + "args": { "channel": "=nodes.compute.item.channel" } } } + ], + "edges": [ + { "from_node": "t", "to_node": "compute" }, + { "from_node": "compute", "to_node": "post" } + ] + })); + assert!( + validate_binding_resolvability(&g).is_empty(), + "{:?}", + validate_binding_resolvability(&g) + ); +} + +#[test] +fn literal_args_unaffected() { + let g = graph(json!({ + "nodes": [ + { "id": "t", "kind": "trigger", "name": "Manual" }, + { "id": "post", "kind": "tool_call", "name": "Post", + "config": { "slug": "SLACK_SEND_MESSAGE", + "args": { "channel": "general", "count": 3, "cc": ["a@b.com"] } } } + ], + "edges": [ { "from_node": "t", "to_node": "post" } ] + })); + assert!(validate_binding_resolvability(&g).is_empty()); +} + +#[test] +fn agent_prompt_binding_unaffected() { + // The check is scoped to `tool_call` `args` only — an agent's own + // config (its prompt) is never inspected, even if it references a + // dangling/unschemad node path. + let g = graph(json!({ + "nodes": [ + { "id": "t", "kind": "trigger", "name": "Manual" }, + { "id": "summarize", "kind": "agent", "name": "Summarize", + "config": { "prompt": "=nodes.missing.item.channel" } } + ], + "edges": [ { "from_node": "t", "to_node": "summarize" } ] + })); + assert!(validate_binding_resolvability(&g).is_empty()); +} + #[test] fn finalize_terminal_status_pending_approval_wins_over_error() { // Precedence: an outstanding pending_approval always wins, even if a step diff --git a/src/openhuman/flows/tools.rs b/src/openhuman/flows/tools.rs index ee7fcd5d6..0a89470e6 100644 --- a/src/openhuman/flows/tools.rs +++ b/src/openhuman/flows/tools.rs @@ -23,7 +23,7 @@ use serde_json::{json, Value}; use tinyflows::model::{Node, NodeKind, WorkflowGraph}; use crate::openhuman::config::Config; -use crate::openhuman::flows::ops::validate_and_migrate_graph; +use crate::openhuman::flows::ops::{validate_and_migrate_graph, validate_binding_resolvability}; use crate::openhuman::tools::traits::{PermissionLevel, Tool, ToolResult}; /// Max characters kept for a `config_hint` before truncation, so a long @@ -176,6 +176,25 @@ impl Tool for ProposeWorkflowTool { } }; + // Enforcing binding-resolvability gate (see + // `ops::validate_binding_resolvability`): reject outright — rather + // than merely warn — a `tool_call` binding that is guaranteed to + // resolve null (or the wrong value) at runtime, so the builder must + // fix the graph before it can even be proposed to the user. + let binding_errors = validate_binding_resolvability(&graph); + if !binding_errors.is_empty() { + tracing::debug!( + target: "flows", + %name, + error_count = binding_errors.len(), + "[flows] propose_workflow: binding-resolvability check rejected the graph" + ); + return Ok(ToolResult::error(format!( + "{}\n\nFix these bindings and call propose_workflow again.", + binding_errors.join("\n\n") + ))); + } + let summary = build_summary(&graph); // Author-time warnings: unfired trigger kinds + unwired REQUIRED // Composio args (see `ops::graph_wiring_warnings`) — surfaced on the diff --git a/src/openhuman/flows/tools_tests.rs b/src/openhuman/flows/tools_tests.rs index 650f47bda..88c0e6c78 100644 --- a/src/openhuman/flows/tools_tests.rs +++ b/src/openhuman/flows/tools_tests.rs @@ -262,3 +262,46 @@ fn display_label_humanizes_the_tool_name() { Some("Propose Workflow") ); } + +// ── enforcing binding-resolvability gate ──────────────────────────────────── + +#[tokio::test] +async fn propose_workflow_rejects_unschemad_agent_binding() { + // The proven live-failure graph: `summarize` has no `output_parser.schema`, + // so `post`'s `args.channel` binding is guaranteed to resolve null at + // runtime. Unlike the advisory dry-run check, propose_workflow must + // REJECT this outright rather than warn (warning_count would have been 0 + // here — nothing stopped this from reaching save_workflow before). + let tmp = TempDir::new().unwrap(); + let tool = ProposeWorkflowTool::new(test_config(&tmp)); + + let graph = json!({ + "nodes": [ + { "id": "t", "kind": "trigger", "name": "Manual" }, + { "id": "summarize", "kind": "agent", "name": "Summarize", + "config": { "agent_ref": "researcher", "prompt": "summarize" } }, + { "id": "notify", "kind": "tool_call", "name": "Notify", + "config": { "slug": "SLACK_SEND_MESSAGE", + "args": { "channel": "=nodes.summarize.item.json.channel" } } } + ], + "edges": [ + { "from_node": "t", "to_node": "summarize" }, + { "from_node": "summarize", "to_node": "notify" } + ] + }); + + let result = tool + .execute(json!({ "name": "Summarize and notify", "graph": graph })) + .await + .unwrap(); + + assert!(result.is_error, "must be rejected: {}", result.output()); + let output = result.output(); + assert!(output.contains("notify"), "{output}"); + assert!(output.contains("channel"), "{output}"); + assert!(output.contains("summarize"), "{output}"); + assert!( + output.contains("output_parser.schema"), + "must name the missing schema as the fix: {output}" + ); +} diff --git a/src/openhuman/tinyflows/caps.rs b/src/openhuman/tinyflows/caps.rs index 088766145..f7f10167b 100644 --- a/src/openhuman/tinyflows/caps.rs +++ b/src/openhuman/tinyflows/caps.rs @@ -836,6 +836,122 @@ fn prepend_system_message(request: &mut Value, system_prompt: &str) { } } +/// A **dry-run-only** [`AgentRunner`] mock that, unlike the vendored crate's +/// `tinyflows::caps::mock::MockAgentRunner`, respects an `agent` node's +/// `config.output_parser.schema` when synthesizing its echo response. +/// +/// `DryRunWorkflowTool` (`flows::builder_tools`) wires this in place of the +/// vendored `MockAgentRunner` so its null-resolution check (every `tool_call` +/// arg that resolves to `null`) doesn't **false-positive** on a CORRECTLY-built +/// agent node. Without it: the vendored `MockAgentRunner` always echoes +/// `{ agent, request, connection }` regardless of schema, and the vendored +/// `agent` node's output-parser sub-port (`tinyflows::nodes::integration::schema`) +/// then fails that shape against ANY declared schema (no field matches) and +/// falls to a one-shot LLM auto-fix that the sandbox's plain `MockLlm` also +/// can't satisfy — so the whole dry run would error out even for a workflow a +/// real run (via [`OpenHumanAgentRunner`], whose completion the same sub-port +/// validates/repairs against the schema) would execute cleanly. +/// +/// When `request` (the resolved node config `run_agent` receives — see +/// [`AgentRunner::run_agent`]) carries a non-null `output_parser.schema` +/// describing an object with `properties`, returns an object with every +/// declared property present, populated with a type-appropriate placeholder +/// (`string` → `""`, `number`/`integer` → `0`, `boolean` → `false`, `object` → +/// `{}`, `array` → `[]`, anything else → `null`; a property with a non-empty +/// `enum` gets its FIRST allowed value instead — see [`placeholder_for_type`]) +/// — enough to satisfy the vendored validator's `type`/`required`/`enum` +/// checks (see `tinyflows::nodes::integration::schema::validate`) without a +/// real model call. With no schema, mirrors the vendored `MockAgentRunner`'s +/// default echo shape so dry-run behavior for schema-less agents is unchanged. +#[derive(Debug, Default, Clone)] +pub struct SchemaAwareMockAgentRunner; + +#[async_trait] +impl AgentRunner for SchemaAwareMockAgentRunner { + async fn run_agent( + &self, + agent_ref: &str, + request: Value, + conn: Option<&str>, + ) -> Result { + let schema = request + .get("output_parser") + .and_then(|parser| parser.get("schema")) + .filter(|schema| !schema.is_null()); + match schema { + Some(schema) => { + let placeholder = placeholder_for_schema(schema); + tracing::debug!( + target: "flows", + agent_ref, + "[flows] dry_run: schema-aware mock agent synthesized a placeholder \ + matching output_parser.schema" + ); + Ok(placeholder) + } + None => { + tracing::debug!( + target: "flows", + agent_ref, + "[flows] dry_run: schema-aware mock agent has no output_parser.schema — \ + mirroring the vendored MockAgentRunner echo shape" + ); + Ok(json!({ "agent": agent_ref, "request": request, "connection": conn })) + } + } + } +} + +/// Builds a placeholder JSON value satisfying `schema`'s `properties`/`type` +/// constraints, for [`SchemaAwareMockAgentRunner`]. Only the shallow, top-level +/// `properties` map is populated — enough for the minimal validator in +/// `tinyflows::nodes::integration::schema` (`type`, `required`, `properties`); +/// deeply-nested `required` constraints on a nested `object`/`array` property +/// are a documented limitation (the placeholder for those is an empty `{}`/`[]`). +fn placeholder_for_schema(schema: &Value) -> Value { + match schema.get("properties").and_then(Value::as_object) { + Some(props) => { + let placeholders = props + .iter() + .map(|(key, subschema)| (key.clone(), placeholder_for_type(subschema))); + Value::Object(placeholders.collect()) + } + // No `properties` to enumerate (e.g. a bare `{"type": "array"}` + // schema) — fall back to a type-only placeholder for the schema itself. + None => placeholder_for_type(schema), + } +} + +/// The placeholder value for one property's subschema, keyed by its +/// declared JSON-Schema `type` (see [`placeholder_for_schema`]). +/// +/// An `enum` constraint is honored FIRST, before falling back to the +/// type-only placeholder: the vendored validator +/// (`tinyflows::nodes::integration::schema::validate`) rejects any value not +/// listed in a schema's `enum`, and a generic type placeholder (e.g. `""` for +/// `{"type": "string", "enum": ["urgent", "normal"]}`) is essentially never +/// one of the allowed values — that would fail the dry run even though a real +/// agent, prompted with the schema, could easily satisfy it. The schema +/// author's own first listed value is always allowed by construction, so it's +/// returned as-is (whatever its JSON type). +fn placeholder_for_type(subschema: &Value) -> Value { + if let Some(first_allowed) = subschema + .get("enum") + .and_then(Value::as_array) + .and_then(|values| values.first()) + { + return first_allowed.clone(); + } + match subschema.get("type").and_then(Value::as_str) { + Some("string") => json!(""), + Some("number" | "integer") => json!(0), + Some("boolean") => json!(false), + Some("object") => json!({}), + Some("array") => json!([]), + _ => Value::Null, + } +} + /// Parses a `"composio::"` `connection_ref` (see the /// node catalog, `my_docs/ohxtf/commons/12-node-catalog-0.2.md`) and returns /// the trailing connection id segment. Values that don't match this shape @@ -1220,8 +1336,10 @@ pub(crate) async fn preflight_composio_args( let first = &missing[0]; Err(EngineError::Capability(format!( "tool_call `{slug}`: required arg(s) {list} missing or resolved to null — wire each from \ - an upstream node's output, e.g. \"{first}\": \"=nodes..item.\". If the \ - value comes from an agent node, give that agent an output schema \ + an upstream node's output, e.g. \"{first}\": \"=nodes..item.json.\" \ + (drop `.json` only if `` is a code/transform/split_out/merge/trigger node — \ + `agent`/`tool_call`/`http_request` nodes wrap their output in a `{{json,text,raw}}` \ + envelope). If the value comes from an agent node, give that agent an output schema \ (config.output_parser.schema) so its fields are addressable." ))) } @@ -2060,6 +2178,150 @@ mod tests { assert_eq!(req, json!("just a string")); } + // ── SchemaAwareMockAgentRunner ─────────────────────────────────────────── + + #[tokio::test] + async fn schema_aware_mock_agent_mirrors_vendored_echo_without_a_schema() { + // No `output_parser.schema` on the request: identical shape to the + // vendored `MockAgentRunner` so schema-less dry runs are unaffected. + let runner = SchemaAwareMockAgentRunner; + let request = json!({ "prompt": "hi" }); + let out = runner + .run_agent("researcher", request.clone(), Some("conn_1")) + .await + .expect("run_agent"); + assert_eq!(out["agent"], "researcher"); + assert_eq!(out["request"], request); + assert_eq!(out["connection"], "conn_1"); + } + + #[tokio::test] + async fn schema_aware_mock_agent_populates_declared_properties() { + let runner = SchemaAwareMockAgentRunner; + let request = json!({ + "prompt": "extract", + "output_parser": { "schema": { "type": "object", + "required": ["email", "count", "active", "meta", "tags"], + "properties": { + "email": { "type": "string" }, + "count": { "type": "integer" }, + "active": { "type": "boolean" }, + "meta": { "type": "object" }, + "tags": { "type": "array" } + } } } + }); + let out = runner + .run_agent("researcher", request, None) + .await + .expect("run_agent"); + assert_eq!(out["email"], ""); + assert_eq!(out["count"], 0); + assert_eq!(out["active"], false); + assert_eq!(out["meta"], json!({})); + assert_eq!(out["tags"], json!([])); + } + + #[tokio::test] + async fn schema_aware_mock_agent_populates_an_enum_property_with_an_allowed_value() { + // A generic string placeholder (`""`) would fail the vendored + // validator's `enum` check even though a real agent could easily + // satisfy it — the mock must pick one of the schema's own allowed + // values (see `placeholder_for_type`'s enum handling). + let runner = SchemaAwareMockAgentRunner; + let request = json!({ + "prompt": "triage", + "output_parser": { "schema": { "type": "object", + "required": ["priority"], + "properties": { + "priority": { "type": "string", "enum": ["urgent", "normal"] } + } } } + }); + let out = runner + .run_agent("researcher", request, None) + .await + .expect("run_agent"); + let allowed = ["urgent", "normal"]; + assert!( + allowed.contains(&out["priority"].as_str().unwrap()), + "expected an allowed enum value, got: {out}" + ); + } + + #[tokio::test] + async fn schema_aware_mock_agent_ignores_null_schema() { + // `output_parser: { schema: null }` (or no `output_parser` at all) is + // treated identically to "no schema" — the vendored echo shape. + let runner = SchemaAwareMockAgentRunner; + let request = json!({ "prompt": "hi", "output_parser": { "schema": null } }); + let out = runner + .run_agent("researcher", request.clone(), None) + .await + .expect("run_agent"); + assert_eq!(out["agent"], "researcher"); + assert_eq!(out["request"], request); + } + + #[test] + fn placeholder_for_schema_falls_back_to_type_without_properties() { + assert_eq!( + placeholder_for_schema(&json!({ "type": "array" })), + json!([]) + ); + assert_eq!( + placeholder_for_schema(&json!({ "type": "string" })), + json!("") + ); + } + + #[test] + fn placeholder_for_type_covers_every_json_schema_type() { + assert_eq!( + placeholder_for_type(&json!({ "type": "string" })), + json!("") + ); + assert_eq!(placeholder_for_type(&json!({ "type": "number" })), json!(0)); + assert_eq!( + placeholder_for_type(&json!({ "type": "integer" })), + json!(0) + ); + assert_eq!( + placeholder_for_type(&json!({ "type": "boolean" })), + json!(false) + ); + assert_eq!( + placeholder_for_type(&json!({ "type": "object" })), + json!({}) + ); + assert_eq!(placeholder_for_type(&json!({ "type": "array" })), json!([])); + assert_eq!(placeholder_for_type(&json!({})), Value::Null); + } + + #[test] + fn placeholder_for_type_prefers_the_first_enum_value_over_the_generic_type() { + // A generic type placeholder (`""`) is essentially never one of an + // enum's allowed values, so it must never be used when `enum` is set. + assert_eq!( + placeholder_for_type(&json!({ "type": "string", "enum": ["urgent", "normal"] })), + json!("urgent") + ); + // The first enum value wins even when its JSON type doesn't match + // `type` (schema authors sometimes skip `type` entirely with `enum`). + assert_eq!( + placeholder_for_type(&json!({ "enum": [1, 2, 3] })), + json!(1) + ); + } + + #[test] + fn placeholder_for_type_ignores_an_empty_enum() { + // An empty `enum` array has no first value to prefer — fall back to + // the type-only placeholder rather than panicking or returning null. + assert_eq!( + placeholder_for_type(&json!({ "type": "string", "enum": [] })), + json!("") + ); + } + fn integration( toolkit: &str, connected: bool, diff --git a/src/openhuman/tinyflows/tests.rs b/src/openhuman/tinyflows/tests.rs index 44eb35f62..72ce1a92e 100644 --- a/src/openhuman/tinyflows/tests.rs +++ b/src/openhuman/tinyflows/tests.rs @@ -420,7 +420,7 @@ async fn preflight_fails_before_dispatch_naming_the_missing_field() { let msg = err.to_string(); assert!(msg.contains("`to`"), "error must name the field: {msg}"); assert!( - msg.contains("=nodes..item."), + msg.contains("=nodes..item.json."), "error must suggest the wiring fix: {msg}" ); assert!(