From a364693fc7fb3a5a4b2ba0d3fd4de953d29501d1 Mon Sep 17 00:00:00 2001 From: Mega Mind <146339422+M3gA-Mind@users.noreply.github.com> Date: Mon, 13 Jul 2026 15:25:09 +0530 Subject: [PATCH] fix(composio): bound composio_connect approval park so it fails fast to a connect prompt (#4758) --- src/openhuman/approval/gate.rs | 194 +++++++++++++++++++++++++- src/openhuman/composio/tools.rs | 84 ++++++++++- src/openhuman/composio/tools_tests.rs | 34 +++++ 3 files changed, 308 insertions(+), 4 deletions(-) diff --git a/src/openhuman/approval/gate.rs b/src/openhuman/approval/gate.rs index cf09480f2..e074b3dd5 100644 --- a/src/openhuman/approval/gate.rs +++ b/src/openhuman/approval/gate.rs @@ -330,6 +330,76 @@ impl ApprovalGate { tool_name: &str, action_summary: &str, args_redacted: serde_json::Value, + ) -> (GateOutcome, Option) { + // No caller-supplied park bound: identical behavior to before. With + // `park_bound = None` the inner never takes the caller-bound abandon + // path, so the out-flag stays `false` and is discarded here. + let mut _park_bound_elapsed = false; + self.intercept_audited_inner( + tool_name, + action_summary, + args_redacted, + None, + &mut _park_bound_elapsed, + ) + .await + } + + /// Like [`Self::intercept_audited`] but the caller may cap how long the + /// gate parks (issue #4756). + /// + /// When `park_bound` is `Some` and shorter than the gate's own effective + /// TTL and it elapses before a decision arrives, the gate abandons the park + /// in a **cancellation-safe** way — it evicts the in-memory waiter and + /// clears the thread/meeting routing mappings (so a later chat/voice reply + /// is not mis-routed to this now-abandoned request) but deliberately LEAVES + /// the `pending_approvals` row open, so a later human card-click can still + /// resolve it in the DB — and returns `None`. This is why callers must bound + /// the park through the gate rather than racing an outer + /// `tokio::time::timeout` against [`Self::intercept_audited`]: dropping the + /// parked future would skip that cleanup and orphan the waiter + routing + /// mappings (chatgpt-codex review on #4756). + /// + /// A `None` bound (or one `>=` the effective TTL) behaves exactly like + /// [`Self::intercept_audited`] and always returns `Some`. + pub async fn intercept_audited_bounded( + &self, + tool_name: &str, + action_summary: &str, + args_redacted: serde_json::Value, + park_bound: Option, + ) -> Option<(GateOutcome, Option)> { + let mut park_bound_elapsed = false; + let resolved = self + .intercept_audited_inner( + tool_name, + action_summary, + args_redacted, + park_bound, + &mut park_bound_elapsed, + ) + .await; + if park_bound_elapsed { + None + } else { + Some(resolved) + } + } + + /// Shared core of [`Self::intercept_audited`] and + /// [`Self::intercept_audited_bounded`]. When `park_bound` is `Some` and + /// shorter than the effective TTL, the park is capped at it; on that bound + /// elapsing the park is abandoned cancellation-safely (waiter evicted, + /// thread/meeting routing cleared, `pending_approvals` row left open) and + /// `*park_bound_elapsed` is set so the bounded caller can render its own + /// fast-path result instead of a `Deny`. + async fn intercept_audited_inner( + &self, + tool_name: &str, + action_summary: &str, + args_redacted: serde_json::Value, + park_bound: Option, + park_bound_elapsed: &mut bool, ) -> (GateOutcome, Option) { // Origin tells us who scheduled this turn. Entry points (web channel, // channel runtime, subconscious, cron, CLI) scope a typed @@ -734,7 +804,19 @@ impl ApprovalGate { self.effective_ttl() }; - let outcome = match tokio::time::timeout(effective_ttl, rx).await { + // Optional caller-supplied park bound (issue #4756). A caller + // (`composio_connect`) can cap how long the gate parks so a turn + // degrades to a fast prompt instead of blocking to the full TTL. + // Bounding must never *extend* the park, so we wait `min(bound, ttl)`; + // the caller-bound abandon path fires only when the bound is what + // elapses (`park_bound_active`). + let park_bound_active = matches!(park_bound, Some(b) if b < effective_ttl); + let wait = match park_bound { + Some(b) => b.min(effective_ttl), + None => effective_ttl, + }; + + let outcome = match tokio::time::timeout(wait, rx).await { Ok(Ok(decision)) => { tracing::info!( request_id = %request_id, @@ -776,6 +858,39 @@ impl ApprovalGate { None, ) } + Err(_elapsed) if park_bound_active => { + // Caller park bound elapsed (#4756) — NOT the gate's own TTL. + // Abandon the park cancellation-safely: evict the in-memory + // waiter and (via `clear_thread`/`clear_meeting` below, on every + // exit) drop the routing mappings so a later chat/voice reply is + // not mis-routed to this now-abandoned request. Deliberately do + // NOT `store::decide(Deny)` — the `pending_approvals` row stays + // open so a later human card-click still resolves it in the DB + // and a re-ask sees it already-connected. Signal the elapse so + // the bounded caller renders its own fast-path result rather than + // a `Deny`. + self.evict_waiter(&request_id); + *park_bound_elapsed = true; + tracing::info!( + request_id = %request_id, + tool = tool_name, + bound_secs = wait.as_secs(), + "[approval::gate] caller park bound elapsed — abandoning park (row left \ + pending for a later card-click; waiter + routing cleared) (#4756)" + ); + // Placeholder outcome; the bounded caller discards it once + // `*park_bound_elapsed` is set (returns `None`). + ( + GateOutcome::Deny { + reason: format!( + "{POLICY_DENIED_MARKER} Approval for '{tool_name}' exceeded the caller \ + park bound ({}s).", + wait.as_secs() + ), + }, + None, + ) + } Err(_elapsed) => { self.evict_waiter(&request_id); // Race: `decide()` may have committed an Approve in @@ -1495,6 +1610,83 @@ mod tests { assert!(gate.pending_for_thread("thread-42").is_none()); } + // ── caller park bound (issue #4756) ────────────────────────────── + // + // A caller (composio_connect) can cap the park via + // `intercept_audited_bounded`. When the bound elapses before the gate's own + // TTL the gate must abandon the park cancellation-safely: return `None`, + // clear the thread→request routing so a later reply is not mis-routed (the + // codex concern), yet LEAVE the `pending_approvals` row open so a later + // card-click still resolves it in the DB. + #[tokio::test] + async fn intercept_audited_bounded_abandons_park_and_leaves_row_pending() { + let (gate, _dir) = test_gate(); // boot-time TTL = 2s + let gate = Arc::new(gate); + + let g = gate.clone(); + let ctx = ApprovalChatContext { + thread_id: "thread-bound".into(), + client_id: "client-1".into(), + }; + let origin = AgentTurnOrigin::WebChat { + thread_id: "thread-bound".into(), + client_id: "client-1".into(), + request_id: Some("req-bound".into()), + }; + // 100ms caller bound — far below the 2s gate TTL — so the bound is what + // elapses, not the gate's own timeout. + let handle = tokio::spawn(async move { + turn_origin::with_origin( + origin, + APPROVAL_CHAT_CONTEXT.scope( + ctx, + g.intercept_audited_bounded( + "shell", + "run ls", + serde_json::json!({}), + Some(Duration::from_millis(100)), + ), + ), + ) + .await + }); + + // While parked, the thread → request mapping is queryable. + let mut tries = 0; + let request_id = loop { + if let Some(r) = gate.pending_for_thread("thread-bound") { + break r; + } + tries += 1; + assert!(tries < 50, "thread mapping never appeared"); + tokio::time::sleep(Duration::from_millis(5)).await; + }; + + // The bound elapses → `None`, so the caller renders its own fast path + // instead of the park resolving to a Deny. + let resolved = handle.await.unwrap(); + assert!( + resolved.is_none(), + "caller park bound must surface as None, not a resolved outcome" + ); + + // Routing is cleared so a later reply is not mis-routed to the abandoned + // request (the codex #4756 concern). + assert!( + gate.pending_for_thread("thread-bound").is_none(), + "thread → request mapping must be cleared on caller-bound abandon" + ); + + // The row is LEFT open — a later human card-click still resolves it. + let decided = gate + .decide(&request_id, ApprovalDecision::ApproveOnce) + .unwrap(); + assert!( + decided.is_some(), + "pending row must survive the abandon so a later card-click resolves it" + ); + } + /// Tests for `effective_ttl` env-override parsing. /// /// These run serially (they mutate the process env) via the shared diff --git a/src/openhuman/composio/tools.rs b/src/openhuman/composio/tools.rs index f7464b8f2..3218f429f 100644 --- a/src/openhuman/composio/tools.rs +++ b/src/openhuman/composio/tools.rs @@ -685,6 +685,43 @@ fn canonicalize_toolkit_slug(slug: &str) -> String { } } +/// Default bound (seconds) for how long [`ComposioConnectTool`] parks on the +/// inline-connect approval card before giving up (issue #4756). +/// +/// The gate's own TTL is up to ten minutes (`DEFAULT_APPROVAL_TTL` in +/// `approval::gate`). That is fine when a human is watching the card, but when +/// the card can't be resolved — a headless/eval run, or a chat turn whose +/// client has since disconnected — `composio_connect` would otherwise block the +/// whole turn for minutes and deliver an empty reply, while the read path +/// (`composio_list_connections`) returns a graceful "not connected" prompt in +/// seconds. Bounding the park keeps the interactive resume-in-turn UX for a +/// present user (a click + OAuth round-trip completes well inside it) while +/// guaranteeing the act path degrades to a fast connect prompt instead of +/// hanging. Generous by design; env-overridable, `0` restores the full gate TTL. +const DEFAULT_COMPOSIO_CONNECT_TIMEOUT_SECS: u64 = 120; + +/// Resolve the connect-card park bound. Reads +/// `OPENHUMAN_COMPOSIO_CONNECT_TIMEOUT_SECS`; `0` means "no composio-side bound" +/// (`None`) → fall back to the gate's own TTL. +fn composio_connect_timeout() -> Option { + parse_composio_connect_timeout( + std::env::var("OPENHUMAN_COMPOSIO_CONNECT_TIMEOUT_SECS") + .ok() + .as_deref(), + ) +} + +/// Pure core of [`composio_connect_timeout`], kept env-free so it is +/// deterministically unit-testable. An absent/unparseable value falls back to +/// [`DEFAULT_COMPOSIO_CONNECT_TIMEOUT_SECS`]; `0` yields `None` (opt out of the +/// composio-side bound). +fn parse_composio_connect_timeout(env_value: Option<&str>) -> Option { + let secs = env_value + .and_then(|v| v.trim().parse::().ok()) + .unwrap_or(DEFAULT_COMPOSIO_CONNECT_TIMEOUT_SECS); + (secs > 0).then(|| std::time::Duration::from_secs(secs)) +} + /// Fresh (uncached) liveness check for `toolkit`. /// /// Tri-state via `Result`: @@ -883,9 +920,50 @@ impl Tool for ComposioConnectTool { } }; let summary = format!("Connect {toolkit} to complete your task"); - let (outcome, _request_id) = gate - .intercept_audited("composio_connect", &summary, json!({ "toolkit": toolkit })) - .await; + // Bound the park (issue #4756). The gate parks up to its full TTL (10 + // min) waiting for the inline connect card to resolve; when nothing + // resolves it — a headless/eval run, or a chat client that has since + // disconnected — that otherwise blocks the whole turn to an empty reply. + // `intercept_audited_bounded` caps the park at `composio_connect_timeout()` + // and, when that bound elapses, abandons the park *inside the gate* in a + // cancellation-safe way (waiter evicted + thread/meeting routing cleared, + // but the `pending_approvals` row left open so a later human card-click + // still resolves it in the DB and a re-ask sees it already-connected). It + // returns `None` on that elapse, so we degrade to a fast, actionable + // connect prompt — matching the read path — rather than hanging. We bound + // through the gate (not an outer `tokio::time::timeout` that would drop + // the parked future and orphan the waiter/routing) per the codex review + // on this PR. The reply is shaped so the agent RELAYS it and does NOT + // immediately retry `composio_connect` (a retry would just park again). + let (outcome, _request_id) = match gate + .intercept_audited_bounded( + "composio_connect", + &summary, + json!({ "toolkit": toolkit }), + composio_connect_timeout(), + ) + .await + { + Some(resolved) => resolved, + None => { + tracing::info!( + toolkit = %toolkit, + "[composio] connect.execute: approval card not resolved within bound — \ + returning a fast connect prompt instead of parking the turn (#4756)" + ); + return Ok(ToolResult::success(serde_json::to_string(&json!({ + "toolkit": toolkit, + "connected": false, + "pending": true, + "reason": format!( + "A Connect card for {toolkit} was raised but wasn't completed in time \ + (no one authorized it). Tell the user to click Connect on the card, or \ + connect {toolkit} in Settings → Connections, then ask again once it's \ + done. Do not call composio_connect again until they confirm." + ), + }))?)); + } + }; match outcome { crate::openhuman::approval::GateOutcome::Allow => { // `Allow` only means the prompt was approved — re-check liveness diff --git a/src/openhuman/composio/tools_tests.rs b/src/openhuman/composio/tools_tests.rs index 888ba10f5..5a8849f03 100644 --- a/src/openhuman/composio/tools_tests.rs +++ b/src/openhuman/composio/tools_tests.rs @@ -1086,3 +1086,37 @@ async fn authorize_in_direct_mode_refuses_with_app_composio_dev_hint() { "must not leak backend-tenant routing artifacts in direct mode: {msg}" ); } + +// ── composio_connect park bound (issue #4756) ──────────────────────── +// +// composio_connect parks on the inline-connect approval card up to the gate's +// full TTL. When nothing resolves it (headless/eval run, or a disconnected +// chat client) that blocked the whole turn to an empty reply, while the read +// path returns a graceful "not connected" prompt fast. The park is now bounded +// by `composio_connect_timeout()`; these cover its pure env parser. + +#[test] +fn parse_composio_connect_timeout_defaults_when_absent_or_garbage() { + let default = std::time::Duration::from_secs(DEFAULT_COMPOSIO_CONNECT_TIMEOUT_SECS); + // Absent → default bound (never unbounded by accident). + assert_eq!(parse_composio_connect_timeout(None), Some(default)); + // Unparseable → default bound. + assert_eq!(parse_composio_connect_timeout(Some("soon")), Some(default)); + assert_eq!(parse_composio_connect_timeout(Some("")), Some(default)); +} + +#[test] +fn parse_composio_connect_timeout_honors_override_and_zero_opt_out() { + // Explicit value → that many seconds. + assert_eq!( + parse_composio_connect_timeout(Some("45")), + Some(std::time::Duration::from_secs(45)) + ); + // Whitespace tolerated. + assert_eq!( + parse_composio_connect_timeout(Some(" 90 ")), + Some(std::time::Duration::from_secs(90)) + ); + // `0` → opt out of the composio-side bound (fall back to the gate TTL). + assert_eq!(parse_composio_connect_timeout(Some("0")), None); +}