diff --git a/Cargo.lock b/Cargo.lock index e58e0f4f6..e643ea41d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -6808,9 +6808,10 @@ dependencies = [ [[package]] name = "tinyagents" -version = "1.5.0" +version = "1.7.1" dependencies = [ "async-trait", + "bytes", "futures", "reqwest 0.12.28", "rhai", @@ -6845,7 +6846,7 @@ dependencies = [ [[package]] name = "tinyflows" -version = "0.3.0" +version = "0.3.2" dependencies = [ "async-trait", "futures-timer", @@ -8219,7 +8220,7 @@ version = "0.1.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22" dependencies = [ - "windows-sys 0.61.2", + "windows-sys 0.48.0", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml index b9c72ea41..8dc2b4fc5 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -42,7 +42,7 @@ crate-type = ["rlib"] tinyplace = "1.0.1" # tinyflows — host-agnostic workflow engine (typed node graph → validate → compile → # run on tinyagents). Powers the "Workflows" feature via the seam in -# `src/openhuman/tinyflows/` + the `flows::` domain. Pulls tinyagents 1.5 transitively +# `src/openhuman/tinyflows/` + the `flows::` domain. Pulls tinyagents 1.7 transitively # (same version openhuman already uses — no conflict). Published on crates.io # and patched below to the vendored submodule. tinyflows = "0.3" @@ -63,7 +63,7 @@ tinyjuice = { version = "0.1", default-features = false } # migration re-points those rows to the crate checkpointer. # The `repl` feature adds the embedded Rhai `.ragsh` session runtime powering # the `rlm` language-workflow tool (`src/openhuman/rlm/`). -tinyagents = { version = "1.5.0", features = ["sqlite", "repl"] } +tinyagents = { version = "1.7", features = ["sqlite", "repl"] } # TinyCortex — Rust core for the memory engine (store/chunks/tree/retrieval/ # queue/ingest/score + long tail), vendored as a git submodule and patched # below to `vendor/tinycortex`. OpenHuman's memory subsystem migrates onto this diff --git a/app/src-tauri/Cargo.lock b/app/src-tauri/Cargo.lock index d018a2ef4..19166d8c6 100644 --- a/app/src-tauri/Cargo.lock +++ b/app/src-tauri/Cargo.lock @@ -115,6 +115,20 @@ dependencies = [ "subtle", ] +[[package]] +name = "ahash" +version = "0.8.12" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5a15f179cd60c4584b8a8c596927aadc462e27f2ca70c04e0071964a73ba7a75" +dependencies = [ + "cfg-if", + "const-random", + "getrandom 0.3.4", + "once_cell", + "version_check", + "zerocopy", +] + [[package]] name = "aho-corasick" version = "1.1.4" @@ -236,7 +250,7 @@ dependencies = [ "objc2-foundation 0.3.2", "parking_lot", "percent-encoding", - "windows-sys 0.60.2", + "windows-sys 0.59.0", "x11rb", ] @@ -2079,7 +2093,7 @@ dependencies = [ "libc", "option-ext", "redox_users 0.5.2", - "windows-sys 0.61.2", + "windows-sys 0.59.0", ] [[package]] @@ -2482,7 +2496,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb" dependencies = [ "libc", - "windows-sys 0.61.2", + "windows-sys 0.59.0", ] [[package]] @@ -4122,7 +4136,7 @@ dependencies = [ "portable-atomic", "portable-atomic-util", "serde_core", - "windows-sys 0.61.2", + "windows-sys 0.59.0", ] [[package]] @@ -4931,6 +4945,15 @@ dependencies = [ "libc", ] +[[package]] +name = "no-std-compat" +version = "0.4.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b93853da6d84c2e3c7d730d6473e8817692dd89be387eb01b94d7f108ecb5b8c" +dependencies = [ + "spin", +] + [[package]] name = "nodrop" version = "0.1.14" @@ -5007,7 +5030,7 @@ version = "0.50.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7957b9740744892f114936ab4a57b3f487491bbeafaf8083688b16841a4240e5" dependencies = [ - "windows-sys 0.61.2", + "windows-sys 0.59.0", ] [[package]] @@ -5492,6 +5515,9 @@ name = "once_cell" version = "1.21.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9f7c3e4beb33f85d45ae3e3a1792185706c8e16d043238c593331cc7cd313b50" +dependencies = [ + "portable-atomic", +] [[package]] name = "once_cell_polyfill" @@ -5636,6 +5662,7 @@ dependencies = [ "tempfile", "thiserror 2.0.18", "tinyagents", + "tinycortex", "tinyflows", "tinyjuice", "tinyplace", @@ -6698,7 +6725,7 @@ dependencies = [ "once_cell", "socket2", "tracing", - "windows-sys 0.60.2", + "windows-sys 0.59.0", ] [[package]] @@ -7133,6 +7160,35 @@ dependencies = [ "bytemuck", ] +[[package]] +name = "rhai" +version = "1.25.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "dd4dd0f8c36625202a4ba553c416c19b719947cd2a31d1bda06126e4a5727daf" +dependencies = [ + "ahash", + "bitflags 2.11.1", + "no-std-compat", + "num-traits", + "once_cell", + "rhai_codegen", + "smallvec", + "smartstring", + "thin-vec", + "web-time", +] + +[[package]] +name = "rhai_codegen" +version = "3.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3cd3a7535e50bf36857e7be7bec276d334e8c2dfa469c2201226fd01638ea5ca" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.117", +] + [[package]] name = "ring" version = "0.17.14" @@ -7256,7 +7312,7 @@ dependencies = [ "errno", "libc", "linux-raw-sys", - "windows-sys 0.61.2", + "windows-sys 0.59.0", ] [[package]] @@ -7315,7 +7371,7 @@ dependencies = [ "security-framework 3.7.0", "security-framework-sys", "webpki-root-certs", - "windows-sys 0.61.2", + "windows-sys 0.59.0", ] [[package]] @@ -8069,6 +8125,17 @@ version = "1.15.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "67b1b7a3b5fe4f1376887184045fcf45c69e92af734b7aaddc05fb777b6fbd03" +[[package]] +name = "smartstring" +version = "1.0.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3fb72c633efbaa2dd666986505016c32c3044395ceaf881518399d2f4127ee29" +dependencies = [ + "autocfg", + "static_assertions", + "version_check", +] + [[package]] name = "socket2" version = "0.6.3" @@ -8188,6 +8255,12 @@ dependencies = [ "system-deps", ] +[[package]] +name = "spin" +version = "0.5.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6e63cff320ae2c57904679ba7cb63280a3dc4613885beafb148ee7bf9aa9042d" + [[package]] name = "spki" version = "0.7.3" @@ -8924,7 +8997,7 @@ dependencies = [ "getrandom 0.4.2", "once_cell", "rustix", - "windows-sys 0.61.2", + "windows-sys 0.59.0", ] [[package]] @@ -8948,6 +9021,12 @@ dependencies = [ "utf-8", ] +[[package]] +name = "thin-vec" +version = "0.2.18" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b0f7e269b48f0a7dd0146680fa24b50cc67fc0373f086a5b2f99bd084639b482" + [[package]] name = "thiserror" version = "1.0.69" @@ -9079,11 +9158,13 @@ dependencies = [ [[package]] name = "tinyagents" -version = "1.5.0" +version = "1.7.1" dependencies = [ "async-trait", + "bytes", "futures", "reqwest 0.12.28", + "rhai", "rusqlite", "serde", "serde_json", @@ -9092,9 +9173,30 @@ dependencies = [ "tokio", ] +[[package]] +name = "tinycortex" +version = "0.1.1" +dependencies = [ + "anyhow", + "async-trait", + "chrono", + "git2", + "parking_lot", + "rand 0.8.6", + "regex", + "rusqlite", + "serde", + "serde_json", + "sha2 0.10.9", + "thiserror 2.0.18", + "toml 0.8.2", + "uuid 1.23.1", + "walkdir", +] + [[package]] name = "tinyflows" -version = "0.3.0" +version = "0.3.2" dependencies = [ "async-trait", "futures-timer", @@ -10485,7 +10587,7 @@ version = "0.1.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22" dependencies = [ - "windows-sys 0.61.2", + "windows-sys 0.48.0", ] [[package]] @@ -11775,7 +11877,3 @@ dependencies = [ "syn 2.0.117", "winnow 1.0.2", ] - -[[patch.unused]] -name = "tinycortex" -version = "0.1.1" diff --git a/docs/plans/tinyflows-integration/README.md b/docs/plans/tinyflows-integration/README.md new file mode 100644 index 000000000..bdb8e6ef7 --- /dev/null +++ b/docs/plans/tinyflows-integration/README.md @@ -0,0 +1,304 @@ +# tinyflows — Completion & Integration Plan (POA) + +> **Status**: proposed · **Date**: 2026-07-04 +> **Scope**: finish integrating the vendored `vendor/tinyflows/` workflow engine (our n8n/Zapier module) into the Rust core and the desktop UI, and close the feature gaps that separate it from a real n8n-class product. +> **Reference product**: [n8n](https://github.com/n8n-io/n8n) — item-based data flow, `=`-expressions, node catalog, editable canvas, credentials, webhook/schedule triggers, executions view. + +--- + +## 1. Where we are today + +The integration is **not greenfield**. Three layers already exist; each is at a different maturity level. + +### 1.1 Engine (`vendor/tinyflows/`, v0.3.0 — mostly done) + +Host-agnostic library crate on the `tinyagents` state-graph runtime. Pipeline: `WorkflowGraph → migrate → validate → compile → engine::run`. + +| Area | State | +| --------------- | -------------------------------------------------------------------------------------------------------------------------------------------------------- | +| Node kinds (12) | ✅ `trigger, agent, tool_call, http_request, code, condition, switch, merge, split_out, transform, output_parser, sub_workflow` | +| Routing | ✅ linear, conditional (ports), parallel fan-out, merge fan-in barrier | +| Data flow | ✅ n8n-style item arrays `{ json, binary?, paired_item? }`; `=`-expressions with **full jq** (`jaq-core`) already wired in `src/expr.rs` | +| Error handling | ✅ per-node `on_error` (stop/continue/route), retry with fixed/exponential backoff, `node_timeout_secs` | +| HITL | ✅ `requires_approval` → interrupt; in-process (`run_resumable`) and cross-process durable resume (`resume_with_checkpointer(thread_id)`) | +| Observability | ✅ `RunObserver` trait (`on_run_start/on_step_finish/on_run_finish`), tracing spans, journaled variants (`GraphEventJournal`) for Langfuse | +| Persistence | By design **none** — host injects `Checkpointer`, `StateStore`, and persists `Run`/`ExecutionStep` | +| Capabilities | ✅ host-injected traits: `LlmProvider`, `ToolInvoker`, `HttpClient`, `CodeRunner`, `StateStore`; opaque `connection_ref` (host resolves secrets) | +| Versioning | ✅ `schema_version` (graph) + per-node `type_version`, `migrate()` pre-parse | +| Triggers | ⚠️ **declarative only** — `manual, schedule, webhook, app_event, form, execute_by_workflow, chat_message, evaluation, system`; the _host_ must fire them | + +Engine-side gaps (see Phase 7): `agent` node sub-ports (chat_model/memory/tool/output_parser) stubbed; `output_parser` is identity passthrough; README/Roadmap lag the code (claim jq and retry backoff are pending when both ship). + +### 1.2 Rust core seam (`src/openhuman/flows/` + `src/openhuman/tinyflows/` — implemented, with holes) + +- **Domain** `flows::` (~3,700 lines + tests): `types.rs` (`Flow` wraps `WorkflowGraph` + `enabled`/`require_approval`/`last_status`; `FlowRun`, `FlowRunStep`, `FlowRunTrigger::{Rpc,Schedule,AppEvent,Resume}`), `store.rs` (SQLite incl. `flow_state` kv), `ops.rs` (validate/migrate + full run/resume under `TrustedAutomation → Workflow` origin, 600 s timeout), `schemas.rs`, `tools.rs` (`ProposeWorkflowTool` — validate-only, never persists), `bus.rs` (`FlowTriggerSubscriber`). +- **RPC surface** (10 methods, wired in `src/core/all.rs`): `openhuman.flows_{create,get,list,update,delete,set_enabled,run,resume,list_runs,get_run}`. +- **Capability seam** `src/openhuman/tinyflows/`: `caps.rs` (LLM/Composio-tools/HTTP/code/state adapters), `observability.rs` (currently `NoopObserver`), `langfuse_export.rs` (post-run trace export). +- **Schedule triggers work end-to-end**: `flows::ops::bind_schedule_trigger` registers a cron `JobType::Flow` → scheduler publishes `DomainEvent::FlowScheduleTick` → `FlowTriggerSubscriber` runs the flow. +- **Composio `app_event` triggers also work end-to-end**: `flows/bus.rs::handle_app_event` matches `DomainEvent::ComposioTriggerReceived { toolkit, trigger }` against enabled `app_event` flows (case-insensitive toolkit/slug match, per-flow concurrency guard) and runs them under `FlowRunTrigger::AppEvent`. Since Composio triggers are delivered by the platform, this **is** our webhook story for third-party apps — no tunnel needed. + +**Known core gaps** (each becomes a workstream below): + +| # | Gap | Evidence | +| --- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | --------------------------------------------------------------------------- | +| G1 | **Raw `webhook` triggers unwired** (mitigated: Composio `app_event` triggers already cover third-party apps) — enabling a webhook-trigger flow only logs a warning; `bus.rs` observes `WebhookIncomingRequest` but explicitly does not dispatch; Composio trigger *subscriptions* are not auto-provisioned on enable | `flows/ops.rs:298` (`log_webhook_trigger_deferred`), `flows/bus.rs:232-242` | +| G2 | **No live run observer** — `NoopObserver`; `FlowRunStep`s reconstructed post-hoc from final state, no per-step timing/attempts, nothing streamed while running | `flows/types.rs:76`, `flows/ops.rs:781-783` (`TODO(0.3)`) | +| G3 | **Credential / connected-account resolution stubbed** — Composio nodes fall back to the ambient signed-in account; toolkit allow-listing hard-rejects real toolkits; HTTP credential resolution unimplemented (`connection_ref` is accepted but unresolvable) | `tinyflows/caps.rs:193,235-261,330,376,408` | +| G4 | **No cancel/deny** — a dismissed approval leaves the run parked `pending_approval` forever; no `flows_cancel`/`flows_deny` RPC | UI comment in `FlowApprovalCard.tsx` | +| G5 | **No JSON-RPC E2E coverage** — zero `openhuman.flows_*` calls in `tests/json_rpc_e2e.rs` (unit tests only) | grep of `tests/*.rs` | +| G6 | Unfired trigger kinds: `chat_message`, `form`, `execute_by_workflow` (as a _trigger_), `evaluation`, `system` have no host dispatcher | `flows/bus.rs` | + +### 1.3 Frontend (`app/src/` — reachable, read-only) + +Shipped: `/flows` nav tab (FlowsPage list: enable toggle, Run, last status), `/flows/:id` **read-only** canvas (`@xyflow/react` v12, custom nodes, minimap), `FlowRunsDrawer` + `FlowRunInspectorDrawer` (2 s polling via `useFlowRunPoller`), chat `WorkflowProposalCard` (agent `propose_workflow` → user "Save & enable" is the **only** creation path), `FlowApprovalCard` HITL notifications. All components have co-located Vitest coverage; i18n namespaces `flows.*`/`flowRuns.*`/`notifications.flow.*` exist across locales. + +**Known UI gaps**: + +| # | Gap | Evidence | +| --- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------ | ---------------------------------------- | +| U1 | **No canvas editing** — `nodesDraggable/Connectable/elementsSelectable: false`; `xyflowToWorkflowGraph` in `graphAdapter.ts` is dead code awaiting the editor ("B5b.2+") | `components/flows/canvas/FlowCanvas.tsx` | +| U2 | **No node config panel** — clicking a node does nothing; `config` shown only as truncated hint chips | — | +| U3 | **No authoring entry** — "New workflow" navigates to `/chat` with a TODO; empty-state copy promises canvas creation that doesn't exist | `pages/FlowsPage.tsx` | +| U4 | No trigger-config UI, no credentials picker, no template gallery, no import/export | — | +| U5 | Approval "Dismiss" is client-only (blocked on G4) | `FlowApprovalCard.tsx` | +| U6 | Run progress is poll-only; no socket push (blocked on G2) | `hooks/useFlowRunPoller.ts` | + +### 1.4 Disambiguation (do not conflate) + +The repo has **three** "workflow" systems. This plan touches only the first: + +1. `flows::` / `openhuman.flows_*` — **tinyflows typed graphs** (this plan). +2. `workflows::` / `openhuman.workflows_*` — WORKFLOW.md/SKILL.md bundle discovery/install (separate product surface under `/skills`). **Slated for decommission** — it is essentially the skills feature wearing the "workflows" name; see Phase 8. Retiring it frees the "Workflows" branding for tinyflows (the `/flows` nav tab already reads "Workflows"). +3. `rlm::` — Rhai `.ragsh` language workflows (`docs/plans/rlm-workflows/`), positioned in `gitbooks/features/orchestration.md` as the _next_ layer on the same substrate. tinyflows remains the shipping visual/typed product; rlm does not replace it. + +--- + +## 2. Gap analysis vs n8n / Zapier + +What an n8n user would expect, mapped to our state: + +| n8n capability | tinyflows/OpenHuman today | Covered by | +| ----------------------------------------------------- | ---------------------------------------------------------------------------- | ----------------- | +| Editable node canvas (drag, connect, add/delete) | Read-only viewer | Phase 3 | +| Node config side-panel with schema-driven forms | Missing | Phase 3 | +| Webhook trigger with live URL | Declared but never fired | Phase 1 | +| Cron/interval trigger | ✅ shipped | — | +| App-event trigger (≈ Zapier polling/instant triggers) | Wired via Composio `ComposioTriggerReceived`, but account resolution stubbed | Phase 1 (G3) | +| Credentials store + per-node credential picker | `connection_ref` plumbing exists; resolution unimplemented; no UI | Phase 2 | +| Executions list + live run view | List + inspector exist; polling, post-hoc steps only | Phase 1–2 | +| Cancel a running/paused execution | Missing | Phase 1 | +| Expressions (`={{ }}` in n8n) | ✅ `=`-prefix + full jq | — | +| Error workflow / retry per node | ✅ on_error/retry/backoff | — | +| Sub-workflows | ✅ `sub_workflow` node (inline graph); _by-id_ reference missing | Phase 4 | +| Templates gallery | Missing | Phase 4 | +| Import/export workflow JSON, n8n import | Missing (feasible: model is deliberately n8n-shaped) | Phase 4 | +| Partial execution / pin data / step-through debug | Missing | Phase 5 (stretch) | +| Waiting/Wait node, human approval | ✅ HITL approval + durable resume | — | +| AI Agent node with attached tools/memory | `agent` node = bare LLM call; sub-ports stubbed | Phase 6 | +| Versioning of workflow definitions | ✅ schema_version + type_version + migrate | — | +| Multi-user sharing/RBAC | Out of scope (desktop, single-user + team domain later) | — | + +Our structural advantages to preserve: the agent can _author_ workflows conversationally (`propose_workflow` → proposal card → single human "Save & enable" persistence gate), everything runs under the security policy/approval-gate substrate, and secrets never enter the engine (opaque `connection_ref`). + +--- + +## 3. Plan of action + +Ordering rationale: **backend correctness first** (triggers, cancellation, live observation, credentials) because every UI phase consumes those RPCs; the editable canvas is the largest UI lift and is independent of trigger work, so it proceeds in parallel from Phase 3. + +### Phase 0 — Substrate hygiene: tinyagents pin & tags + +Every workflow is, at run time, a **unique tinyagents state graph**: `compile()` wraps the `WorkflowGraph`, and the engine lowers it per run into a tinyagents graph (nodes → graph nodes, output ports → conditional edges, fan-in → waiting edges); run state (`MergeReducer` over a single JSON value), durable checkpointing, HITL interrupts, and the event journal are all tinyagents primitives, keyed per run by `thread_id`/`run_id`. So the health of the tinyagents pin is the health of the whole feature. + +Audit (2026-07-04): + +- ✅ **Single unified copy** — both Cargo worlds (`Cargo.toml:355`, `app/src-tauri/Cargo.toml:214`) patch `tinyagents` to `vendor/tinyagents`; exactly one `tinyagents 1.5.0` in both lockfiles. tinyflows' `tinyagents = "1.2"` requirement unifies onto the same copy (semver-compatible). No duplication, no type-identity risk. +- ⚠️ **Not on a tag** — the submodule is pinned at `df391c4` = `v1.5.0-13-gdf391c4`: 13 untagged commits past v1.5.0 (REPL host-embedding / cancel work taken early for the rlm feature). The crate still self-reports `1.5.0`, so the version string understates the vendored code. +- ⚠️ **Two minor versions behind** — upstream is at **v1.7.1** (60 commits past v1.5.0). The 13 early-adopted commits all landed upstream (via PR #19 etc.; `ReplSession::set_cancel_flag` verified present in v1.7.1), so retagging loses nothing. +- ⚠️ Same pattern on tinyflows itself: submodule at `v0.3.0-1-g438f8fc` (one commit past tag). + +Work items: + +1. Bump `vendor/tinyagents` submodule to the **v1.7.1 tag**; bump root requirement `tinyagents = { version = "1.7", features = ["sqlite", "repl"] }` (verify both features still exist in 1.7); `cargo update -p tinyagents` in both lockfiles. +2. Review the v1.5→v1.7 changelog for API breaks in the seams we touch: `Checkpointer`/`DurabilityMode`, `GraphEventJournal`/`GraphObservation`, interrupt/resume semantics (used by `flows::ops` + `src/openhuman/tinyagents/` + rlm). +3. Retag `vendor/tinyflows` on a proper release (v0.3.1) whose `Cargo.toml` requires `tinyagents = "1.7"` so the version story is coherent end-to-end. +4. Gate with `cargo check` both worlds, `pnpm test:rust`, and the flows/rlm unit suites; adopt a standing rule: **vendored tiny\* submodules pin release tags, never floating commits** (early-adopting an upstream PR requires a pre-release tag). + +### Phase 1 — Backend completion (triggers, lifecycle, observability) + +**1a. External event triggers — Composio-first (G1).** + +Product decision: **Composio triggers are the webhook story.** `app_event` dispatch already works end-to-end (see §1.2) and the platform handles inbound delivery, auth, and NAT traversal — so we do *not* build our own tunnel/webhook infrastructure for third-party app events. Remaining work is lifecycle and coverage, not plumbing: + +- **Trigger subscription lifecycle**: enabling an `app_event` flow should ensure the corresponding Composio trigger subscription exists upstream (create on `set_enabled(true)` via the `composio` domain, tear down on disable/delete when no other flow uses it) instead of assuming the user pre-configured it in Composio. Reconcile at boot like `bind_schedule_trigger` does. +- **Payload hygiene**: trigger payloads are untrusted input — route through `prompt_injection` screening before any `agent` node consumes them (applies to all external-event runs). +- **Trigger catalog surfacing**: expose the Composio trigger catalog (toolkit → available trigger slugs + payload schemas) over RPC so the UI (Phase 3 trigger config) and the builder agent (Phase 5 `search_tool_catalog`) offer real, connectable events instead of free-text slugs. +- **Raw `webhook` trigger kind → demoted**: for arbitrary custom HTTP callers not covered by a Composio toolkit, keep the `webhook` kind but defer implementation (generic inbound via `webhooks::ops`/`create_tunnel` stays a backlog item). Until then, validation must **warn loudly** at save/enable time that a `webhook`-trigger flow will not fire, instead of today's silent log line (`flows/ops.rs:298`). + +**1b. Run lifecycle: cancel + deny (G4).** + +- New RPCs `openhuman.flows_cancel_run(run_id)` (terminal `cancelled` status; abort the tokio task / drop the checkpointed thread) and deny semantics on resume: `flows_resume(id, thread_id, approvals, rejections)` → rejected node routes to its `error` port or fails the run. +- Sweep: TTL for parked `pending_approval` runs (align with the 10-min approval-gate TTL, configurable per flow). + +**1c. Live run observation (G2).** + +- Implement a real `RunObserver` in `src/openhuman/tinyflows/observability.rs`: `on_step_finish` → persist `FlowRunStep` incrementally (timing, attempt count, status, output) via `flows::store`, and publish a new `DomainEvent::FlowRunProgress { run_id, node_id, status }`. +- Bridge to the frontend socket (`socket` domain) so the inspector can subscribe instead of polling (UI lands in Phase 2/3; keep polling as fallback). +- Wire the journaled run variants so Langfuse export happens per-step, not only post-run. + +**1d. Remaining trigger kinds (G6).** + +Full taxonomy, the generic event-dispatcher design, and the **cron → workflows migration** live in the companion doc **[triggers.md](triggers.md)**. Summary: + +- `form` + one-shot/interval schedule sugar first (cheap wins). +- A **trigger catalog** (curated allowlist of `DomainEvent`s with stable payload schemas + jq filters) feeding one generic `system`-kind dispatcher in `FlowTriggerSubscriber` — instead of a bespoke handler per event. Guardrails: loop prevention via run provenance depth, per-flow rate limits, payload hygiene/risk classes. +- `chat_message` from the channels pipeline (untrusted-input guardrails required first). +- **Cron unification**: legacy `JobType::Shell`/`Agent` cron jobs become one-node scheduled flows (dual-write bridge, then backfill migration); cron remains only as the internal tick engine. "Everything automated is a workflow." +- Until a kind's dispatcher ships, validation warns loudly at save/enable instead of silently never firing. + +**1e. E2E tests (G5).** + +- Extend `tests/json_rpc_e2e.rs` (+ `scripts/test-rust-with-mock.sh`): full lifecycle round-trip (create → run → steps → resume → cancel → delete), schedule-tick dispatch, webhook dispatch, approval park/deny. This is the merge gate for 1a–1d. + +**Deliverable**: every trigger kind either fires or loudly warns; runs can be cancelled/denied; run steps stream. + +### Phase 2 — Credentials & connections (G3) + +The `connection_ref` seam exists end-to-end but resolves nothing. This phase makes integration nodes real. + +- **Composio connected accounts**: resolve `connection_ref` → specific Composio connected account (today: ambient-account fallback). Fix toolkit allow-listing in `caps.rs:193` so real toolkits pass policy instead of being hard-rejected. Surface account choice in the node config (`connection_ref` = connected-account id). +- **HTTP credentials**: back `connection_ref` for `http_request` nodes with the existing `credentials` domain (header/bearer/basic templates stored encrypted; injected server-side in `OpenHumanHttp::request` — never returned to the UI or the engine). +- **RPCs**: `flows_list_connections()` (aggregate Composio connected accounts + stored HTTP credentials, ids + display names only) for the UI picker. +- **Policy**: `code` and `http_request` node execution must respect the `[autonomy]` tier / `classify_command`-equivalent gating (network class); document the matrix in `flows/ops.rs`. + +**Deliverable**: an `http_request` or Composio `tool_call` node can be pointed at a chosen account/credential by reference, with secrets never leaving the core. + +### Phase 3 — Editable canvas (U1–U2) — the n8n builder + +Largest UI phase; runs in parallel with Phase 2 once Phase 1c's RPCs exist. + +- **3a. Edit mode in `FlowCanvas`**: flip the readonly defaults behind an `editable` prop; enable drag (persist `position`), connect (port-aware: derive valid source/target handles from node kind — reuse `graphAdapter` port logic), delete nodes/edges, and a node palette (12 kinds with the existing emoji/accent metadata). Wire the already-written `xyflowToWorkflowGraph` as the save path → `flows_update`. +- **3b. Node config panel**: right-hand drawer on node select. v1 pragmatic approach: per-kind form components for the high-traffic kinds (`trigger` schedule/webhook config, `http_request` method/url/headers/body, `agent` prompt/model, `tool_call` slug/args, `condition`/`switch` expression, `transform` set-map, `code` editor with language toggle) + a raw-JSON escape hatch for the rest. `=`-expression fields get a monospace input with an "expression" affordance (full expression-editor with live preview is a stretch goal — needs a `flows_eval_expr` RPC against sample run data). +- **3c. Validation UX**: new RPC `openhuman.flows_validate(graph)` (thin wrapper over `ops::validate_and_migrate_graph`, same path `propose_workflow` uses) → inline canvas errors (missing trigger, cycle, invalid config on node X) before save. +- **3d. Draft/dirty state**: local draft in component state; explicit Save; unsaved-changes guard. No autosave in v1 (a saved+enabled flow is live — accidental saves fire schedules). +- **3e. Live run overlay**: subscribe to `FlowRunProgress` socket events; animate node status on the canvas during a run (n8n's signature interaction) and in the inspector. + +**Deliverable**: create/edit a workflow entirely on the canvas; watch it execute live. + +### Phase 4 — Authoring entry points, templates, interop (U3–U4) + +- **4a. "New workflow"**: replace the `/chat` TODO with a chooser — _Start from scratch_ (blank canvas with a trigger node), _Describe it_ (interim: prefill chat composer → `propose_workflow`; superseded by the Phase 5 in-place prompt bar), _From template_. +- **4b. Sub-workflow by id**: extend `sub_workflow`/`execute_by_workflow` to reference a saved `flow_id` (engine currently only inlines a child graph) — engine change upstreamed to `vendor/tinyflows` + host resolver in `caps`/ops. +- **4c. Templates**: ship 5–10 curated `WorkflowGraph` JSONs (bundled resources, like agent prompts): e.g. "Daily digest to channel", "Webhook → agent triage → notify", "Scheduled scrape → transform → memory". Gallery UI on FlowsPage empty state. +- **4d. Import/export**: export flow JSON; import with `migrate()` + validate. **n8n importer** (host-side, best-effort): map n8n workflow JSON → `WorkflowGraph` for the overlapping vocabulary (IF→condition, Switch, Merge, SplitOut, HTTP Request, Code, Schedule/Webhook triggers; `={{...}}` → `=` jq where trivially translatable); unmapped nodes land as annotated placeholder nodes rather than failing the import. +- **4e. Proposal card upgrade**: "Open in canvas" action on `WorkflowProposalCard` (review/edit before Save & enable) — keeps the single persistence gate. + +### Phase 5 — Prompt-first authoring: the Workflow Builder agent + +The product stance ("the agent builds it, you approve it") gets a first-class surface: users prompt **from the Workflows UI itself** — not by wandering into `/chat` — and a dedicated, tool-scoped agent designs and iterates on the graph in place. This is the differentiator over n8n, so it deserves its own phase rather than being a bullet under 4a. + +**5a. A dedicated `workflow-builder` agent definition.** + +The harness already supports exactly this shape: `AgentDefinition` (`src/openhuman/agent/harness/definition.rs`) with `ToolScope::Named`, registered as a builtin in `harness/builtin_definitions.rs` and overridable by user TOML via `definition_loader.rs`. + +- New builtin definition `workflow-builder` (Worker tier): system prompt specialized for workflow design — knows the 12 node kinds, `=`/jq expression semantics, port/edge rules, trigger kinds and which ones are live, error-handling config (`on_error`/`retry`), and the "propose, never persist" invariant. Prompt ships in `src/openhuman/agent/prompts/` like the other bundled prompts. +- `ToolScope::Named` — deliberately narrow toolset (see 5b): no shell, no file writes, no channel sends. `external_effect = false` end-to-end; the only side effects it can cause are validated *proposals*. +- Reachable two ways: (1) directly from the Flows UI prompt surface (5c), spawned with the flow/draft as context; (2) by delegation from the main agent via the existing delegation tools (`agent/tools/delegate_to_personality.rs`, archetype/skill delegation in `agent_orchestration::tools`) so "set up a workflow that…" in normal chat routes to the specialist automatically. + +**5b. The builder's tool belt** (new tools live in `flows/tools.rs` per the tool-ownership rule, registered in `tools/ops.rs`): + +| Tool | Status | Purpose | +| --- | --- | --- | +| `propose_workflow` | ✅ exists | Validate a full graph + emit proposal payload (unchanged invariant: never persists) | +| `revise_workflow` | new | Take the current draft graph + an instruction diff → emit an updated proposal; enables iterative refinement instead of regenerate-from-scratch | +| `list_workflows` / `get_workflow` | new (read-only) | Inspect existing flows so the agent can reference, clone, or avoid duplicating them | +| `get_workflow_run` | new (read-only) | Read a failed run's steps so the agent can debug/repair a workflow from an error report | +| `list_flow_connections` | new | Enumerate Composio connected accounts + stored HTTP credentials (ids/names only — Phase 2's RPC surfaced as a tool) so generated nodes carry valid `connection_ref`s | +| `search_tool_catalog` | new | Search the Composio/tools registry for real tool slugs, so `tool_call` nodes are grounded in tools that actually exist (today the agent can hallucinate slugs) | +| `dry_run_workflow` | new | Execute the *draft* graph against mock/sandboxed capabilities (`tinyflows` `mock` feature or capped real caps with `requires_approval` forced on) and return step results — lets the agent self-verify before proposing | + +All read-only tools return `PermissionLevel::None`; `dry_run_workflow` is gated by autonomy tier since `code`/`http_request` nodes could execute. + +**5c. Prompt surface in the Flows UI.** + +- **FlowsPage prompt bar**: a "Describe a workflow…" composer at the top of `/flows` (and as the empty-state hero). Submitting spawns a `workflow-builder` turn in a dedicated thread; the resulting proposal renders inline (reuse `WorkflowProposalCard`) with **Open in canvas** and **Save & enable**. +- **Canvas copilot panel**: on `/flows/:id` (and on drafts), a side panel chat bound to the same agent with the current graph injected as context. Each agent proposal updates a **draft overlay** on the canvas (diff-style: added nodes highlighted, removed ones ghosted) — accept/reject applies it to the local draft from Phase 3d. This is `revise_workflow` in a loop: "add a Slack notification on failure", "make the schedule weekdays only", "split this into a sub-workflow". +- **Repair entry point**: from a failed run in `FlowRunInspectorDrawer`, "Fix with agent" opens the copilot with the run's failing step context preloaded (`get_workflow_run`). +- Plumbing reuses the existing chat runtime (`ChatRuntimeProvider` already parses `propose_workflow` outputs into `pendingWorkflowProposalsByThread`); the new work is thread scoping per draft/flow, the canvas diff overlay, and routing turns to the `workflow-builder` definition instead of the main agent. + +**5d. Invariants** (carry over from the current design, enforced in review): + +- The agent **never persists or enables** a flow — `flows_create`/`flows_update`/`set_enabled` remain UI-only actions behind an explicit user click. +- Proposals are always re-validated server-side at save time (`flows_validate` path), never trusted from the client. +- `dry_run_workflow` output is labeled as sandbox output in the UI so users don't mistake it for a live run. + +**Deliverable**: a user types "every Monday, summarize my unread Slack messages and email me" into the Flows page, watches the graph appear on the canvas, iterates in plain language, then clicks Save & enable. + +### Phase 6 — Polish & debug tooling (stretch) + +- Partial execution ("run from this node" with pinned upstream data) — needs engine support for seeding node state; scope with tinyflows maintainers. +- Run diff/inspector niceties: per-item data browser (n8n's table/JSON toggle), input↔output pairing via `paired_item`. +- `flows_duplicate`, run retention/pruning policy, per-flow run-history limits. +- Desktop E2E (WDIO) spec: create → run → inspect happy path. + +### Phase 7 — Engine (vendor/tinyflows) upstream work + +Tracked separately since it's a submodule with its own release cadence (host pins `0.3`, patched to path): + +1. `agent` node sub-ports (chat_model/memory/tool/output_parser wiring) — unlocks n8n-style "AI Agent with tools" composition. +2. `output_parser`: schema validation + LLM auto-fix (currently identity). +3. Sub-workflow by reference (4b) — config `workflow_id` alternative to inline `workflow`. +4. Docs truth-up: README/Roadmap still claim jq and retry backoff are pending; both are implemented (`src/expr.rs` routes to jaq; `engine.rs` has fixed/exponential backoff + `node_timeout_secs`). +5. Optional: cancellation token support in `engine::run` (cleaner than task-abort for 1b). + +### Phase 8 — Decommission the legacy `workflows::` bundle domain + +The `workflows::` domain (WORKFLOW.md/SKILL.md bundle discovery/install, RPC `openhuman.workflows_*`) predates tinyflows and is functionally the **skills** feature under a different name. Keeping two things called "workflows" confuses users, agents, and contributors alike. Plan: fold what's unique into `skills`, delete the rest, and hand the name to tinyflows. + +- **Audit consumers first**: `agent/tools/run_workflow.rs` (agent tool that runs WORKFLOW.md bundles — decide: retire, or repoint to `flows_run`/skills), the `/skills` UI surfaces (`WorkflowsTab`, `CreateWorkflowForm`, `WorkflowRunnerBody`, `WorkflowNew.tsx`, `WorkflowsRun.tsx`, `DevWorkflowPanel`, `workflowsApi.ts`), `about_app`, gitbooks, and the rlm plan's references to `run_workflow` as a composition surface. +- **Migrate**: bundle discovery/install semantics that skills doesn't already cover move into the `skills` domain (it is metadata-only post-QuickJS-removal, so this is mostly file-format and registry work). +- **Deprecate then delete**: mark `openhuman.workflows_*` deprecated for one release (RPC responses carry a deprecation notice), then remove `src/openhuman/workflows/`, its controllers from `src/core/all.rs`, the frontend clients/pages, and the `/workflows/new`//`workflows/run` routes (bare `/workflows` already redirects to `/settings/automations`). +- **Not in scope**: `openhuman.workflow_run_*` (`agent_orchestration`'s declarative run ledger) is a different system and untouched here — though its name should also be revisited once "Workflows" ≡ tinyflows. +- **Naming end-state**: one user-facing concept — **Workflows = tinyflows graphs** at `/flows`; skills are skills. + +--- + +## 4. Missing-features summary (checklist) + +Substrate: ☐ tinyagents submodule → v1.7.1 tag (+ root req bump, both lockfiles) · ☐ v1.5→v1.7 API-break review · ☐ tinyflows retag (v0.3.1, `tinyagents = "1.7"`) · ☐ tags-only submodule policy. + +Backend: ☐ Composio trigger-subscription lifecycle (auto-provision on enable) · ☐ trigger catalog RPC · ☐ loud validation warning for unfired trigger kinds · ☐ raw-webhook dispatch (deferred backlog) · ☐ `flows_cancel_run` · ☐ resume-with-rejection/deny · ☐ live `RunObserver` + incremental step persistence · ☐ `FlowRunProgress` socket events · ☐ `flows_validate` RPC · ☐ `flows_list_connections` · ☐ Composio connected-account resolution · ☐ HTTP credential resolution · ☐ toolkit allow-list fix · ☐ `chat_message` trigger dispatch · ☐ sub-workflow by id · ☐ parked-run TTL sweep · ☐ JSON-RPC E2E suite. + +Frontend: ☐ editable canvas (drag/connect/palette/delete) · ☐ node config panels · ☐ trigger config UI (cron builder, webhook URL display, app-event picker) · ☐ credentials picker · ☐ new-workflow chooser · ☐ template gallery · ☐ import/export + n8n import · ☐ live canvas run overlay (socket) · ☐ approval deny (real) · ☐ "Open in canvas" from proposal card · ☐ WDIO E2E spec. + +Agent authoring: ☐ `workflow-builder` builtin `AgentDefinition` + specialized prompt · ☐ `revise_workflow` tool · ☐ read-only `list_workflows`/`get_workflow`/`get_workflow_run` tools · ☐ `list_flow_connections` tool · ☐ `search_tool_catalog` tool · ☐ `dry_run_workflow` (sandboxed) · ☐ delegation routing from main agent · ☐ FlowsPage prompt bar · ☐ canvas copilot panel with draft diff overlay · ☐ "Fix with agent" from failed runs. + +Engine (upstream): ☐ agent sub-ports · ☐ output_parser validation · ☐ `workflow_id` sub-workflows · ☐ docs truth-up · ☐ cancellation. + +Decommission: ☐ `run_workflow` tool disposition · ☐ bundle semantics folded into `skills` · ☐ `workflows_*` RPC deprecation release · ☐ delete `src/openhuman/workflows/` + frontend pages/clients · ☐ docs/about_app rename sweep. + +--- + +## 5. Sequencing, estimates, risks + +``` +Phase 0 (tinyagents pin) ██ ~2–3 d do first; touches both lockfiles +Phase 1 (backend) ██████ ~2–3 wk unblocks everything +Phase 2 (credentials) ████ ~1–2 wk parallel w/ Phase 3 +Phase 3 (canvas) ████████ ~3–4 wk largest UI lift +Phase 4 (authoring) ████ ~2 wk +Phase 5 (builder agent) ██████ ~2–3 wk needs Phase 3 canvas + Phase 2 connections +Phase 6 (polish) ██ opportunistic +Phase 7 (engine) ── continuous, PR'd to vendor submodule ── +Phase 8 (decommission) ── independent; audit early, delete after a deprecation release ── +``` + +**Risks / decisions to confirm:** + +1. **External-event delivery model** — decided: Composio triggers are the webhook story (platform handles delivery/auth/NAT). Raw custom webhooks (own tunnel via `webhooks::ops::create_tunnel`) stay deferred; revisit only if users need non-Composio HTTP callers. +2. **Prompt-injection surface** — webhook/app_event payloads feeding `agent` nodes are untrusted; must go through the `prompt_injection` domain. Non-negotiable before G1 ships. +3. **Editable canvas vs. agent-first authoring** — product stance so far is "the agent builds it, you approve it" (`gitbooks/features/workflows.md`). Phase 5 makes that stance first-class in the Flows UI itself (prompt bar + canvas copilot); the Phase 3 hand-editing canvas is the escape hatch, not the headline. Keep the prompt-first path primary in onboarding copy so the two don't compete. +4. **GPL-3.0 licensing** of `vendor/tinyflows` — already vendored/linked; re-confirm distribution posture before expanding surface (flag for legal, not an engineering blocker). +5. **Submodule cadence** — engine changes (Phase 7) must land in the tinyflows repo and re-vendor; keep host-side workarounds (e.g. by-id sub-workflow resolution in caps) so UI phases aren't blocked on vendor releases. + +**Per-repo conventions that apply to all phases**: new RPCs go through domain `schemas.rs` + controller registry (no `dispatch.rs` branches); ≥80 % diff coverage gate; verbose `[flows]`-prefixed debug logging on every new path; i18n keys added to `en.ts` **and all 13 locales**; update `src/openhuman/about_app/` and `gitbooks/features/workflows.md` as features land. diff --git a/docs/plans/tinyflows-integration/triggers.md b/docs/plans/tinyflows-integration/triggers.md new file mode 100644 index 000000000..526736a37 --- /dev/null +++ b/docs/plans/tinyflows-integration/triggers.md @@ -0,0 +1,140 @@ +# Trigger taxonomy & the cron → workflows unification + +> Companion to [README.md](README.md) (the tinyflows completion POA). This doc answers two questions: **what kinds of triggers can this platform offer**, and **how do we migrate cron so that "everything automated is a workflow"**. + +## 0. Thesis + +Almost every automated behavior in OpenHuman is already shaped like *event → conditions → actions*: cron jobs, Composio app events, channel-message reactions, subconscious escalations, meeting follow-ups, notification triage. Today each has bespoke plumbing. tinyflows gives us one uniform substrate (a tinyagents graph per run, durable, observable, approval-gated), and the event bus (`src/core/event_bus/events.rs`, ~130 `DomainEvent` variants) is already the spine every one of those signals travels on. + +**End state**: *Workflows are the single user-facing automation concept.* A trigger is just a subscription that seeds a flow run; the trigger catalog is a curated projection of the event bus plus time and manual entry points. + +## 1. What the engine already models + +`tinyflows::model::TriggerKind` (declarative — the host fires them): `manual`, `schedule`, `webhook`, `app_event`, `form`, `execute_by_workflow`, `chat_message`, `evaluation`, `system`. + +Host dispatch status today (`flows/bus.rs`): + +| Kind | Status | +| --- | --- | +| `manual` | ✅ `flows_run` RPC / UI Run button | +| `schedule` | ✅ cron `JobType::Flow` → `FlowScheduleTick` | +| `app_event` | ✅ `ComposioTriggerReceived` matching (toolkit + trigger_slug) | +| `webhook` | ⏸ deferred (Composio is the webhook story — README Phase 1a) | +| `chat_message`, `form`, `execute_by_workflow`, `evaluation`, `system` | ❌ no dispatcher | + +The gap is not engine capability — it's host dispatchers and a catalog. Crucially, `system` is a free config surface: we can build **one generic event-trigger dispatcher** instead of nine bespoke ones. + +## 2. Proposed trigger catalog + +Organized by source. "Backing event" = the existing `DomainEvent` variant(s) that would feed the dispatcher. + +### 2.1 Time (`schedule`) + +Already live via cron. Expose the full richness cron already has (`CronJob` in `cron/types.rs`): + +- **Cron expression** (today's path) and **interval** ("every 15 min") sugar. +- **One-shot** ("run once at T") — cron's `delete_after_run` already models this; surface it in trigger config. +- **Windowed schedules** (weekdays only, quiet hours) — respects `scheduler_gate`. + +### 2.2 Manual & parameterized (`manual`, `form`) + +- **Manual**: Run button / RPC / agent-initiated (`flows_run` with input). +- **Form**: a manual run with a typed input schema. Trigger config declares fields; the UI renders a form; submission becomes the trigger payload. This is n8n's Form trigger and gives every workflow a shareable "mini-app" entry point. Cheap to ship: validation + a generated form in the canvas/run dialog. + +### 2.3 External apps (`app_event`, `webhook`) + +- **Composio triggers** (live): Gmail new message, Slack event, GitHub push, calendar event, … the platform delivers; we match toolkit/slug. Work remaining is subscription lifecycle + catalog surfacing (README Phase 1a). +- **Raw webhook** (deferred): non-Composio HTTP callers via `webhooks::ops` tunnels. Backlog until demanded. + +### 2.4 Conversation (`chat_message`) + +Backing events: `ChannelInboundMessage` / `ChannelMessageReceived` / `ChannelReactionReceived`. + +- **Message received** with filters: provider/channel id, sender, regex or keyword, mention-of-me, DM vs group. +- **Reaction added** (e.g. ✅ on a message → file it as a task). +- Guardrail: these runs consume untrusted text — `prompt_injection` screening before any `agent` node, and a per-flow debounce/concurrency cap (the `app_event` dispatcher's guard generalizes). + +### 2.5 Agent & task lifecycle (`system`) + +Backing events: `AgentTurnCompleted`, `SubagentCompleted`/`SubagentFailed`, `ApprovalDecided`, `TaskSourceTaskIngested`, `TaskRunReclaimed`, `ThreadGoalUpdated`. + +Examples: "when any subagent fails, post a summary to my ops channel"; "when a task is ingested from a task source, enrich and label it"; "when an approval is denied, notify the requester thread". + +### 2.6 Memory & knowledge (`system`) + +Backing events: `MemoryStored`, `MemoryIngestionCompleted`, `MemoryDiffComputed`, `DocumentCanonicalized`, `TreeSummarizerHourCompleted`. + +Examples: "when a document is canonicalized, generate an abstract and store it"; "when the hourly summary lands, check for action items". + +### 2.7 Meetings & voice (`system`) + +Backing events: `MeetingSessionCreated`, `MeetingSummaryGenerated`, `BackendMeetTranscript`, `PttTranscriptCommitted`, `MeetingAutoJoinTriggered`. + +Examples: "when a meeting summary is generated, extract action items → create tasks → email attendees" — the flagship demo workflow. + +### 2.8 Notifications & devices (`system`) + +Backing events: `NotificationIngested`/`NotificationTriaged`, `DevicePaired`, `DevicePeerOnline`/`Offline`. + +Examples: triage automations ("when a notification from app X is ingested, decide urgent/ignore"), presence automations ("when my phone comes online, sync …"). + +### 2.9 Platform & health (`system`) + +Backing events: `SystemStartup`, `HealthChanged`, `HealthRestarted`, `AutonomyConfigChanged`, `McpServerDisconnected`, `ProviderApiKeyRejected`, `EmbeddingModelUnhealthy`. + +Examples: self-healing/maintenance workflows ("on health degradation, run doctor and report"). These should default to `require_approval` off but notification-heavy. + +### 2.10 Subconscious escalations (`system`) + +Backing events: `SubconsciousTriggerProcessed`, `TriggerEvaluated`, `TriggerEscalated`. + +The subconscious domain is itself a mini automation engine (evaluate → escalate). Long-term convergence candidate: an escalation's *action* becomes "run flow X", making workflows the actuator layer for subconscious signals rather than a parallel system. + +### 2.11 Workflow-to-workflow (`execute_by_workflow`, `evaluation`) + +- `execute_by_workflow`: fired when another flow's `sub_workflow`/by-id call targets this flow (README Phase 4b) — enables composition and shared "library" flows. +- `evaluation`: reserved for eval-harness runs of a flow (regression-test a workflow against recorded inputs); pairs with `dry_run_workflow` (README Phase 5b). + +## 3. Design: one generic event dispatcher, not N bespoke ones + +Adding a hand-written `handle_*` per event (like `handle_app_event`) doesn't scale to §2.5–2.10. Proposal: + +1. **Trigger catalog registry** (new, `flows/trigger_catalog.rs`): a curated allowlist of `DomainEvent` variants exposed as triggers — `{ key: "meeting.summary_generated", event: , payload_schema, description, risk_class }`. Only cataloged events are subscribable; the raw bus is never exposed wholesale. +2. **Generic dispatcher** in `FlowTriggerSubscriber`: one match arm per cataloged domain that projects the event into a stable JSON payload, then matches enabled flows whose trigger is `kind: system, config: { event: "", filter: "=" }`. The existing jq engine (`tinyflows::expr`) evaluates the optional filter against the payload — no new expression language. +3. **RPC** `flows_list_trigger_catalog()` → feeds the Phase 3 trigger-config UI and the Phase 5 builder agent (so prompts like "when a meeting ends…" ground to a real catalog key). +4. **Guardrails** (non-negotiable, enforced in the dispatcher): + - **Loop prevention**: flow runs emit events too; runs triggered by a flow-originated event carry a provenance chain with a max depth (reuse tinyagents `root_run_id` lineage); a flow can never trigger itself. + - **Debounce/rate limits**: per-flow trigger token bucket + the existing per-flow concurrency guard. + - **Payload hygiene**: catalog projection strips PII-ish fields by default; `prompt_injection` screening before agent nodes. + - **Risk classes**: catalog entries tagged (e.g. conversation events = untrusted-input class → forces screening; health events = internal class); validation surfaces the class in the UI. + +## 4. Cron → workflows migration + +Today `cron` runs three job types (`cron/types.rs::JobType`): `Shell`, `Agent`, and `Flow` (already just a tick-publisher for flows). The inversion: **cron stops being a product surface and becomes the schedule service for workflows.** + +**M1 — Model mapping.** Every legacy job is a one-node flow: + +- `JobType::Shell { command }` → flow: `trigger(schedule)` → `code` node (or a dedicated `shell` tool_call, subject to the same `classify_command` gating cron uses today). +- `JobType::Agent { prompt, agent_id, model, session_target }` → flow: `trigger(schedule)` → `agent` node (prompt/model in config; `agent_id` maps to the Phase 5+ agent-definition reference once engine sub-ports land — until then, definition prompt inlined). +- `DeliveryConfig` / `session_target` (post result to main thread vs isolated) → flow-level `delivery` setting or an explicit terminal "notify/post" node. Prefer the explicit node: it's visible on the canvas, which is the whole point. +- `delete_after_run` → one-shot schedule config (§2.1). + +**M2 — Dual-write bridge.** `cron_add`/`cron_update` agent tools and RPCs keep working but create flows under the hood (a `legacy_cron` tag preserves round-tripping for `cron_list`). New UI creation always goes through flows. + +**M3 — Backfill migration.** One-time migration (per-user, on upgrade): convert existing enabled shell/agent jobs into flows, preserving ids in metadata, next_run continuity, and run history linkage (`CronRun` rows stay readable; new runs are `FlowRun`s). Feature-flag + rollback window before deleting the legacy paths. + +**M4 — Surface consolidation.** Settings → Automations lists flows only; cron internals (`scheduler.rs`, jobs table) remain as the tick engine for `JobType::Flow`, which becomes the *only* job type. The `cron_*` tool family shrinks to a thin alias the agent can still call, documented as "creates a scheduled workflow". + +**What cron keeps**: the single tokio scheduler loop, persistence of next-fire times, catch-up semantics, `scheduler_gate`. That's infrastructure — workflows are the product on top of it. + +**Sequencing**: M1/M2 fit inside README Phase 1 (they're mostly `flows`-side); M3/M4 after the canvas + run inspector are solid (users must be able to *see* what their migrated jobs became), i.e. post-Phase 3. + +## 5. Rollout order for new trigger kinds + +1. `form` + one-shot/interval schedule sugar (cheap, immediately useful). +2. Trigger catalog + generic `system` dispatcher, seeded with a small vetted set: `meeting.summary_generated`, `notification.triaged`, `subagent.failed`, `task.ingested`, `document.canonicalized`. +3. `chat_message` (needs the untrusted-input guardrails first). +4. Cron migration M1–M2, then M3–M4. +5. `execute_by_workflow` (with Phase 4b sub-workflow-by-id), `evaluation` last. + +Each new dispatcher lands with: E2E test in `tests/json_rpc_e2e.rs`, catalog entry with payload schema, validation warning removal (the kind moves from "declared but never fires" to live), and a template (README Phase 4c) demonstrating it. diff --git a/src/core/event_bus/events.rs b/src/core/event_bus/events.rs index 92246a445..a38dbee0a 100644 --- a/src/core/event_bus/events.rs +++ b/src/core/event_bus/events.rs @@ -451,6 +451,23 @@ pub enum DomainEvent { /// Identifier of the `flows::Flow` to run. flow_id: String, }, + /// Live per-step progress of an in-flight `flows_run` / `flows_resume` + /// (issue G2, live run observation). Published from + /// `flows::observability::FlowRunObserver::on_step_finish` as each + /// non-trigger node settles, so the Workflows UI can show a run advancing + /// node-by-node instead of only polling the settled `flow_runs` row. The + /// durable `flow_runs` row (updated incrementally by the same observer and + /// finalized at settle) remains the source of truth — this event is a + /// best-effort progress feed (broadcast bridges drop on lag), which is why + /// the frontend keeps its 2s poller as a fallback. + FlowRunProgress { + /// The run's stable identifier (== the tinyflows checkpointer thread id). + run_id: String, + /// The node whose step just finished. + node_id: String, + /// Step outcome: `"success"` | `"error"`. + status: String, + }, // ── Skills ────────────────────────────────────────────────────────── /// A skill was loaded into the runtime. @@ -1309,7 +1326,8 @@ impl DomainEvent { | Self::CronJobCompleted { .. } | Self::CronDeliveryRequested { .. } | Self::ProactiveMessageRequested { .. } - | Self::FlowScheduleTick { .. } => "cron", + | Self::FlowScheduleTick { .. } + | Self::FlowRunProgress { .. } => "cron", Self::WorkflowLoaded { .. } | Self::WorkflowStopped { .. } @@ -1470,6 +1488,7 @@ impl DomainEvent { Self::CronDeliveryRequested { .. } => "CronDeliveryRequested", Self::ProactiveMessageRequested { .. } => "ProactiveMessageRequested", Self::FlowScheduleTick { .. } => "FlowScheduleTick", + Self::FlowRunProgress { .. } => "FlowRunProgress", Self::WorkflowLoaded { .. } => "WorkflowLoaded", Self::WorkflowStopped { .. } => "WorkflowStopped", Self::WorkflowStartFailed { .. } => "WorkflowStartFailed", diff --git a/src/core/socketio.rs b/src/core/socketio.rs index 1c9203a84..fb3c5b768 100644 --- a/src/core/socketio.rs +++ b/src/core/socketio.rs @@ -1048,6 +1048,30 @@ pub fn spawn_web_channel_bridge(io: SocketIo) { }); let _ = io_memory_sync.emit("init:completed", &payload); } + // Live per-step progress of an in-flight flow run (issue G2). + // Best-effort: the durable `flow_runs` row is the source of + // truth and the Workflows UI keeps a 2s poller as fallback, so + // a dropped event here (broadcast lag) only delays the live + // update, never corrupts run history. + crate::core::event_bus::DomainEvent::FlowRunProgress { + run_id, + node_id, + status, + } => { + let payload = serde_json::json!({ + "run_id": run_id, + "node_id": node_id, + "status": status, + }); + log::debug!( + "[socketio] broadcast flow_run_progress run_id={} node_id={} status={}", + run_id, + node_id, + status + ); + let _ = io_memory_sync.emit("flow:run_progress", &payload); + let _ = io_memory_sync.emit("flow_run_progress", &payload); + } _ => {} } } diff --git a/src/openhuman/flows/mod.rs b/src/openhuman/flows/mod.rs index dbeebe87d..d70743a04 100644 --- a/src/openhuman/flows/mod.rs +++ b/src/openhuman/flows/mod.rs @@ -9,6 +9,7 @@ pub mod bus; pub mod ops; +mod run_registry; mod schemas; mod store; pub mod tools; @@ -23,5 +24,9 @@ pub use schemas::{ // (`src/openhuman/tinyflows/caps.rs`) lives in a sibling domain and needs // them to implement `tinyflows::caps::StateStore` without duplicating the // `flow_state` table's persistence logic. -pub use store::{kv_get, kv_set}; -pub use types::{Flow, FlowRun, FlowRunStep, FlowRunTrigger}; +// `upsert_flow_run_step` is likewise re-exported for the tinyflows seam: the +// live run observer (`tinyflows::observability::FlowRunObserver`, issue G2) +// lives in the sibling `tinyflows` domain and persists each finished step onto +// the `flow_runs` row through this function as the run executes. +pub use store::{kv_get, kv_set, upsert_flow_run_step}; +pub use types::{Flow, FlowRun, FlowRunStep, FlowRunTrigger, FlowValidation}; diff --git a/src/openhuman/flows/ops.rs b/src/openhuman/flows/ops.rs index ecd4c9cfe..e289125b2 100644 --- a/src/openhuman/flows/ops.rs +++ b/src/openhuman/flows/ops.rs @@ -12,6 +12,7 @@ use tinyflows::model::{TriggerKind, WorkflowGraph}; use crate::openhuman::agent::turn_origin::{with_origin, AgentTurnOrigin, TrustedAutomationSource}; use crate::openhuman::config::Config; use crate::openhuman::flows::bus; +use crate::openhuman::flows::run_registry; use crate::openhuman::flows::store; use crate::openhuman::flows::types::{FlowRunStep, FlowRunTrigger}; use crate::openhuman::flows::{Flow, FlowRun}; @@ -22,6 +23,16 @@ use crate::rpc::RpcOutcome; /// call must never let the RPC block indefinitely — this caps the whole run. const FLOW_RUN_TIMEOUT_SECS: u64 = 600; +/// How long a run may sit parked at a human-in-the-loop approval gate +/// (`pending_approval`) before the TTL sweep expires it to a terminal +/// `"cancelled"` (issue G4). Aligned with the agent tool-call `ApprovalGate`'s +/// 10-minute fail-closed TTL (`src/openhuman/approval/`), so a flow HITL gate a +/// human never answers doesn't wedge a run — and its durable checkpoint — +/// forever. The two are distinct mechanisms (flow runs execute as +/// `TrustedAutomation { Workflow }`, which the tool-call gate lets through), so +/// this is a dedicated flows-side TTL, not a reuse of the approval store's. +const FLOW_PARKED_TTL_SECS: i64 = 600; + /// Runs a raw graph JSON value through `tinyflows::migrate::migrate` (upgrade /// an older-schema definition to current), deserializes it, and rejects a /// structurally invalid graph via `tinyflows::validate::validate` — so a bad @@ -40,6 +51,122 @@ pub(crate) fn validate_and_migrate_graph(graph_json: Value) -> Result &'static str { + match kind { + TriggerKind::Manual => "manual", + TriggerKind::Schedule => "schedule", + TriggerKind::Webhook => "webhook", + TriggerKind::AppEvent => "app_event", + TriggerKind::Form => "form", + TriggerKind::ExecuteByWorkflow => "execute_by_workflow", + TriggerKind::ChatMessage => "chat_message", + TriggerKind::Evaluation => "evaluation", + TriggerKind::System => "system", + } +} + +/// Whether a flow's trigger kind currently produces *automatic* runs in this +/// host. Only three kinds fire today: +/// - `manual` — runnable on demand via `flows_run` (no automatic dispatch, but +/// that's the whole contract of a manual trigger — never a surprise). +/// - `schedule` — a `cron` job drives `FlowScheduleTick` (see +/// [`bind_schedule_trigger`]). +/// - `app_event` — matched against `ComposioTriggerReceived` at dispatch time +/// (see `flows::bus::FlowTriggerSubscriber`). +/// +/// Everything else (`webhook`, `chat_message`, `form`, `execute_by_workflow`, +/// `evaluation`, `system`) is *accepted and saved* but has no wired dispatch +/// path yet — enabling such a flow silently produces a flow that never runs +/// itself. [`graph_trigger_warnings`] turns that silence into a loud warning. +fn trigger_kind_fires(kind: &TriggerKind) -> bool { + matches!( + kind, + TriggerKind::Manual | TriggerKind::Schedule | TriggerKind::AppEvent + ) +} + +/// Produces host-side, **non-fatal** validation warnings for a graph — today +/// exactly one: "this trigger kind does not fire automatically yet". Returns +/// an empty vec when the trigger fires (`manual`/`schedule`/`app_event`), when +/// the graph has no single resolvable trigger node, or when the trigger has no +/// `trigger_kind` discriminator (a legacy/manual-only graph authored before +/// B2 simply never self-fires — not a warnable surprise, matching +/// `bus::extract_trigger_kind`'s "no automatic binding" treatment). +/// +/// This lives host-side (NOT in `tinyflows::validate`, which is host-agnostic +/// and only does structural checks) because "which trigger kinds this host has +/// wired" is an OpenHuman fact, not a property of the portable graph. +pub(crate) fn graph_trigger_warnings(graph: &WorkflowGraph) -> Vec { + let Some(trigger) = graph.trigger() else { + return Vec::new(); + }; + let Some(kind_value) = trigger.config.get("trigger_kind") else { + return Vec::new(); + }; + let kind: TriggerKind = match serde_json::from_value(kind_value.clone()) { + Ok(k) => k, + Err(_) => return Vec::new(), + }; + if trigger_kind_fires(&kind) { + return Vec::new(); + } + let label = trigger_kind_label(&kind); + vec![format!( + "Trigger kind '{label}' does not fire automatically yet — this flow will be saved and \ + can be enabled, but nothing will run it on its own until that trigger is wired up. Run \ + it manually with flows_run, or switch to a `schedule` or `app_event` trigger." + )] +} + +/// 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 +/// ([`graph_trigger_warnings`]). Backs `openhuman.flows_validate` (PHASE 3c): +/// an authoring surface can call this to preview validity + warnings before a +/// save. Pure (no persistence, no config) — `valid == false` is a normal +/// result, NOT an `Err`; `Err` is reserved for internal serialization faults +/// (there are none on this path today). +pub fn flows_validate(graph_json: Value) -> RpcOutcome { + use crate::openhuman::flows::FlowValidation; + tracing::debug!(target: "flows", "[flows] flows_validate: validating candidate graph"); + match validate_and_migrate_graph(graph_json) { + Ok(graph) => { + let warnings = graph_trigger_warnings(&graph); + for warning in &warnings { + tracing::warn!(target: "flows", warning = %warning, "[flows] flows_validate: non-fatal validation warning"); + } + tracing::debug!( + target: "flows", + node_count = graph.nodes.len(), + warning_count = warnings.len(), + "[flows] flows_validate: graph is structurally valid" + ); + RpcOutcome::single_log( + FlowValidation { + valid: true, + errors: Vec::new(), + warnings, + }, + "flow validated", + ) + } + Err(error) => { + tracing::debug!(target: "flows", %error, "[flows] flows_validate: graph is structurally invalid"); + RpcOutcome::single_log( + FlowValidation { + valid: false, + errors: vec![error], + warnings: Vec::new(), + }, + "flow validation failed", + ) + } + } +} + /// Creates a new flow from a name and a raw graph JSON value. /// /// `store::create_flow` defaults new flows to `enabled = true` — this binds @@ -195,10 +322,24 @@ pub async fn flows_set_enabled( unbind_trigger(config, &flow); } - Ok(RpcOutcome::single_log( - flow, - format!("flow {id} enabled={enabled}"), - )) + let mut logs = vec![format!("flow {id} enabled={enabled}")]; + // When enabling, loudly surface any unfired-trigger-kind warning in the + // result (a structured `warning:`-prefixed log), not just a silent tracing + // line — so an enable of a flow that will never fire itself (webhook, + // chat_message, form, …) is impossible to miss at the call site. + if enabled { + for warning in graph_trigger_warnings(&flow.graph) { + tracing::warn!( + target: "flows", + flow_id = %id, + warning = %warning, + "[flows] flows_set_enabled: enabling a flow whose trigger kind does not fire yet" + ); + logs.push(format!("warning: {warning}")); + } + } + + Ok(RpcOutcome::new(flow, logs)) } /// Registers the automatic-dispatch side effect for `flow`'s trigger kind, if @@ -444,9 +585,15 @@ pub async fn flows_run( start_flow_run_row(config, &thread_id, flow_id); + // Register this run as in-flight (issue G4) so a concurrent + // `flows_cancel_run` can signal it to abort. The guard deregisters on any + // exit from this fn (including the early returns below). + let (cancel_token, _run_guard) = run_registry::register(&thread_id); + // Record a failed attempt so `last_run_at`/`last_status` reflect reality // (a stop-policy engine/capability failure or a timeout) rather than - // leaving the prior success/pending state on the flow. + // leaving the prior success/pending state on the flow. Preserve whatever + // steps the observer persisted live (don't wipe them back to `[]`). let record_failed = |error: &str| { if let Err(rec_err) = store::record_run(config, flow_id, "failed") { tracing::warn!( @@ -456,7 +603,8 @@ pub async fn flows_run( "[flows] flows_run: failed to record failed run" ); } - finish_flow_run_row(config, &thread_id, "failed", &[], &[], Some(error)); + let observed = current_persisted_steps(config, &thread_id); + finish_flow_run_row(config, &thread_id, "failed", &observed, &[], Some(error)); }; let origin = workflow_origin(flow_id, flow.require_approval); @@ -465,35 +613,67 @@ pub async fn flows_run( // post-run Langfuse export reads back. Process-local and dropped with the // run — never persisted. let journal = Arc::new(tinyflows::engine::InMemoryGraphEventJournal::new()); + // Live run observer (issue G2): persists each finished step into the + // `flow_runs` row as it happens and streams a `FlowRunProgress` event to + // the frontend, so the durable + journaled path also reports live. + let observer: Arc = Arc::new( + crate::openhuman::tinyflows::observability::FlowRunObserver::new( + Arc::new(config.clone()), + flow_id, + thread_id.clone(), + ), + ); let run = with_origin( origin, - tinyflows::engine::run_with_checkpointer_journaled( + tinyflows::engine::run_with_checkpointer_journaled_observed( &compiled, input, &caps, checkpointer, &thread_id, journal.clone(), + &observer, ), ); - let journaled = match tokio::time::timeout( - std::time::Duration::from_secs(FLOW_RUN_TIMEOUT_SECS), - run, - ) - .await - { - Ok(Ok(journaled)) => journaled, - Ok(Err(e)) => { - record_failed(&e.to_string()); - tracing::warn!(target: "flows", flow_id = %flow_id, error = %e, "[flows] flows_run: run failed"); - return Err(e.to_string()); - } - Err(_elapsed) => { - let msg = format!("flow run timed out after {FLOW_RUN_TIMEOUT_SECS}s"); - record_failed(&msg); - tracing::warn!(target: "flows", flow_id = %flow_id, timeout_secs = FLOW_RUN_TIMEOUT_SECS, "[flows] flows_run: run timed out"); - return Err(msg); + let timed = tokio::time::timeout(std::time::Duration::from_secs(FLOW_RUN_TIMEOUT_SECS), run); + tokio::pin!(timed); + // Race the run against a cancellation signal (issue G4). `biased` checks the + // cancel arm first so a `flows_cancel_run` that lands right as the run + // settles still wins deterministically. + let journaled = tokio::select! { + biased; + _ = cancel_token.cancelled() => { + tracing::info!(target: "flows", flow_id = %flow_id, thread_id = %thread_id, "[flows] flows_run: cancelled mid-run"); + if let Err(e) = store::record_run(config, flow_id, "cancelled") { + tracing::warn!(target: "flows", flow_id = %flow_id, error = %e, "[flows] flows_run: failed to record cancelled run"); + } + let observed = current_persisted_steps(config, &thread_id); + finish_flow_run_row(config, &thread_id, "cancelled", &observed, &[], Some("run cancelled")); + drop_checkpoint(config, &thread_id).await; + return Ok(RpcOutcome::single_log( + json!({ + "output": Value::Null, + "pending_approvals": Vec::::new(), + "thread_id": thread_id, + "cancelled": true, + }), + format!("flow run cancelled: {thread_id}"), + )); } + result = &mut timed => match result { + Ok(Ok(journaled)) => journaled, + Ok(Err(e)) => { + record_failed(&e.to_string()); + tracing::warn!(target: "flows", flow_id = %flow_id, error = %e, "[flows] flows_run: run failed"); + return Err(e.to_string()); + } + Err(_elapsed) => { + let msg = format!("flow run timed out after {FLOW_RUN_TIMEOUT_SECS}s"); + record_failed(&msg); + tracing::warn!(target: "flows", flow_id = %flow_id, timeout_secs = FLOW_RUN_TIMEOUT_SECS, "[flows] flows_run: run timed out"); + return Err(msg); + } + }, }; let outcome = journaled.outcome; @@ -507,7 +687,7 @@ pub async fn flows_run( config, &thread_id, status, - &reconstruct_steps(&outcome.output), + &settle_steps(config, &thread_id, &outcome.output), &outcome.pending_approvals, None, ); @@ -564,6 +744,7 @@ pub async fn flows_resume( flow_id: &str, thread_id: &str, approvals: Vec, + rejections: Vec, ) -> Result, String> { let flow = store::get_flow(config, flow_id) .map_err(|e| e.to_string())? @@ -586,8 +767,21 @@ pub async fn flows_resume( run_record.status )); } + // A gate can't be both approved and denied in the same resume — that's an + // ambiguous instruction, reject it up front. + if let Some(dup) = approvals.iter().find(|a| rejections.contains(a)) { + return Err(format!( + "gate '{dup}' cannot be both approved and rejected in the same resume" + )); + } + // Same host-side guard the approvals path uses (see this fn's doc): the + // engine trusts whatever the resume delivers, so require that the caller's + // approvals/rejections actually name a currently-pending gate before ever + // touching the engine. A denial (issue G4) is enforced the same way — a + // rejection naming a pending gate is a valid resume just as an approval is. let matches_pending = approvals .iter() + .chain(rejections.iter()) .any(|a| run_record.pending_approvals.contains(a)); if !matches_pending { tracing::warn!( @@ -595,12 +789,13 @@ pub async fn flows_resume( flow_id = %flow_id, %thread_id, ?approvals, + ?rejections, pending = ?run_record.pending_approvals, - "[flows] flows_resume: rejected — caller approvals name none of the pending gates" + "[flows] flows_resume: rejected — caller approvals/rejections name none of the pending gates" ); return Err(format!( - "no pending approval matches: approvals {approvals:?} do not name any of the \ - currently pending gates {:?} for run '{thread_id}'", + "no pending approval matches: approvals {approvals:?} / rejections {rejections:?} do \ + not name any of the currently pending gates {:?} for run '{thread_id}'", run_record.pending_approvals )); } @@ -617,6 +812,7 @@ pub async fn flows_resume( flow_id = %flow_id, %thread_id, approval_count = approvals.len(), + rejection_count = rejections.len(), "[flows] flows_resume: resuming checkpointed run" ); @@ -624,15 +820,30 @@ pub async fn flows_resume( // Same per-run journal as `flows_run`: the resumed execution mints a new // tinyagents run id, so its observation slice is read under that id. let journal = Arc::new(tinyflows::engine::InMemoryGraphEventJournal::new()); + // Live observer (issue G2): the resumed run fires `on_step_finish` for each + // node that runs after the interrupt boundary, so downstream steps are + // persisted + streamed live too, keyed by the same `thread_id`/run row. + let observer: Arc = Arc::new( + crate::openhuman::tinyflows::observability::FlowRunObserver::new( + Arc::new(config.clone()), + flow_id, + thread_id.to_string(), + ), + ); + // `rejections` (issue G4 — deny semantics): a denied gate routes to its + // `error` port (recovery branch) or, if it has none, fails the run. The + // empty-rejections case is byte-for-byte the prior approve-only resume. let run = with_origin( origin, - tinyflows::engine::resume_with_checkpointer_journaled( + tinyflows::engine::resume_with_checkpointer_journaled_observed( &compiled, &caps, checkpointer, thread_id, approvals, + rejections, journal.clone(), + &observer, ), ); @@ -645,14 +856,23 @@ pub async fn flows_resume( Ok(Ok(journaled)) => journaled, Ok(Err(e)) => { let _ = store::record_run(config, flow_id, "failed"); - finish_flow_run_row(config, thread_id, "failed", &[], &[], Some(&e.to_string())); + let observed = current_persisted_steps(config, thread_id); + finish_flow_run_row( + config, + thread_id, + "failed", + &observed, + &[], + Some(&e.to_string()), + ); tracing::warn!(target: "flows", flow_id = %flow_id, %thread_id, error = %e, "[flows] flows_resume: run failed"); return Err(e.to_string()); } Err(_elapsed) => { let msg = format!("flow resume timed out after {FLOW_RUN_TIMEOUT_SECS}s"); let _ = store::record_run(config, flow_id, "failed"); - finish_flow_run_row(config, thread_id, "failed", &[], &[], Some(&msg)); + let observed = current_persisted_steps(config, thread_id); + finish_flow_run_row(config, thread_id, "failed", &observed, &[], Some(&msg)); tracing::warn!(target: "flows", flow_id = %flow_id, %thread_id, timeout_secs = FLOW_RUN_TIMEOUT_SECS, "[flows] flows_resume: run timed out"); return Err(msg); } @@ -669,7 +889,7 @@ pub async fn flows_resume( config, thread_id, status, - &reconstruct_steps(&outcome.output), + &settle_steps(config, thread_id, &outcome.output), &outcome.pending_approvals, None, ); @@ -706,12 +926,15 @@ pub async fn flows_resume( } /// Lists the most recent runs for a flow (newest first), for the B3 -/// run-history inspector. +/// run-history inspector. Runs a lazy parked-run TTL sweep first (see +/// [`sweep_expired_parked_runs`]) so the listing reflects any run that has now +/// aged out of `pending_approval`. pub async fn flows_list_runs( config: &Config, flow_id: &str, limit: usize, ) -> Result>, String> { + sweep_expired_parked_runs(config).await; let runs = store::list_flow_runs(config, flow_id, limit).map_err(|e| e.to_string())?; Ok(RpcOutcome::single_log( runs, @@ -719,8 +942,11 @@ pub async fn flows_list_runs( )) } -/// Loads a single flow run record by id (== `thread_id`). +/// Loads a single flow run record by id (== `thread_id`). Runs the lazy +/// parked-run TTL sweep first so a stale parked run is reported as `cancelled` +/// rather than perpetually `pending_approval`. pub async fn flows_get_run(config: &Config, run_id: &str) -> Result, String> { + sweep_expired_parked_runs(config).await; let run = store::get_flow_run(config, run_id) .map_err(|e| e.to_string())? .ok_or_else(|| format!("flow run '{run_id}' not found"))?; @@ -730,6 +956,133 @@ pub async fn flows_get_run(config: &Config, run_id: &str) -> Result usize { + let now = Utc::now(); + let cutoff = (now - chrono::Duration::seconds(FLOW_PARKED_TTL_SECS)).to_rfc3339(); + let now_str = now.to_rfc3339(); + let error_msg = format!("parked run expired after {FLOW_PARKED_TTL_SECS}s awaiting approval"); + + let swept = match store::expire_parked_runs(config, &cutoff, &now_str, &error_msg) { + Ok(swept) => swept, + Err(e) => { + tracing::warn!(target: "flows", error = %e, "[flows] parked-run TTL sweep failed (read continues)"); + return 0; + } + }; + for (run_id, flow_id) in &swept { + if let Err(e) = store::record_run(config, flow_id, "cancelled") { + tracing::warn!(target: "flows", run_id, flow_id, error = %e, "[flows] TTL sweep: failed to update flow summary for expired run"); + } + drop_checkpoint(config, run_id).await; + } + if !swept.is_empty() { + tracing::info!(target: "flows", count = swept.len(), ttl_secs = FLOW_PARKED_TTL_SECS, "[flows] parked-run TTL sweep expired stale runs"); + } + swept.len() +} + +/// Cancels a flow run (issue G4), settling it to a terminal `"cancelled"` +/// status and dropping its durable checkpoint so the aborted thread can never +/// be resumed. +/// +/// Two cases, distinguished by [`run_registry::cancel`]: +/// - **In-flight** (a `flows_run` / `flows_resume` currently executing its run +/// future): the token is signalled and that run's own cancellation arm writes +/// the terminal row + drops the checkpoint as it unwinds — we don't write the +/// row here, to avoid two writers racing the same `flow_runs` row. +/// - **Parked / stale** (a `pending_approval` run awaiting a human decision, or +/// a `running` row whose task is gone): no live task exists to unwind, so +/// this settles the row terminally itself and drops the checkpoint. +/// +/// A run that is already terminal (`completed` / `failed` / `cancelled`) is a +/// clear error, not a silent no-op. +pub async fn flows_cancel_run(config: &Config, run_id: &str) -> Result, String> { + let run = store::get_flow_run(config, run_id) + .map_err(|e| e.to_string())? + .ok_or_else(|| format!("flow run '{run_id}' not found"))?; + + if matches!(run.status.as_str(), "completed" | "failed" | "cancelled") { + return Err(format!( + "flow run '{run_id}' is already terminal (status: {}) — nothing to cancel", + run.status + )); + } + + let signalled = run_registry::cancel(run_id); + tracing::info!( + target: "flows", + run_id, + flow_id = %run.flow_id, + signalled, + prior_status = %run.status, + "[flows] flows_cancel_run: cancelling run" + ); + + if signalled { + // The in-flight run's cancellation arm owns the terminal write + the + // checkpoint drop; we've signalled it and return. Its settle is + // eventual (the run future unwinds), so report "requested". + return Ok(RpcOutcome::single_log( + json!({ "run_id": run_id, "cancelled": true, "was_in_flight": true }), + format!("flow run {run_id} cancellation requested"), + )); + } + + // Not in flight: settle the row terminally and drop the checkpoint here. + if let Err(e) = store::record_run(config, &run.flow_id, "cancelled") { + tracing::warn!(target: "flows", run_id, flow_id = %run.flow_id, error = %e, "[flows] flows_cancel_run: failed to record cancelled status on flow summary"); + } + let observed = current_persisted_steps(config, run_id); + finish_flow_run_row( + config, + run_id, + "cancelled", + &observed, + &[], + Some("run cancelled"), + ); + drop_checkpoint(config, run_id).await; + + Ok(RpcOutcome::single_log( + json!({ "run_id": run_id, "cancelled": true, "was_in_flight": false }), + format!("flow run {run_id} cancelled"), + )) +} + +/// Best-effort drop of a run's durable tinyagents checkpoint thread, so a +/// cancelled (or expired) run can never be resumed from its persisted interrupt +/// boundary. Logged, never fatal — the `flow_runs` row's terminal status is the +/// authoritative "not resumable" signal (the `flows_resume` guard already +/// rejects any non-`pending_approval` status); dropping the checkpoint is +/// belt-and-suspenders that also reclaims the storage. +async fn drop_checkpoint(config: &Config, thread_id: &str) { + use tinyflows::engine::Checkpointer as _; + match crate::openhuman::tinyflows::open_flow_checkpointer(config) { + Ok(checkpointer) => match checkpointer.delete_thread(thread_id).await { + Ok(()) => { + tracing::debug!(target: "flows", thread_id, "[flows] dropped durable checkpoint for cancelled/expired run") + } + Err(e) => { + tracing::warn!(target: "flows", thread_id, error = %e, "[flows] failed to drop durable checkpoint") + } + }, + Err(e) => { + tracing::warn!(target: "flows", thread_id, error = %e, "[flows] could not open checkpointer to drop checkpoint"); + } + } +} + /// Builds the `TrustedAutomation { Workflow }` origin scoped around every /// `flows_run` / `flows_resume` invocation. See `flows_run`'s doc for why /// this applies uniformly regardless of caller. @@ -775,14 +1128,14 @@ fn finish_flow_run_row( } /// Reconstructs a lean per-node step list from a settled run's -/// `output["nodes"]` map. tinyflows 0.2's durable path installs a -/// `NoopObserver` (see `tinyflows/observability.rs`), so there is no live -/// step stream to persist — this is the B2 "good enough" substitute the -/// spec calls for; a richer per-step `RunObserver` is a tinyflows 0.3 item. +/// `output["nodes"]` map. /// -/// // TODO(0.3): replace this reconstruction with a real `RunObserver` that -/// // streams `node_id`/`status`/`output`/`duration_ms` as each node -/// // finishes, once the durable run path supports installing one. +/// As of issue G2 (live run observation) this is no longer the primary source +/// of run steps — `flows::observability::FlowRunObserver` persists each step +/// live as it finishes (with real `status`/`duration_ms`). This reconstruction +/// is now only a **fallback**, used by [`settle_steps`] to fill in any node the +/// observer didn't emit an `on_step_finish` for (notably the trigger node), +/// and as the whole-run source when the observer saw nothing at all. fn reconstruct_steps(output: &Value) -> Vec { let Some(nodes) = output.get("nodes").and_then(Value::as_object) else { return Vec::new(); @@ -793,10 +1146,64 @@ fn reconstruct_steps(output: &Value) -> Vec { node_id: node_id.clone(), output: slot.get("items").cloned().unwrap_or(Value::Null), port: slot.get("port").and_then(Value::as_str).map(str::to_string), + // Reconstructed post-hoc: no live status/timing (see FlowRunStep). + status: None, + duration_ms: None, }) .collect() } +/// Reads back whatever steps the live [`FlowRunObserver`] has already persisted +/// onto the run's row. Best-effort: a read failure yields an empty list (the +/// caller still writes a terminal row), never propagating an error into the +/// run's settle path. +/// +/// [`FlowRunObserver`]: crate::openhuman::tinyflows::observability::FlowRunObserver +fn current_persisted_steps(config: &Config, run_id: &str) -> Vec { + store::get_flow_run(config, run_id) + .ok() + .flatten() + .map(|run| run.steps) + .unwrap_or_default() +} + +/// Assembles the final step list to persist at settle: the live steps the +/// observer already recorded (carrying real `status`/`duration_ms`), plus any +/// node present in the post-hoc [`reconstruct_steps`] projection that the +/// observer never emitted a step for — the trigger node, or (defensively) an +/// observer that missed a step. If the observer recorded nothing at all +/// (e.g. a run that paused immediately at a gate before any node finished), +/// falls back wholesale to the reconstruction. +fn settle_steps(config: &Config, run_id: &str, output: &Value) -> Vec { + let reconstructed = reconstruct_steps(output); + let persisted = current_persisted_steps(config, run_id); + if persisted.is_empty() { + tracing::debug!( + target: "flows", + run_id, + reconstructed = reconstructed.len(), + "[flows] settle_steps: no live-observed steps — using post-hoc reconstruction" + ); + return reconstructed; + } + let mut merged = persisted; + let mut filled = 0usize; + for step in reconstructed { + if !merged.iter().any(|s| s.node_id == step.node_id) { + merged.push(step); + filled += 1; + } + } + tracing::debug!( + target: "flows", + run_id, + step_count = merged.len(), + filled_from_reconstruction = filled, + "[flows] settle_steps: merged live-observed steps with post-hoc reconstruction" + ); + merged +} + /// Milliseconds since the Unix epoch, for `CoreNotificationEvent::timestamp_ms`. fn now_ms() -> u64 { std::time::SystemTime::now() diff --git a/src/openhuman/flows/ops_tests.rs b/src/openhuman/flows/ops_tests.rs index bd5ed76f7..47c2bd619 100644 --- a/src/openhuman/flows/ops_tests.rs +++ b/src/openhuman/flows/ops_tests.rs @@ -475,7 +475,7 @@ async fn flows_resume_continues_a_paused_run_to_completion() { serde_json::from_value(run.value["pending_approvals"].clone()).unwrap(); assert_eq!(pending, vec!["gate".to_string()]); - let resumed = flows_resume(&config, &created.value.id, &thread_id, pending) + let resumed = flows_resume(&config, &created.value.id, &thread_id, pending, vec![]) .await .unwrap(); assert_eq!(resumed.value["pending_approvals"], json!([])); @@ -506,7 +506,7 @@ async fn flows_resume_continues_a_paused_run_to_completion() { async fn flows_resume_missing_flow_errors() { let tmp = TempDir::new().unwrap(); let config = test_config(&tmp); - let err = flows_resume(&config, "missing", "thread-1", vec![]) + let err = flows_resume(&config, "missing", "thread-1", vec![], vec![]) .await .expect_err("must error"); assert!(err.contains("not found")); @@ -541,7 +541,7 @@ async fn flows_resume_with_empty_approvals_is_rejected_and_does_not_complete_the .unwrap(); let thread_id = run.value["thread_id"].as_str().unwrap().to_string(); - let err = flows_resume(&config, &created.value.id, &thread_id, vec![]) + let err = flows_resume(&config, &created.value.id, &thread_id, vec![], vec![]) .await .expect_err("an empty approvals list must not silently approve the pending gate"); assert!( @@ -586,6 +586,7 @@ async fn flows_resume_with_mismatched_approvals_is_rejected() { &created.value.id, &thread_id, vec!["not-a-real-gate".to_string()], + vec![], ) .await .expect_err("approvals naming no actually-pending gate must be rejected"); @@ -615,6 +616,7 @@ async fn flows_resume_with_the_correct_gate_completes_and_runs_downstream() { &created.value.id, &thread_id, vec!["gate".to_string()], + vec![], ) .await .unwrap(); @@ -628,6 +630,149 @@ async fn flows_resume_with_the_correct_gate_completes_and_runs_downstream() { assert_eq!(reloaded.value.last_status.as_deref(), Some("completed")); } +// ── flows_resume deny semantics (issue G4) ──────────────────────────────── + +/// A gate with BOTH a `main` edge (to `downstream`) and an `error` edge (to +/// `recover`): denying the gate routes to `recover`, not `downstream`. +fn approval_gated_graph_with_error_port() -> Value { + json!({ + "name": "approval-gated-error-port", + "nodes": [ + { "id": "t", "kind": "trigger", "name": "Trigger" }, + { "id": "gate", "kind": "output_parser", "name": "Gate", "config": { "requires_approval": true } }, + { "id": "downstream", "kind": "output_parser", "name": "Downstream" }, + { "id": "recover", "kind": "output_parser", "name": "Recover" } + ], + "edges": [ + { "from_node": "t", "to_node": "gate" }, + { "from_node": "gate", "from_port": "main", "to_node": "downstream" }, + { "from_node": "gate", "from_port": "error", "to_node": "recover" } + ] + }) +} + +#[tokio::test] +async fn flows_resume_denying_a_gate_routes_to_its_error_port() { + let tmp = TempDir::new().unwrap(); + let config = test_config(&tmp); + let created = flows_create( + &config, + "gated-deny".to_string(), + approval_gated_graph_with_error_port(), + false, + ) + .await + .unwrap(); + + let run = flows_run( + &config, + &created.value.id, + json!({ "x": 1 }), + FlowRunTrigger::Rpc, + ) + .await + .unwrap(); + let thread_id = run.value["thread_id"].as_str().unwrap().to_string(); + + // Deny the gate: no approvals, `gate` in rejections. + let resumed = flows_resume( + &config, + &created.value.id, + &thread_id, + vec![], + vec!["gate".to_string()], + ) + .await + .unwrap(); + + assert_eq!(resumed.value["pending_approvals"], json!([])); + assert_eq!( + resumed.value["output"]["nodes"]["recover"]["items"][0]["json"]["error"]["node"], + json!("gate"), + "a denied gate must route its error item to the `error`-port recovery node" + ); + assert!( + resumed.value["output"]["nodes"]["downstream"].is_null(), + "the main branch must not run when the gate is denied" + ); + + let reloaded = flows_get(&config, &created.value.id).await.unwrap(); + assert_eq!(reloaded.value.last_status.as_deref(), Some("completed")); + + let run_row = flows_get_run(&config, &thread_id).await.unwrap(); + assert_eq!(run_row.value.status, "completed"); + assert!(run_row.value.pending_approvals.is_empty()); +} + +#[tokio::test] +async fn flows_resume_denying_a_gate_with_no_error_port_fails_the_run() { + let tmp = TempDir::new().unwrap(); + let config = test_config(&tmp); + // `approval_gated_graph()` has only a `main` edge out of the gate — no + // `error` port to route a denial to, so the whole run must fail. + let created = flows_create(&config, "gated".to_string(), approval_gated_graph(), false) + .await + .unwrap(); + + let run = flows_run( + &config, + &created.value.id, + json!({ "x": 1 }), + FlowRunTrigger::Rpc, + ) + .await + .unwrap(); + let thread_id = run.value["thread_id"].as_str().unwrap().to_string(); + + let err = flows_resume( + &config, + &created.value.id, + &thread_id, + vec![], + vec!["gate".to_string()], + ) + .await + .expect_err("denying a gate with no error port must fail the run"); + assert!( + err.contains("denied"), + "expected a denial error, got: {err}" + ); + + let reloaded = flows_get(&config, &created.value.id).await.unwrap(); + assert_eq!(reloaded.value.last_status.as_deref(), Some("failed")); + let run_row = flows_get_run(&config, &thread_id).await.unwrap(); + assert_eq!(run_row.value.status, "failed"); +} + +#[tokio::test] +async fn flows_resume_rejects_a_gate_named_in_both_approvals_and_rejections() { + let tmp = TempDir::new().unwrap(); + let config = test_config(&tmp); + let created = flows_create(&config, "gated".to_string(), approval_gated_graph(), false) + .await + .unwrap(); + + let run = flows_run(&config, &created.value.id, json!({}), FlowRunTrigger::Rpc) + .await + .unwrap(); + let thread_id = run.value["thread_id"].as_str().unwrap().to_string(); + + let err = flows_resume( + &config, + &created.value.id, + &thread_id, + vec!["gate".to_string()], + vec!["gate".to_string()], + ) + .await + .expect_err("a gate cannot be both approved and rejected"); + assert!(err.contains("cannot be both approved and rejected")); + + // The run must be untouched (still pending), never half-resumed. + let run_row = flows_get_run(&config, &thread_id).await.unwrap(); + assert_eq!(run_row.value.status, "pending_approval"); +} + #[tokio::test] async fn flows_resume_of_a_non_paused_run_errors_clearly() { let tmp = TempDir::new().unwrap(); @@ -643,7 +788,7 @@ async fn flows_resume_of_a_non_paused_run_errors_clearly() { .unwrap(); let thread_id = run.value["thread_id"].as_str().unwrap().to_string(); - let err = flows_resume(&config, &created.value.id, &thread_id, vec![]) + let err = flows_resume(&config, &created.value.id, &thread_id, vec![], vec![]) .await .expect_err("resuming an already-completed run must be a clear error, not a silent no-op"); assert!( @@ -665,6 +810,7 @@ async fn flows_resume_with_no_recorded_run_for_thread_id_errors_clearly() { &created.value.id, "thread-that-was-never-started", vec![], + vec![], ) .await .expect_err("must error when no run is recorded for this thread_id"); @@ -808,3 +954,406 @@ async fn flows_run_does_not_notify_when_run_completes_without_pending_approvals( "a fully-completed run must not publish a pending-approval notification" ); } + +// ── Live run observation (issue G2) ─────────────────────────────────────── + +use crate::openhuman::tinyflows::observability::FlowRunObserver; +use std::sync::Arc as StdArc; +// `RunObserver` must be in scope to call `on_step_finish` on the observer. +use tinyflows::observability::{ExecutionStep, RunObserver as _, StepStatus}; + +/// trigger -> output_parser passthrough: the parser is a non-trigger node, so +/// the engine fires `on_step_finish` for it, exercising live persistence. +fn passthrough_graph() -> Value { + json!({ + "name": "passthrough", + "nodes": [ + { "id": "t", "kind": "trigger", "name": "Trigger" }, + { "id": "p", "kind": "output_parser", "name": "Parse" } + ], + "edges": [ { "from_node": "t", "to_node": "p" } ] + }) +} + +#[tokio::test] +async fn observer_persists_each_step_incrementally() { + // The observer no-ops until the run's start row exists (mirrors + // `start_flow_run_row`), so seed a flow + a running run row first. + let tmp = TempDir::new().unwrap(); + let config = test_config(&tmp); + let created = flows_create(&config, "obs".to_string(), passthrough_graph(), false) + .await + .unwrap(); + let run_id = format!("flow:{}:run-under-test", created.value.id); + store::insert_flow_run( + &config, + &run_id, + &created.value.id, + &run_id, + "2026-01-01T00:00:00Z", + ) + .unwrap(); + + let observer = FlowRunObserver::new( + StdArc::new(config.clone()), + created.value.id.clone(), + &run_id, + ); + observer.on_step_finish(&ExecutionStep { + node_id: "a".to_string(), + status: StepStatus::Success, + output: json!([{ "json": { "ok": true } }]), + duration_ms: 7, + }); + observer.on_step_finish(&ExecutionStep { + node_id: "b".to_string(), + status: StepStatus::Error, + output: Value::Null, + duration_ms: 3, + }); + + // The store now holds both live steps with real status + timing — proof of + // incremental persistence (post-hoc reconstruction leaves status None). + let row = store::get_flow_run(&config, &run_id).unwrap().unwrap(); + assert_eq!(row.steps.len(), 2, "both live steps should be persisted"); + let a = row.steps.iter().find(|s| s.node_id == "a").unwrap(); + assert_eq!(a.status.as_deref(), Some("success")); + assert_eq!(a.duration_ms, Some(7)); + let b = row.steps.iter().find(|s| s.node_id == "b").unwrap(); + assert_eq!(b.status.as_deref(), Some("error")); + assert_eq!(b.duration_ms, Some(3)); + + // Re-firing the same node id replaces its entry rather than duplicating it. + observer.on_step_finish(&ExecutionStep { + node_id: "a".to_string(), + status: StepStatus::Success, + output: json!([{ "json": { "ok": true } }]), + duration_ms: 42, + }); + let row = store::get_flow_run(&config, &run_id).unwrap().unwrap(); + assert_eq!(row.steps.len(), 2, "re-firing a node must not duplicate it"); + let a = row.steps.iter().find(|s| s.node_id == "a").unwrap(); + assert_eq!( + a.duration_ms, + Some(42), + "the step should be replaced in place" + ); +} + +#[tokio::test] +async fn flows_run_persists_live_steps_with_status_and_timing() { + let tmp = TempDir::new().unwrap(); + let config = test_config(&tmp); + let created = flows_create( + &config, + "passthrough".to_string(), + passthrough_graph(), + false, + ) + .await + .unwrap(); + + let run = flows_run( + &config, + &created.value.id, + json!({ "x": 1 }), + FlowRunTrigger::Rpc, + ) + .await + .unwrap(); + let thread_id = run.value["thread_id"].as_str().unwrap().to_string(); + + let row = flows_get_run(&config, &thread_id).await.unwrap(); + assert_eq!(row.value.status, "completed"); + + // The non-trigger node 'p' was observed live: it carries a real status + + // timing that only the live observer (not post-hoc reconstruction) sets. + let p = row + .value + .steps + .iter() + .find(|s| s.node_id == "p") + .expect("the output_parser step should be persisted"); + assert_eq!(p.status.as_deref(), Some("success")); + assert!( + p.duration_ms.is_some(), + "a live-observed step should carry executor timing" + ); + + // The trigger node emits no `on_step_finish`; `settle_steps` fills it in + // from the post-hoc reconstruction, so it carries no live status. + let t = row + .value + .steps + .iter() + .find(|s| s.node_id == "t") + .expect("the trigger step should be reconstructed at settle"); + assert!( + t.status.is_none(), + "the trigger step is reconstructed post-hoc, not observed live" + ); +} + +// ── flows_cancel_run (issue G4) ─────────────────────────────────────────── + +#[tokio::test] +async fn flows_cancel_run_cancels_a_parked_pending_approval_run() { + let tmp = TempDir::new().unwrap(); + let config = test_config(&tmp); + let created = flows_create(&config, "gated".to_string(), approval_gated_graph(), false) + .await + .unwrap(); + + // Run pauses at the gate → a durable `pending_approval` row with no live + // task (the run future already returned): the not-in-flight cancel path. + let run = flows_run(&config, &created.value.id, json!({}), FlowRunTrigger::Rpc) + .await + .unwrap(); + let thread_id = run.value["thread_id"].as_str().unwrap().to_string(); + assert_eq!( + flows_get_run(&config, &thread_id) + .await + .unwrap() + .value + .status, + "pending_approval" + ); + + let cancelled = flows_cancel_run(&config, &thread_id).await.unwrap(); + assert_eq!(cancelled.value["cancelled"], json!(true)); + assert_eq!( + cancelled.value["was_in_flight"], + json!(false), + "a parked run has no live task, so the cancel settles the row directly" + ); + + // The run row and the flow summary both reach the terminal `cancelled`. + let run_row = flows_get_run(&config, &thread_id).await.unwrap(); + assert_eq!(run_row.value.status, "cancelled"); + assert!(run_row.value.pending_approvals.is_empty()); + assert_eq!(run_row.value.error.as_deref(), Some("run cancelled")); + + let reloaded = flows_get(&config, &created.value.id).await.unwrap(); + assert_eq!(reloaded.value.last_status.as_deref(), Some("cancelled")); + + // A cancelled run can no longer be resumed — the status guard rejects it. + let err = flows_resume( + &config, + &created.value.id, + &thread_id, + vec!["gate".to_string()], + vec![], + ) + .await + .expect_err("a cancelled run must not be resumable"); + assert!(err.contains("not pending approval") || err.contains("no paused run")); +} + +#[tokio::test] +async fn flows_cancel_run_of_an_already_completed_run_errors() { + let tmp = TempDir::new().unwrap(); + let config = test_config(&tmp); + let created = flows_create(&config, "demo".to_string(), trigger_only_graph(), false) + .await + .unwrap(); + + let run = flows_run(&config, &created.value.id, json!({}), FlowRunTrigger::Rpc) + .await + .unwrap(); + let thread_id = run.value["thread_id"].as_str().unwrap().to_string(); + + let err = flows_cancel_run(&config, &thread_id) + .await + .expect_err("cancelling an already-completed run must be a clear error"); + assert!(err.contains("already terminal"), "got: {err}"); +} + +#[tokio::test] +async fn flows_cancel_run_missing_run_errors() { + let tmp = TempDir::new().unwrap(); + let config = test_config(&tmp); + let err = flows_cancel_run(&config, "no-such-run") + .await + .expect_err("must error for an unknown run"); + assert!(err.contains("not found")); +} + +// ── parked-run TTL sweep (issue G4) ─────────────────────────────────────── + +#[tokio::test] +async fn parked_run_ttl_sweep_expires_stale_runs_but_spares_fresh_ones() { + let tmp = TempDir::new().unwrap(); + let config = test_config(&tmp); + let created = flows_create(&config, "gated".to_string(), approval_gated_graph(), false) + .await + .unwrap(); + + // Seed a parked run whose "parked since" (finished_at) is far in the past, + // so it is well beyond the TTL. + let stale_id = format!("flow:{}:stale-run", created.value.id); + let ancient = "2000-01-01T00:00:00+00:00"; + store::insert_flow_run(&config, &stale_id, &created.value.id, &stale_id, ancient).unwrap(); + store::finish_flow_run( + &config, + &stale_id, + "pending_approval", + ancient, + &[], + &["gate".to_string()], + None, + ) + .unwrap(); + + // A genuinely fresh parked run (just paused now) must survive the sweep. + let fresh = flows_run(&config, &created.value.id, json!({}), FlowRunTrigger::Rpc) + .await + .unwrap(); + let fresh_id = fresh.value["thread_id"].as_str().unwrap().to_string(); + + let swept = sweep_expired_parked_runs(&config).await; + assert_eq!(swept, 1, "only the stale parked run should be swept"); + + let stale_row = store::get_flow_run(&config, &stale_id).unwrap().unwrap(); + assert_eq!(stale_row.status, "cancelled"); + assert!( + stale_row.error.unwrap_or_default().contains("expired"), + "an expired run's error must note the TTL expiry" + ); + + let fresh_row = store::get_flow_run(&config, &fresh_id).unwrap().unwrap(); + assert_eq!( + fresh_row.status, "pending_approval", + "a run parked within the TTL must not be swept" + ); + + // The swept run is no longer resumable. + let err = flows_resume( + &config, + &created.value.id, + &stale_id, + vec!["gate".to_string()], + vec![], + ) + .await + .expect_err("an expired parked run must not be resumable"); + assert!(err.contains("not pending approval") || err.contains("no paused run")); +} + +// --------------------------------------------------------------------------- +// Unfired-trigger-kind warnings (PHASE 1a validation + PHASE 3c flows_validate) +// --------------------------------------------------------------------------- + +fn webhook_trigger_graph() -> Value { + json!({ + "name": "hooked", + "nodes": [ + { + "id": "t", + "kind": "trigger", + "name": "Trigger", + "config": { "trigger_kind": "webhook" } + } + ], + "edges": [] + }) +} + +#[test] +fn flows_validate_warns_on_unfired_webhook_trigger() { + let outcome = flows_validate(webhook_trigger_graph()); + assert!(outcome.value.valid, "a webhook graph is structurally valid"); + assert!(outcome.value.errors.is_empty()); + assert_eq!( + outcome.value.warnings.len(), + 1, + "an unfired webhook trigger must produce exactly one warning: {:?}", + outcome.value.warnings + ); + assert!( + outcome.value.warnings[0].contains("webhook") + && outcome.value.warnings[0].contains("does not fire"), + "warning must name the kind and explain it does not fire: {:?}", + outcome.value.warnings + ); +} + +#[test] +fn flows_validate_does_not_warn_on_schedule_trigger() { + let outcome = flows_validate(schedule_trigger_graph("0 9 * * *")); + assert!(outcome.value.valid); + assert!( + outcome.value.warnings.is_empty(), + "a schedule trigger fires — it must not warn: {:?}", + outcome.value.warnings + ); +} + +#[test] +fn flows_validate_reports_error_for_graph_without_trigger() { + let graph = json!({ + "name": "bad", + "nodes": [ { "id": "a", "kind": "output_parser", "name": "A" } ], + "edges": [] + }); + let outcome = flows_validate(graph); + assert!(!outcome.value.valid); + assert_eq!(outcome.value.errors.len(), 1); + assert!(outcome.value.errors[0].contains("trigger")); + assert!( + outcome.value.warnings.is_empty(), + "an invalid graph reports no warnings" + ); +} + +#[tokio::test] +async fn flows_set_enabled_surfaces_unfired_trigger_warning_at_enable() { + let tmp = TempDir::new().unwrap(); + let config = test_config(&tmp); + + let created = flows_create( + &config, + "hooked".to_string(), + webhook_trigger_graph(), + false, + ) + .await + .unwrap(); + + // Re-enable (create already enables) to exercise the enable path's warning. + let enabled = flows_set_enabled(&config, &created.value.id, true) + .await + .unwrap(); + assert!(enabled.value.enabled); + assert!( + enabled + .logs + .iter() + .any(|l| l.starts_with("warning:") && l.contains("webhook")), + "enabling a webhook-trigger flow must surface a loud warning log, got: {:?}", + enabled.logs + ); +} + +#[tokio::test] +async fn flows_set_enabled_schedule_flow_has_no_warning() { + let tmp = TempDir::new().unwrap(); + let config = test_config(&tmp); + + let created = flows_create( + &config, + "scheduled".to_string(), + schedule_trigger_graph("0 9 * * *"), + false, + ) + .await + .unwrap(); + + let enabled = flows_set_enabled(&config, &created.value.id, true) + .await + .unwrap(); + assert!( + !enabled.logs.iter().any(|l| l.starts_with("warning:")), + "a schedule-trigger flow must not surface an unfired-trigger warning: {:?}", + enabled.logs + ); +} diff --git a/src/openhuman/flows/run_registry.rs b/src/openhuman/flows/run_registry.rs new file mode 100644 index 000000000..1fe4c5e16 --- /dev/null +++ b/src/openhuman/flows/run_registry.rs @@ -0,0 +1,108 @@ +//! Process-local registry of in-flight flow runs, keyed by `run_id` +//! (== the run's checkpointer `thread_id`), so `flows_cancel_run` (issue G4) +//! can signal a synchronously-executing run to abort. +//! +//! A `flows_run` / `flows_resume` executes inline inside its RPC await (or a +//! fire-and-forget `tokio::spawn` from `flows::bus`), so there is no +//! `JoinHandle` a caller can reach. Instead each active run [`register`]s a +//! [`tokio_util::sync::CancellationToken`] here for the duration of the run and +//! `tokio::select!`s its future against the token's `cancelled()`. A separate +//! `flows_cancel_run` RPC looks the token up by `run_id` and [`cancel`]s it, +//! tripping the run's select arm. +//! +//! The registration is RAII: [`register`] returns a [`RunGuard`] that removes +//! the entry on `Drop` (including on panic / early return), so a finished run +//! can never leave a stale token wedged in the map. + +use std::collections::HashMap; +use std::sync::{LazyLock, Mutex}; + +use tokio_util::sync::CancellationToken; + +/// The live in-flight runs: `run_id` → its cancellation token. +static IN_FLIGHT: LazyLock>> = + LazyLock::new(|| Mutex::new(HashMap::new())); + +/// Registers `run_id` as in-flight and returns both a clone of its +/// cancellation token (to `select!` on) and a [`RunGuard`] that deregisters it +/// on drop. Hold the guard for the whole run. +/// +/// A duplicate `run_id` (should not happen — thread ids are UUID-suffixed) +/// replaces the prior token; the returned guard still removes exactly this +/// `run_id` on drop. +pub(crate) fn register(run_id: &str) -> (CancellationToken, RunGuard) { + let token = CancellationToken::new(); + IN_FLIGHT + .lock() + .unwrap_or_else(|e| e.into_inner()) + .insert(run_id.to_string(), token.clone()); + tracing::debug!(target: "flows", run_id, "[flows] run_registry: registered in-flight run"); + ( + token, + RunGuard { + run_id: run_id.to_string(), + }, + ) +} + +/// Signals the in-flight run keyed by `run_id` to cancel, if one is registered. +/// Returns `true` when a live run was signalled, `false` when no run with that +/// id is currently in flight (e.g. it already settled, or is a parked +/// `pending_approval` row with no executing task). +pub(crate) fn cancel(run_id: &str) -> bool { + let guard = IN_FLIGHT.lock().unwrap_or_else(|e| e.into_inner()); + match guard.get(run_id) { + Some(token) => { + token.cancel(); + tracing::info!(target: "flows", run_id, "[flows] run_registry: signalled in-flight run to cancel"); + true + } + None => { + tracing::debug!(target: "flows", run_id, "[flows] run_registry: no in-flight run to cancel"); + false + } + } +} + +/// RAII guard that removes a run's entry from the in-flight registry on drop. +pub(crate) struct RunGuard { + run_id: String, +} + +impl Drop for RunGuard { + fn drop(&mut self) { + IN_FLIGHT + .lock() + .unwrap_or_else(|e| e.into_inner()) + .remove(&self.run_id); + tracing::debug!(target: "flows", run_id = %self.run_id, "[flows] run_registry: deregistered run"); + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn cancel_signals_a_registered_run_then_deregisters_on_drop() { + let run_id = "flow:reg-test:run-1"; + let (token, guard) = register(run_id); + assert!(!token.is_cancelled()); + + // A live run is signalled and its token trips. + assert!(cancel(run_id), "a registered run must be signalled"); + assert!(token.is_cancelled(), "the run's token must be cancelled"); + + // Dropping the guard removes it; a second cancel finds nothing. + drop(guard); + assert!( + !cancel(run_id), + "after the guard drops the run must no longer be in flight" + ); + } + + #[test] + fn cancel_of_unknown_run_is_false() { + assert!(!cancel("flow:never-registered:run-x")); + } +} diff --git a/src/openhuman/flows/schemas.rs b/src/openhuman/flows/schemas.rs index 21e0788fa..62b3b945c 100644 --- a/src/openhuman/flows/schemas.rs +++ b/src/openhuman/flows/schemas.rs @@ -68,6 +68,7 @@ fn run_output_fields() -> Vec { pub fn all_controller_schemas() -> Vec { vec![ schemas("create"), + schemas("validate"), schemas("get"), schemas("list"), schemas("update"), @@ -75,6 +76,7 @@ pub fn all_controller_schemas() -> Vec { schemas("set_enabled"), schemas("run"), schemas("resume"), + schemas("cancel_run"), schemas("list_runs"), schemas("get_run"), ] @@ -86,6 +88,10 @@ pub fn all_registered_controllers() -> Vec { schema: schemas("create"), handler: handle_create, }, + RegisteredController { + schema: schemas("validate"), + handler: handle_validate, + }, RegisteredController { schema: schemas("get"), handler: handle_get, @@ -114,6 +120,10 @@ pub fn all_registered_controllers() -> Vec { schema: schemas("resume"), handler: handle_resume, }, + RegisteredController { + schema: schemas("cancel_run"), + handler: handle_cancel_run, + }, RegisteredController { schema: schemas("list_runs"), handler: handle_list_runs, @@ -149,6 +159,40 @@ pub fn schemas(function: &str) -> ControllerSchema { ], outputs: vec![flow_output()], }, + "validate" => ControllerSchema { + namespace: "flows", + function: "validate", + description: "Validate a tinyflows graph without saving it: reports structural \ + validity plus non-fatal warnings (e.g. a trigger kind that does not \ + fire automatically yet).", + inputs: vec![FieldSchema { + name: "graph", + ty: TypeSchema::Json, + comment: "A tinyflows WorkflowGraph (nodes + edges) to validate and migrate.", + required: true, + }], + outputs: vec![ + FieldSchema { + name: "valid", + ty: TypeSchema::Bool, + comment: "True when the graph is structurally valid.", + required: true, + }, + FieldSchema { + name: "errors", + ty: TypeSchema::Array(Box::new(TypeSchema::String)), + comment: "Structural validation errors; empty when `valid`.", + required: true, + }, + FieldSchema { + name: "warnings", + ty: TypeSchema::Array(Box::new(TypeSchema::String)), + comment: "Non-fatal warnings (e.g. an unfired trigger kind); the graph is \ + still saveable/enable-able.", + required: true, + }, + ], + }, "get" => ControllerSchema { namespace: "flows", function: "get", @@ -277,6 +321,15 @@ pub fn schemas(function: &str) -> ControllerSchema { comment: "Node ids being approved; defaults to an empty list.", required: false, }, + FieldSchema { + name: "rejections", + ty: TypeSchema::Option(Box::new(TypeSchema::Array(Box::new( + TypeSchema::String, + )))), + comment: "Node ids being denied; each routes to its `error` port (or fails \ + the run if it has none). Defaults to an empty list.", + required: false, + }, ], outputs: vec![FieldSchema { name: "result", @@ -287,6 +340,49 @@ pub fn schemas(function: &str) -> ControllerSchema { required: true, }], }, + "cancel_run" => ControllerSchema { + namespace: "flows", + function: "cancel_run", + description: "Cancel a flow run: settle it to a terminal `cancelled` status, abort \ + the in-flight run task if one is executing, and drop its durable \ + checkpoint so it can't be resumed.", + inputs: vec![FieldSchema { + name: "run_id", + ty: TypeSchema::String, + comment: "Identifier of the run to cancel (== its checkpoint thread id).", + required: true, + }], + outputs: vec![FieldSchema { + name: "result", + ty: TypeSchema::Object { + fields: vec![ + FieldSchema { + name: "run_id", + ty: TypeSchema::String, + comment: "Identifier of the run that was cancelled.", + required: true, + }, + FieldSchema { + name: "cancelled", + ty: TypeSchema::Bool, + comment: + "True once the run is cancelled or its cancellation requested.", + required: true, + }, + FieldSchema { + name: "was_in_flight", + ty: TypeSchema::Bool, + comment: + "True when a live run task was signalled to abort; false when \ + a parked/stale run row was settled directly.", + required: true, + }, + ], + }, + comment: "Cancellation result payload.", + required: true, + }], + }, "list_runs" => ControllerSchema { namespace: "flows", function: "list_runs", @@ -357,6 +453,14 @@ fn handle_create(params: Map) -> ControllerFuture { }) } +fn handle_validate(params: Map) -> ControllerFuture { + Box::pin(async move { + // No config load: validation is pure (no persistence, no workspace). + let graph = read_required::(¶ms, "graph")?; + to_json(ops::flows_validate(graph)) + }) +} + fn handle_get(params: Map) -> ControllerFuture { Box::pin(async move { let config = config_rpc::load_config_with_timeout().await?; @@ -438,7 +542,25 @@ fn handle_resume(params: Map) -> ControllerFuture { .transpose() .map_err(|e| format!("invalid 'approvals': {e}"))? .unwrap_or_default(); - to_json(ops::flows_resume(&config, id.trim(), thread_id.trim(), approvals).await?) + let rejections: Vec = params + .get("rejections") + .filter(|v| !v.is_null()) + .cloned() + .map(serde_json::from_value) + .transpose() + .map_err(|e| format!("invalid 'rejections': {e}"))? + .unwrap_or_default(); + to_json( + ops::flows_resume(&config, id.trim(), thread_id.trim(), approvals, rejections).await?, + ) + }) +} + +fn handle_cancel_run(params: Map) -> ControllerFuture { + Box::pin(async move { + let config = config_rpc::load_config_with_timeout().await?; + let run_id = read_required::(¶ms, "run_id")?; + to_json(ops::flows_cancel_run(&config, run_id.trim()).await?) }) } @@ -489,6 +611,7 @@ mod tests { names, vec![ "create", + "validate", "get", "list", "update", @@ -496,6 +619,7 @@ mod tests { "set_enabled", "run", "resume", + "cancel_run", "list_runs", "get_run", ] @@ -505,12 +629,13 @@ mod tests { #[test] fn all_registered_controllers_has_handler_per_schema() { let controllers = all_registered_controllers(); - assert_eq!(controllers.len(), 10); + assert_eq!(controllers.len(), 12); let names: Vec<_> = controllers.iter().map(|c| c.schema.function).collect(); assert_eq!( names, vec![ "create", + "validate", "get", "list", "update", @@ -518,6 +643,7 @@ mod tests { "set_enabled", "run", "resume", + "cancel_run", "list_runs", "get_run", ] diff --git a/src/openhuman/flows/store.rs b/src/openhuman/flows/store.rs index 3e1a483d2..26f4520e5 100644 --- a/src/openhuman/flows/store.rs +++ b/src/openhuman/flows/store.rs @@ -450,6 +450,98 @@ pub fn finish_flow_run( }) } +/// Incrementally upserts a single [`FlowRunStep`] onto a live `flow_runs` +/// row's `steps_json`, keyed by `node_id` — used by the run observer +/// (`flows::observability::FlowRunObserver`) to persist each node's step **as +/// it finishes** (issue G2, live run observation) rather than only rebuilding +/// the whole step list at settle. +/// +/// Read-modify-write under a single connection (the WAL + `busy_timeout=5000` +/// this store opens with tolerates the concurrent settle write). A re-run of +/// the same `node_id` (a retry, or a resumed run re-touching a node) replaces +/// its prior entry rather than duplicating it, so the persisted list stays one +/// entry per node. No-op if the run's start row hasn't been inserted yet +/// (nothing to update) — mirrors the best-effort contract of the run-row +/// writers in `flows::ops`. +pub fn upsert_flow_run_step(config: &Config, run_id: &str, step: &FlowRunStep) -> Result<()> { + use rusqlite::OptionalExtension; + with_connection(config, |conn| { + let existing: Option = conn + .query_row( + "SELECT steps_json FROM flow_runs WHERE id = ?1", + params![run_id], + |row| row.get(0), + ) + .optional() + .context("Failed to read flow run steps for incremental upsert")?; + let Some(raw) = existing else { + tracing::debug!(target: "flows", run_id, node = %step.node_id, "[flows] upsert_flow_run_step: no run row yet — skipping incremental step persist"); + return Ok(()); + }; + let mut steps: Vec = + serde_json::from_str(&raw).context("Failed to deserialize existing flow run steps")?; + match steps.iter_mut().find(|s| s.node_id == step.node_id) { + Some(slot) => *slot = step.clone(), + None => steps.push(step.clone()), + } + let steps_json = + serde_json::to_string(&steps).context("Failed to serialize flow run steps")?; + conn.execute( + "UPDATE flow_runs SET steps_json = ?1 WHERE id = ?2", + params![steps_json, run_id], + ) + .context("Failed to persist incremental flow run step")?; + tracing::debug!(target: "flows", run_id, node = %step.node_id, step_count = steps.len(), "[flows] persisted incremental flow run step"); + Ok(()) + }) +} + +/// Expires every parked `pending_approval` run whose "parked since" timestamp +/// (`COALESCE(finished_at, started_at)` — a run's `finished_at` is stamped when +/// it pauses at a gate) is strictly older than `cutoff` (an RFC3339 instant), +/// transitioning it to a terminal `"cancelled"` status stamped `now` with +/// `error_msg`. Returns the `(run_id, flow_id)` of each swept run so the caller +/// can update the flow summary + drop the durable checkpoint (issue G4 — +/// parked-run TTL). +/// +/// RFC3339 timestamps produced by `chrono::Utc::…to_rfc3339()` all carry the +/// same `+00:00` offset, so a lexicographic `<` is a valid chronological +/// comparison here. Best-effort by contract at the call site: the update runs +/// under the same WAL + `busy_timeout` connection as every other write. +pub fn expire_parked_runs( + config: &Config, + cutoff: &str, + now: &str, + error_msg: &str, +) -> Result> { + with_connection(config, |conn| { + let mut stmt = conn.prepare( + "SELECT id, flow_id FROM flow_runs + WHERE status = 'pending_approval' + AND COALESCE(finished_at, started_at) < ?1", + )?; + let stale: Vec<(String, String)> = stmt + .query_map(params![cutoff], |row| Ok((row.get(0)?, row.get(1)?)))? + .collect::>()?; + drop(stmt); + + for (run_id, _flow_id) in &stale { + // Re-check the status in the WHERE so a run resumed/cancelled + // between the SELECT and here is not clobbered. + conn.execute( + "UPDATE flow_runs SET status = 'cancelled', finished_at = ?1, error = ?2 \ + WHERE id = ?3 AND status = 'pending_approval'", + params![now, error_msg, run_id], + ) + .context("Failed to expire parked flow run")?; + } + if !stale.is_empty() { + tracing::info!(target: "flows", swept = stale.len(), "[flows] expired parked pending_approval runs past TTL"); + } + Ok(stale) + }) +} + /// Loads one flow run by id (== thread_id). pub fn get_flow_run(config: &Config, id: &str) -> Result> { with_connection(config, |conn| { diff --git a/src/openhuman/flows/store_tests.rs b/src/openhuman/flows/store_tests.rs index 7bbab20e8..d9e6a4144 100644 --- a/src/openhuman/flows/store_tests.rs +++ b/src/openhuman/flows/store_tests.rs @@ -268,6 +268,7 @@ fn flow_run_insert_finish_get_round_trip() { node_id: "t".to_string(), output: serde_json::json!([{"json": {"x": 1}}]), port: None, + ..Default::default() }]; finish_flow_run( &config, diff --git a/src/openhuman/flows/types.rs b/src/openhuman/flows/types.rs index c0047151f..1702f4eb8 100644 --- a/src/openhuman/flows/types.rs +++ b/src/openhuman/flows/types.rs @@ -35,6 +35,29 @@ impl FlowRunTrigger { } } +/// The result of validating a candidate `tinyflows` graph without persisting +/// it — returned by `openhuman.flows_validate` (PHASE 3c) and used to surface +/// structural errors and non-fatal warnings (e.g. "this trigger kind never +/// fires automatically yet") to an authoring surface *before* a flow is saved. +/// +/// A graph is `valid` when it passes `tinyflows::validate::validate` after +/// migration; `errors` carries the single structural error when it does not. +/// `warnings` is orthogonal to validity — a `valid` graph can still carry +/// warnings (it saves and enables fine, it just won't behave as an author +/// might expect), and an invalid graph reports no warnings (there's nothing to +/// warn about a graph that won't compile). +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Default)] +pub struct FlowValidation { + /// True when the graph is structurally valid (migrates + validates). + pub valid: bool, + /// Structural validation errors (empty when `valid`). Today at most one — + /// `tinyflows::validate::validate` returns the first error it hits. + pub errors: Vec, + /// Non-fatal warnings: the graph is accepted, but something about it is + /// worth flagging (e.g. an unfired trigger kind). Never blocks save/enable. + pub warnings: Vec, +} + /// A saved automation workflow: a `tinyflows` graph plus OpenHuman-side /// bookkeeping (enablement, run history summary). #[derive(Debug, Clone, Serialize, Deserialize)] @@ -65,28 +88,37 @@ pub struct Flow { pub require_approval: bool, } -/// One reconstructed step of a persisted [`FlowRun`] (issue B2, run-history -/// inspector). tinyflows 0.2's durable path installs a `NoopObserver` (see -/// `src/openhuman/tinyflows/observability.rs`), so there is no live per-step -/// stream to persist — instead, `flows::ops` reconstructs a lean step list -/// straight from `RunOutcome.output["nodes"]` after the run settles: each -/// entry is a node id plus its emitted items and the output port it took, if -/// any. There is no per-step timing or input/attempt data in 0.2. +/// One step of a persisted [`FlowRun`] (run-history inspector). /// -/// // TODO(0.3): a richer `RunObserver` that streams per-step -/// // node_id/status/output/duration_ms live (see -/// // `tinyflows::observability::ExecutionStep`) would let this type carry -/// // real timing/attempt data instead of being reconstructed post-hoc. +/// As of issue G2 (live run observation) these are persisted **incrementally** +/// as each non-trigger node finishes, by +/// `flows::observability::FlowRunObserver::on_step_finish`, which maps a live +/// `tinyflows::observability::ExecutionStep` (carrying real `status` + +/// `duration_ms`) onto this type. The prior post-hoc reconstruction from +/// `RunOutcome.output["nodes"]` (see `flows::ops::reconstruct_steps`) now only +/// fills in steps the observer missed (e.g. a trigger node, which does not +/// emit an `on_step_finish`) — those carry no `status`/`duration_ms` and keep +/// the `port` the reconstruction recovers. #[derive(Debug, Clone, Serialize, Deserialize, Default, PartialEq)] pub struct FlowRunStep { /// The node's id within the flow's graph. pub node_id: String, - /// The node's emitted items for this run (`output["nodes"][id]["items"]`). + /// The node's emitted items for this run (`output["nodes"][id]["items"]`, + /// or the live `ExecutionStep.output` when observed incrementally). pub output: serde_json::Value, /// The output port the node routed on, if it picked one (branching / - /// switch nodes) — `output["nodes"][id]["port"]`. + /// switch nodes) — `output["nodes"][id]["port"]`. Only recovered by the + /// post-hoc reconstruction; the live observer does not carry a port. #[serde(default, skip_serializing_if = "Option::is_none")] pub port: Option, + /// Live step outcome, when this step was observed incrementally: + /// `"success"` | `"error"`. `None` for a step recovered post-hoc. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub status: Option, + /// Wall-clock duration of the node's executor in milliseconds, when + /// observed incrementally. `None` for a step recovered post-hoc. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub duration_ms: Option, } /// A persisted record of one `flows_run` / `flows_resume` invocation, for the @@ -100,7 +132,12 @@ pub struct FlowRun { pub flow_id: String, /// The tinyflows checkpointer thread id (needed to `flows_resume`). pub thread_id: String, - /// `"running"` | `"completed"` | `"pending_approval"` | `"failed"`. + /// Run status. Not an enum (kept a free-form `String` for forward-compat + /// with statuses added by newer builds), but the vocabulary is fixed: + /// `"running"` | `"completed"` | `"pending_approval"` | `"failed"` | + /// `"cancelled"` (issue G4 — a run cancelled via `flows_cancel_run`, or a + /// parked `pending_approval` run swept by the TTL expiry). All of + /// `completed` / `failed` / `cancelled` are terminal. pub status: String, /// RFC3339 timestamp when the run started. pub started_at: String, @@ -190,6 +227,7 @@ mod tests { node_id: "t".to_string(), output: serde_json::json!([{"json": {"hello": "world"}}]), port: None, + ..Default::default() }], pending_approvals: Vec::new(), error: None, @@ -208,6 +246,7 @@ mod tests { node_id: "n".to_string(), output: serde_json::Value::Null, port: None, + ..Default::default() }; let v = serde_json::to_value(&step).unwrap(); assert!(v.get("port").is_none()); diff --git a/src/openhuman/tinyagents/middleware.rs b/src/openhuman/tinyagents/middleware.rs index 8025efa8e..de50ff7b9 100644 --- a/src/openhuman/tinyagents/middleware.rs +++ b/src/openhuman/tinyagents/middleware.rs @@ -3487,6 +3487,7 @@ mod tests { "big".to_string(), TaToolPolicy::classified().with_runtime(tinyagents::harness::tool::ToolRuntime { timeout_ms: None, + timeout: tinyagents::harness::tool::ToolTimeout::Inherit, max_retries: None, idempotent: false, cancelable: true, @@ -3516,6 +3517,7 @@ mod tests { "capped".to_string(), TaToolPolicy::classified().with_runtime(tinyagents::harness::tool::ToolRuntime { timeout_ms: None, + timeout: tinyagents::harness::tool::ToolTimeout::Inherit, max_retries: None, idempotent: false, cancelable: true, diff --git a/src/openhuman/tinyagents/observability.rs b/src/openhuman/tinyagents/observability.rs index 955b93cf1..861218d32 100644 --- a/src/openhuman/tinyagents/observability.rs +++ b/src/openhuman/tinyagents/observability.rs @@ -916,6 +916,12 @@ impl EventListener for OpenhumanEventBridge { tool_name, input, output, + // tinyagents 1.7 added started_at_ms/duration_ms/output_bytes/error + // to this event. The bridge keeps sourcing success/duration/size + // from its own capture side channel (failure_map/tool_started_at) + // below, so the new crate-provided fields are intentionally ignored + // here to preserve existing behavior. TODO: adopt them directly. + .. } => { let iteration = self.iteration(); // The crate event carries no success/error, so read what the @@ -1105,8 +1111,12 @@ mod tests { sink.emit(AgentEvent::ToolCompleted { call_id: "c1".into(), tool_name: "echo".to_string(), + started_at_ms: None, input: None, output: None, + duration_ms: None, + output_bytes: None, + error: None, }); sink.emit(AgentEvent::UsageRecorded { usage: Usage::new(100, 40), @@ -1157,6 +1167,7 @@ mod tests { }); sink.emit(AgentEvent::ModelCompleted { call_id: "m1".into(), + started_at_ms: None, usage: Some(Usage::new(1_000, 50)), input: Some(serde_json::json!([ {"role": "system", "content": "You are OpenHuman."} @@ -1223,6 +1234,7 @@ mod tests { sink.subscribe(bridge.clone()); sink.emit(AgentEvent::ModelCompleted { call_id: "m1".into(), + started_at_ms: None, usage: Some(Usage::new(10, 5)), input: None, output: None, @@ -1257,8 +1269,12 @@ mod tests { sink.emit(AgentEvent::ToolCompleted { call_id: "t1".into(), tool_name: "echo".to_string(), + started_at_ms: None, input: Some(serde_json::json!({"text": "ping"})), output: Some(serde_json::Value::String("pong".to_string())), + duration_ms: None, + output_bytes: None, + error: None, }); let mut seen = None; diff --git a/src/openhuman/tinyagents/tools.rs b/src/openhuman/tinyagents/tools.rs index dcddc1cb1..0538b315e 100644 --- a/src/openhuman/tinyagents/tools.rs +++ b/src/openhuman/tinyagents/tools.rs @@ -153,6 +153,12 @@ pub(crate) fn tool_policy_from_openhuman_tool( }) .with_runtime(ToolRuntime { timeout_ms, + // tinyagents 1.7 added a structured `timeout: ToolTimeout` field + // alongside the numeric `timeout_ms`. Inherit the run/global policy + // to preserve prior behavior (the numeric `timeout_ms` above stays + // the only per-tool deadline signal). Fully-qualified because the + // local `ToolTimeout` import here is openhuman's, not the harness's. + timeout: tinyagents::harness::tool::ToolTimeout::Inherit, max_retries: None, idempotent: tool.is_concurrency_safe(&serde_json::Value::Null), cancelable: true, diff --git a/src/openhuman/tinyflows/observability.rs b/src/openhuman/tinyflows/observability.rs index 8d97229b7..784e8d5a6 100644 --- a/src/openhuman/tinyflows/observability.rs +++ b/src/openhuman/tinyflows/observability.rs @@ -1,14 +1,32 @@ -//! Log-only [`tinyflows::observability::RunObserver`] for the `flows::` +//! Host [`tinyflows::observability::RunObserver`] impls for the `flows::` //! domain. //! //! `tinyflows` emits structured run/step records; the host decides what to do -//! with them. B1 only logs — persisting steps/runs for a run-history view is -//! B2+ (see `my_docs/ohxtf/b1-engine-seam-domain/07-execution-and-hitl.md`). -//! Note the durable path (`engine::run_with_checkpointer`, what `flows_run` -//! uses) installs a `NoopObserver` internally in 0.2, so this observer is -//! wired for `run_with_observer` call sites, not the durable run path itself. +//! with them. Two observers live here: +//! +//! - [`TracingRunObserver`] — log-only, for `run_with_observer` call sites and +//! as a simple example. +//! - [`FlowRunObserver`] — the real one used by the durable run path (issue +//! G2, live run observation). As each non-trigger node finishes it (1) +//! persists a [`FlowRunStep`] incrementally into the run's `flow_runs` row +//! via `flows::store::upsert_flow_run_step` and (2) publishes a +//! [`DomainEvent::FlowRunProgress`] so the frontend socket bridge can stream +//! the run advancing node-by-node. Both effects are best-effort: a step that +//! fails to persist is logged, never fatal to the run. +//! +//! The durable + journaled run path used by `flows_run`/`flows_resume` now +//! accepts an observer (tinyflows 0.3.1's +//! `run_with_checkpointer_journaled_observed` / +//! `resume_with_checkpointer_journaled_observed`), so [`FlowRunObserver`] sees +//! steps live instead of the run being reconstructed only after it settles. -use tinyflows::observability::{ExecutionStep, Run, RunObserver}; +use std::sync::Arc; + +use tinyflows::observability::{ExecutionStep, Run, RunObserver, StepStatus}; + +use crate::core::event_bus::{publish_global, DomainEvent}; +use crate::openhuman::config::Config; +use crate::openhuman::flows::{upsert_flow_run_step, FlowRunStep}; /// Logs run/step lifecycle events with grep-friendly `[flows]` prefixes. pub struct TracingRunObserver { @@ -43,11 +61,124 @@ impl RunObserver for TracingRunObserver { } } +/// Maps a live [`StepStatus`] to the stable `"success"`/`"error"` string +/// persisted on [`FlowRunStep::status`] and published on +/// [`DomainEvent::FlowRunProgress`]. `StepStatus` does not derive `PartialEq`, +/// so this matches rather than compares. +fn step_status_str(status: &StepStatus) -> &'static str { + match status { + StepStatus::Success => "success", + StepStatus::Error => "error", + } +} + +/// The durable-run observer (issue G2): persists each finished step +/// incrementally and streams a progress event to the frontend. +/// +/// Held by the engine as `Arc` and cloned into node handlers +/// that run across threads, so it must be cheap and non-blocking-ish — it does +/// one small SQLite read-modify-write per step (WAL + `busy_timeout` absorb the +/// concurrent settle write) and one fire-and-forget broadcast publish. Never +/// logs step output/PII: only ids, status, and timing are logged/emitted. +pub struct FlowRunObserver { + config: Arc, + flow_id: String, + /// The run's checkpointer thread id, which doubles as the `flow_runs` row + /// id and the `run_id` on the emitted progress event. + run_id: String, +} + +impl FlowRunObserver { + /// Builds an observer bound to one `flows_run` / `flows_resume` invocation. + pub fn new(config: Arc, flow_id: impl Into, run_id: impl Into) -> Self { + Self { + config, + flow_id: flow_id.into(), + run_id: run_id.into(), + } + } +} + +impl RunObserver for FlowRunObserver { + fn on_run_start(&self, run_id: &str) { + tracing::debug!( + target: "flows", + flow_id = %self.flow_id, + run_id = %self.run_id, + engine_run_id = %run_id, + "[flows] observer: run start" + ); + } + + fn on_step_finish(&self, step: &ExecutionStep) { + let status = step_status_str(&step.status); + tracing::debug!( + target: "flows", + flow_id = %self.flow_id, + run_id = %self.run_id, + node = %step.node_id, + status, + duration_ms = step.duration_ms, + "[flows] observer: step finished — persisting incrementally" + ); + + // Persist the live step. `duration_ms` is `u128` on the engine side; + // clamp into `u64` (a node executor taking > 584 million years is not a + // real concern, but never panic on cast). + let flow_step = FlowRunStep { + node_id: step.node_id.clone(), + output: step.output.clone(), + port: None, + status: Some(status.to_string()), + duration_ms: Some(u64::try_from(step.duration_ms).unwrap_or(u64::MAX)), + }; + if let Err(e) = upsert_flow_run_step(&self.config, &self.run_id, &flow_step) { + tracing::warn!( + target: "flows", + flow_id = %self.flow_id, + run_id = %self.run_id, + node = %step.node_id, + error = %e, + "[flows] observer: failed to persist incremental step (run continues; post-hoc reconstruction will fill it in at settle)" + ); + } + + // Best-effort progress feed to the frontend socket bridge. The durable + // `flow_runs` row remains the source of truth; this event may be + // dropped under broadcast lag, so the UI keeps its 2s poller fallback. + tracing::debug!( + target: "flows", + flow_id = %self.flow_id, + run_id = %self.run_id, + node = %step.node_id, + status, + "[flows] observer: publishing FlowRunProgress" + ); + publish_global(DomainEvent::FlowRunProgress { + run_id: self.run_id.clone(), + node_id: step.node_id.clone(), + status: status.to_string(), + }); + } + + fn on_run_finish(&self, run: &Run) { + tracing::debug!( + target: "flows", + flow_id = %self.flow_id, + run_id = %self.run_id, + engine_run_id = %run.id, + steps = run.steps.len(), + status = ?run.status, + "[flows] observer: run finished" + ); + } +} + #[cfg(test)] mod tests { use super::*; use serde_json::Value; - use tinyflows::observability::{RunStatus, StepStatus}; + use tinyflows::observability::RunStatus; #[test] fn callbacks_do_not_panic() { @@ -67,4 +198,16 @@ mod tests { steps: Vec::new(), }); } + + #[test] + fn step_status_maps_to_stable_strings() { + assert_eq!(step_status_str(&StepStatus::Success), "success"); + assert_eq!(step_status_str(&StepStatus::Error), "error"); + } + + // The end-to-end proof that `FlowRunObserver::on_step_finish` persists each + // step into the `flow_runs` row lives in `flows::ops_tests` + // (`observer_persists_each_step_incrementally` and the run-driven + // `flows_run_persists_live_steps_with_status_and_timing`), where the flows + // `store` internals are in scope for seeding/asserting rows. } diff --git a/tests/json_rpc_e2e.rs b/tests/json_rpc_e2e.rs index 2c254cdc0..548c13ace 100644 --- a/tests/json_rpc_e2e.rs +++ b/tests/json_rpc_e2e.rs @@ -12256,6 +12256,480 @@ async fn json_rpc_workflows_lifecycle_round_trip() { rpc_join.abort(); } +// ── openhuman.flows_* (tinyflows) JSON-RPC E2E — PHASE 1e (G5) ───────────── +// +// These exercise the tinyflows `flows_*` controller surface end-to-end over +// HTTP JSON-RPC (distinct from the legacy markdown `workflows_*` namespace +// covered above). Flows persist to a SQLite store + durable checkpoint under +// the temp `$HOME/.openhuman` workspace, so each test isolates `HOME`. +// +// All flows ops wrap their value with a log line, so `into_cli_compatible_json` +// emits `{ "result": , "logs": [...] }` — peel that envelope with +// `peel_logs_envelope` before asserting on the inner payload. + +/// Shared boot for a flows E2E: isolates `HOME`, seeds a minimal config against +/// a mock upstream, and stands up the core HTTP router. Returns the rpc base +/// URL plus the join handles + tempdir the caller must keep alive/abort. +async fn boot_flows_rpc_env() -> ( + String, + tempfile::TempDir, + tokio::task::JoinHandle>, + tokio::task::JoinHandle>, + Vec, +) { + let tmp = tempdir().expect("tempdir"); + let home = tmp.path(); + let openhuman_home = home.join(".openhuman"); + + let guards = vec![ + EnvVarGuard::set_to_path("HOME", home), + EnvVarGuard::unset("OPENHUMAN_WORKSPACE"), + EnvVarGuard::unset("BACKEND_URL"), + EnvVarGuard::unset("VITE_BACKEND_URL"), + EnvVarGuard::unset("OPENHUMAN_API_URL"), + ]; + + let (api_addr, api_join) = serve_on_ephemeral(mock_upstream_router()).await; + let api_origin = format!("http://{api_addr}"); + write_min_config(openhuman_home.as_path(), &api_origin); + + let (rpc_addr, rpc_join) = serve_on_ephemeral(build_core_http_router(false)).await; + let rpc_base = format!("http://{rpc_addr}"); + + (rpc_base, tmp, api_join, rpc_join, guards) +} + +/// The smallest valid graph with a human-in-the-loop approval gate: +/// `trigger → gate(requires_approval) → downstream`. A run pauses at `gate`; +/// approving it via `flows_resume` runs `downstream`. +fn approval_gated_graph_json() -> Value { + json!({ + "name": "approval-gated", + "nodes": [ + { "id": "t", "kind": "trigger", "name": "Trigger" }, + { "id": "gate", "kind": "output_parser", "name": "Gate", "config": { "requires_approval": true } }, + { "id": "downstream", "kind": "output_parser", "name": "Downstream" } + ], + "edges": [ + { "from_node": "t", "to_node": "gate" }, + { "from_node": "gate", "to_node": "downstream" } + ] + }) +} + +/// Full flows lifecycle over JSON-RPC: +/// create → run (park on approval) → inspect steps via get_run/list_runs → +/// resume (approve) to completion → run again → cancel_run the parked run → +/// delete. Note: there is no dedicated `flows_cancel` for an in-flight run; +/// cancel operates on a run id (checkpoint thread id), so we cancel a *fresh* +/// parked run rather than the already-completed one (cancelling a terminal run +/// is an error). +#[tokio::test] +async fn json_rpc_flows_lifecycle_round_trip() { + let _env_lock = json_rpc_e2e_env_lock(); + let (rpc_base, _tmp, api_join, rpc_join, _guards) = boot_flows_rpc_env().await; + + // 1. Create an approval-gated flow. + let create = post_json_rpc( + &rpc_base, + 9301, + "openhuman.flows_create", + json!({ "name": "Gated Demo", "graph": approval_gated_graph_json() }), + ) + .await; + let flow = peel_logs_envelope(assert_no_jsonrpc_error(&create, "flows_create")); + let flow_id = flow + .get("id") + .and_then(Value::as_str) + .expect("flow id in create result") + .to_string(); + assert_eq!(flow.get("name").and_then(Value::as_str), Some("Gated Demo")); + assert_eq!(flow.get("enabled").and_then(Value::as_bool), Some(true)); + + // 2. List reflects exactly the new flow. + let list = post_json_rpc(&rpc_base, 9302, "openhuman.flows_list", json!({})).await; + let flows = peel_logs_envelope(assert_no_jsonrpc_error(&list, "flows_list")) + .as_array() + .expect("flows array"); + assert_eq!(flows.len(), 1, "exactly one flow after create"); + assert_eq!( + flows[0].get("id").and_then(Value::as_str), + Some(flow_id.as_str()) + ); + + // 3. Run — pauses at the approval gate. Capture the thread id. + let run = post_json_rpc( + &rpc_base, + 9303, + "openhuman.flows_run", + json!({ "id": flow_id, "input": { "x": 1 } }), + ) + .await; + let run_out = peel_logs_envelope(assert_no_jsonrpc_error(&run, "flows_run")); + let thread_id = run_out + .get("thread_id") + .and_then(Value::as_str) + .expect("thread_id in run result") + .to_string(); + let pending = run_out + .get("pending_approvals") + .and_then(Value::as_array) + .expect("pending_approvals array"); + assert!( + pending.iter().any(|v| v == "gate"), + "run must park on the `gate` approval, got: {pending:?}" + ); + // Downstream must NOT have run while the gate is pending. + assert!(run_out["output"]["nodes"]["downstream"].is_null()); + + // 4. Inspect the run via get_run — status pending_approval, steps recorded. + let get_run = post_json_rpc( + &rpc_base, + 9304, + "openhuman.flows_get_run", + json!({ "run_id": thread_id }), + ) + .await; + let run_row = peel_logs_envelope(assert_no_jsonrpc_error(&get_run, "flows_get_run")); + assert_eq!( + run_row.get("status").and_then(Value::as_str), + Some("pending_approval") + ); + assert_eq!( + run_row.get("thread_id").and_then(Value::as_str), + Some(thread_id.as_str()) + ); + assert!( + run_row.get("steps").and_then(Value::as_array).is_some(), + "get_run returns a steps array" + ); + + // 5. list_runs surfaces the same single run. + let list_runs = post_json_rpc( + &rpc_base, + 9305, + "openhuman.flows_list_runs", + json!({ "id": flow_id }), + ) + .await; + let runs = peel_logs_envelope(assert_no_jsonrpc_error(&list_runs, "flows_list_runs")) + .as_array() + .expect("runs array"); + assert_eq!(runs.len(), 1, "exactly one run after a single flows_run"); + assert_eq!( + runs[0].get("id").and_then(Value::as_str), + Some(thread_id.as_str()) + ); + + // 6. Resume with the gate approved — completes and runs downstream. + let resume = post_json_rpc( + &rpc_base, + 9306, + "openhuman.flows_resume", + json!({ "id": flow_id, "thread_id": thread_id, "approvals": ["gate"] }), + ) + .await; + let resume_out = peel_logs_envelope(assert_no_jsonrpc_error(&resume, "flows_resume")); + assert_eq!(resume_out["pending_approvals"], json!([])); + assert!( + !resume_out["output"]["nodes"]["downstream"]["items"].is_null(), + "downstream should run once the gate is approved via resume" + ); + + // The persisted run row must now be terminal `completed`. + let get_run2 = post_json_rpc( + &rpc_base, + 9307, + "openhuman.flows_get_run", + json!({ "run_id": thread_id }), + ) + .await; + let run_row2 = peel_logs_envelope(assert_no_jsonrpc_error(&get_run2, "flows_get_run")); + assert_eq!( + run_row2.get("status").and_then(Value::as_str), + Some("completed") + ); + + // 7. Run again to get a fresh parked run, then cancel it. + let run2 = post_json_rpc( + &rpc_base, + 9308, + "openhuman.flows_run", + json!({ "id": flow_id, "input": { "x": 2 } }), + ) + .await; + let run2_out = peel_logs_envelope(assert_no_jsonrpc_error(&run2, "flows_run")); + let thread_id_2 = run2_out + .get("thread_id") + .and_then(Value::as_str) + .expect("thread_id in second run") + .to_string(); + + let cancel = post_json_rpc( + &rpc_base, + 9309, + "openhuman.flows_cancel_run", + json!({ "run_id": thread_id_2 }), + ) + .await; + let cancel_out = peel_logs_envelope(assert_no_jsonrpc_error(&cancel, "flows_cancel_run")); + assert_eq!( + cancel_out.get("cancelled").and_then(Value::as_bool), + Some(true) + ); + assert_eq!( + cancel_out.get("run_id").and_then(Value::as_str), + Some(thread_id_2.as_str()) + ); + // A parked run has no live task, so cancel settles the row directly. + assert_eq!( + cancel_out.get("was_in_flight").and_then(Value::as_bool), + Some(false) + ); + + // The cancelled run is terminal and no longer resumable. + let get_run3 = post_json_rpc( + &rpc_base, + 9310, + "openhuman.flows_get_run", + json!({ "run_id": thread_id_2 }), + ) + .await; + let run_row3 = peel_logs_envelope(assert_no_jsonrpc_error(&get_run3, "flows_get_run")); + assert_eq!( + run_row3.get("status").and_then(Value::as_str), + Some("cancelled") + ); + + let resume_cancelled = post_json_rpc( + &rpc_base, + 9311, + "openhuman.flows_resume", + json!({ "id": flow_id, "thread_id": thread_id_2, "approvals": ["gate"] }), + ) + .await; + assert_jsonrpc_error(&resume_cancelled, "resume of a cancelled run must error"); + + // 8. Delete the flow; the list is empty again. + let delete = post_json_rpc( + &rpc_base, + 9312, + "openhuman.flows_delete", + json!({ "id": flow_id }), + ) + .await; + let del_out = peel_logs_envelope(assert_no_jsonrpc_error(&delete, "flows_delete")); + assert_eq!(del_out.get("removed").and_then(Value::as_bool), Some(true)); + + let after = post_json_rpc(&rpc_base, 9313, "openhuman.flows_list", json!({})).await; + assert!( + peel_logs_envelope(assert_no_jsonrpc_error(&after, "flows_list")) + .as_array() + .expect("flows array") + .is_empty(), + "no flows after delete" + ); + + api_join.abort(); + rpc_join.abort(); +} + +/// Approval park + DENY over JSON-RPC (issue G4): a gate with both a `main` +/// edge (→ `downstream`) and an `error` edge (→ `recover`). Resuming with the +/// gate in `rejections` routes the denied gate's error item to `recover`, and +/// `downstream` must not run. +#[tokio::test] +async fn json_rpc_flows_resume_deny_routes_to_error_port() { + let _env_lock = json_rpc_e2e_env_lock(); + let (rpc_base, _tmp, api_join, rpc_join, _guards) = boot_flows_rpc_env().await; + + let graph = json!({ + "name": "approval-gated-error-port", + "nodes": [ + { "id": "t", "kind": "trigger", "name": "Trigger" }, + { "id": "gate", "kind": "output_parser", "name": "Gate", "config": { "requires_approval": true } }, + { "id": "downstream", "kind": "output_parser", "name": "Downstream" }, + { "id": "recover", "kind": "output_parser", "name": "Recover" } + ], + "edges": [ + { "from_node": "t", "to_node": "gate" }, + { "from_node": "gate", "from_port": "main", "to_node": "downstream" }, + { "from_node": "gate", "from_port": "error", "to_node": "recover" } + ] + }); + + let create = post_json_rpc( + &rpc_base, + 9401, + "openhuman.flows_create", + json!({ "name": "Deny Demo", "graph": graph }), + ) + .await; + let flow_id = peel_logs_envelope(assert_no_jsonrpc_error(&create, "flows_create")) + .get("id") + .and_then(Value::as_str) + .expect("flow id") + .to_string(); + + let run = post_json_rpc( + &rpc_base, + 9402, + "openhuman.flows_run", + json!({ "id": flow_id, "input": { "x": 1 } }), + ) + .await; + let thread_id = peel_logs_envelope(assert_no_jsonrpc_error(&run, "flows_run")) + .get("thread_id") + .and_then(Value::as_str) + .expect("thread_id") + .to_string(); + + // Deny the gate: no approvals, `gate` in rejections. + let resume = post_json_rpc( + &rpc_base, + 9403, + "openhuman.flows_resume", + json!({ "id": flow_id, "thread_id": thread_id, "rejections": ["gate"] }), + ) + .await; + let resume_out = peel_logs_envelope(assert_no_jsonrpc_error(&resume, "flows_resume")); + assert_eq!(resume_out["pending_approvals"], json!([])); + assert_eq!( + resume_out["output"]["nodes"]["recover"]["items"][0]["json"]["error"]["node"], + json!("gate"), + "a denied gate must route its error item to the `error`-port recovery node" + ); + assert!( + resume_out["output"]["nodes"]["downstream"].is_null(), + "the main branch must not run when the gate is denied" + ); + + let get_run = post_json_rpc( + &rpc_base, + 9404, + "openhuman.flows_get_run", + json!({ "run_id": thread_id }), + ) + .await; + let run_row = peel_logs_envelope(assert_no_jsonrpc_error(&get_run, "flows_get_run")); + assert_eq!( + run_row.get("status").and_then(Value::as_str), + Some("completed") + ); + + api_join.abort(); + rpc_join.abort(); +} + +/// `openhuman.flows_validate` over JSON-RPC (PHASE 3c): a structurally valid +/// graph whose trigger kind does not fire automatically (`webhook`) validates +/// clean but returns a loud, non-fatal warning; a `schedule` trigger (which +/// does fire) warns nothing; a graph with no trigger is `valid: false` with a +/// structural error and no warnings. +#[tokio::test] +async fn json_rpc_flows_validate_reports_warnings_and_errors() { + let _env_lock = json_rpc_e2e_env_lock(); + let (rpc_base, _tmp, api_join, rpc_join, _guards) = boot_flows_rpc_env().await; + + // 1. Webhook trigger — structurally valid, but does not fire yet → warning. + let webhook_graph = json!({ + "name": "webhook-flow", + "nodes": [ + { "id": "t", "kind": "trigger", "name": "Trigger", "config": { "trigger_kind": "webhook" } } + ], + "edges": [] + }); + let validate = post_json_rpc( + &rpc_base, + 9501, + "openhuman.flows_validate", + json!({ "graph": webhook_graph }), + ) + .await; + let v = peel_logs_envelope(assert_no_jsonrpc_error(&validate, "flows_validate")); + assert_eq!(v.get("valid").and_then(Value::as_bool), Some(true)); + assert!( + v.get("errors") + .and_then(Value::as_array) + .expect("errors") + .is_empty(), + "a structurally valid webhook graph has no errors" + ); + let warnings = v + .get("warnings") + .and_then(Value::as_array) + .expect("warnings"); + assert_eq!( + warnings.len(), + 1, + "webhook trigger emits exactly one warning" + ); + assert!( + warnings[0] + .as_str() + .is_some_and(|w| w.contains("does not fire automatically")), + "webhook warning must explain the trigger will not fire on its own, got: {warnings:?}" + ); + + // 2. Schedule trigger — fires automatically → no warning. + let schedule_graph = json!({ + "name": "scheduled-flow", + "nodes": [ + { "id": "t", "kind": "trigger", "name": "Trigger", "config": { "trigger_kind": "schedule", "schedule": "0 9 * * *" } } + ], + "edges": [] + }); + let validate_sched = post_json_rpc( + &rpc_base, + 9502, + "openhuman.flows_validate", + json!({ "graph": schedule_graph }), + ) + .await; + let vs = peel_logs_envelope(assert_no_jsonrpc_error(&validate_sched, "flows_validate")); + assert_eq!(vs.get("valid").and_then(Value::as_bool), Some(true)); + assert!( + vs.get("warnings") + .and_then(Value::as_array) + .expect("warnings") + .is_empty(), + "a schedule trigger fires automatically — no unfired-kind warning" + ); + + // 3. No trigger node — structurally invalid. + let invalid_graph = json!({ + "name": "no-trigger", + "nodes": [ { "id": "a", "kind": "output_parser", "name": "A" } ], + "edges": [] + }); + let validate_bad = post_json_rpc( + &rpc_base, + 9503, + "openhuman.flows_validate", + json!({ "graph": invalid_graph }), + ) + .await; + let vb = peel_logs_envelope(assert_no_jsonrpc_error(&validate_bad, "flows_validate")); + assert_eq!(vb.get("valid").and_then(Value::as_bool), Some(false)); + assert!( + !vb.get("errors") + .and_then(Value::as_array) + .expect("errors") + .is_empty(), + "a graph without a trigger must report a structural error" + ); + assert!( + vb.get("warnings") + .and_then(Value::as_array) + .expect("warnings") + .is_empty(), + "an invalid graph reports errors, not warnings" + ); + + api_join.abort(); + rpc_join.abort(); +} + /// Task 4 / #3090: when a web-chat request is sent with /// `speak_reply: true`, `run_chat_task` should drive the agent's final text /// through `voice::reply_speech::synthesize_reply` after the turn completes. diff --git a/vendor/tinyagents b/vendor/tinyagents index df391c468..3e81e4935 160000 --- a/vendor/tinyagents +++ b/vendor/tinyagents @@ -1 +1 @@ -Subproject commit df391c46891cb69ffd39f69023f497a5d95bec0e +Subproject commit 3e81e4935efeb05901915d2e54e1d4aa71a4185a diff --git a/vendor/tinyflows b/vendor/tinyflows index 438f8fce2..96635d981 160000 --- a/vendor/tinyflows +++ b/vendor/tinyflows @@ -1 +1 @@ -Subproject commit 438f8fce2dc54cfd98e57d9a2fe93d5e07aad27c +Subproject commit 96635d9812cff48bbe4244563a94eeb7bfe730da