From afdc268040e8d81ef3f60d164debe16505eecd38 Mon Sep 17 00:00:00 2001 From: oxoxDev <164490987+oxoxDev@users.noreply.github.com> Date: Wed, 13 May 2026 09:15:06 +0530 Subject: [PATCH] fix(observability): drop transient upstream HTTP from Sentry (429/408/502/503/504) (#1529) Co-authored-by: Claude Opus 4.7 (1M context) Co-authored-by: Steven Enamakel --- Cargo.toml | 7 ++ src/core/observability.rs | 142 ++++++++++++++++++++- src/main.rs | 13 ++ src/openhuman/agent/harness/tool_loop.rs | 43 +++++-- src/openhuman/providers/ops.rs | 71 ++++++----- src/openhuman/providers/reliable.rs | 42 ++++++- src/openhuman/providers/reliable_tests.rs | 29 +++++ tests/observability_smoke.rs | 144 ++++++++++++++++++++++ 8 files changed, 450 insertions(+), 41 deletions(-) create mode 100644 tests/observability_smoke.rs diff --git a/Cargo.toml b/Cargo.toml index 73a4cfdf0..cfc5e25ed 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -156,6 +156,13 @@ landlock = { version = "0.4", optional = true } rppal = { version = "0.22", optional = true } [dev-dependencies] +# Enable sentry's TestTransport for runtime smoke of the observability +# before_send filter (see tests/observability_smoke.rs). `default-features +# = false` here is load-bearing — sentry's default feature set pulls in +# actix-web / actix-http / actix-server / sentry-actix and ~13 transitive +# crates we never use (and that bloat the dev Cargo.lock noticeably). +# TestTransport only needs the `test` feature. +sentry = { version = "0.47.0", default-features = false, features = ["test"] } [features] sandbox-landlock = ["dep:landlock"] diff --git a/src/core/observability.rs b/src/core/observability.rs index 417e32544..e0a5afb2c 100644 --- a/src/core/observability.rs +++ b/src/core/observability.rs @@ -1,4 +1,6 @@ -//! Centralised error reporting for the core. +//! Centralised error reporting for the core, plus a Sentry +//! `before_send` filter that drops per-attempt transient-upstream +//! provider failures. //! //! Wraps `tracing::error!` (which the global subscriber forwards to Sentry via //! `sentry-tracing`) inside a `sentry::with_scope` so each captured event @@ -18,6 +20,21 @@ use std::fmt::Display; /// anything you'd want to facet on (`error_kind`, `tool_name`, `method`). pub type Tag<'a> = (&'a str, &'a str); +/// HTTP status codes that the reliable-provider layer already handles via +/// retry + fallback, so per-attempt Sentry reports add noise without signal: +/// +/// - **408** Request Timeout +/// - **429** Too Many Requests +/// - **502** Bad Gateway +/// - **503** Service Unavailable +/// - **504** Gateway Timeout +/// +/// Single source of truth for both the call-site classifier +/// (`openhuman::providers::ops::should_report_provider_http_failure`) and the +/// `before_send` filter (`is_transient_provider_http_failure`). Update here +/// and both sites pick it up — keeps the two layers from drifting. +pub const TRANSIENT_PROVIDER_HTTP_STATUSES: &[u16] = &[408, 429, 502, 503, 504]; + #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum ExpectedErrorKind { LocalAiDisabled, @@ -122,6 +139,37 @@ fn report_error_message(message: &str, domain: &str, operation: &str, extra: &[T ); } +/// Returns true when a Sentry event is a per-attempt provider HTTP failure +/// that the reliable-provider layer already handles via retry + fallback. +/// +/// The primary suppression lives at the call site +/// (`openhuman::providers::ops::should_report_provider_http_failure`), +/// which short-circuits transient codes before `report_error` ever fires. +/// This helper is intended for use inside the `sentry::ClientOptions` +/// `before_send` hook as defense-in-depth — it catches any future call +/// site that emits a `tracing::error!` with the same shape but bypasses +/// the classifier. +/// +/// Match criteria (all required): +/// - tag `domain == "llm_provider"` — pins the filter to provider-originated +/// events so an unrelated subsystem emitting `failure=non_2xx`/`status=503` +/// for its own reasons doesn't get silently dropped +/// - tag `failure == "non_2xx"` (the marker set by `ops::api_error`) +/// - tag `status` parses to one of [`TRANSIENT_PROVIDER_HTTP_STATUSES`] +pub fn is_transient_provider_http_failure(event: &sentry::protocol::Event<'_>) -> bool { + let tags = &event.tags; + if tags.get("domain").map(String::as_str) != Some("llm_provider") { + return false; + } + if tags.get("failure").map(String::as_str) != Some("non_2xx") { + return false; + } + let Some(status_u16) = tags.get("status").and_then(|s| s.parse::().ok()) else { + return false; + }; + TRANSIENT_PROVIDER_HTTP_STATUSES.contains(&status_u16) +} + #[cfg(test)] mod tests { use super::*; @@ -178,6 +226,98 @@ mod tests { ); } + fn event_with_tags(pairs: &[(&str, &str)]) -> sentry::protocol::Event<'static> { + let mut event = sentry::protocol::Event::default(); + let mut tags: std::collections::BTreeMap = + std::collections::BTreeMap::new(); + for (k, v) in pairs { + tags.insert((*k).to_string(), (*v).to_string()); + } + event.tags = tags; + event + } + + #[test] + fn transient_filter_drops_429_408_502_503_504() { + for status in ["429", "408", "502", "503", "504"] { + let event = event_with_tags(&[ + ("domain", "llm_provider"), + ("failure", "non_2xx"), + ("status", status), + ]); + assert!( + is_transient_provider_http_failure(&event), + "status {status} must be classified as transient and filtered" + ); + } + } + + #[test] + fn transient_filter_keeps_permanent_failures() { + for status in ["400", "401", "403", "404", "500"] { + let event = event_with_tags(&[ + ("domain", "llm_provider"), + ("failure", "non_2xx"), + ("status", status), + ]); + assert!( + !is_transient_provider_http_failure(&event), + "status {status} must NOT be filtered — it's actionable" + ); + } + } + + #[test] + fn transient_filter_keeps_aggregate_all_exhausted() { + let event = event_with_tags(&[ + ("domain", "llm_provider"), + ("failure", "all_exhausted"), + ("status", "503"), + ]); + assert!( + !is_transient_provider_http_failure(&event), + "aggregate all_exhausted events must surface (they are the cascade signal)" + ); + } + + #[test] + fn transient_filter_keeps_events_with_no_status_tag() { + let event = event_with_tags(&[("domain", "llm_provider"), ("failure", "non_2xx")]); + assert!( + !is_transient_provider_http_failure(&event), + "missing status tag must not be silently dropped" + ); + } + + // Regression guard: the filter must scope to provider events only. Other + // subsystems emit `failure=non_2xx` (e.g. + // `providers/compatible.rs` uses the same marker for OAI-compatible + // error paths, but every site goes through `report_error(.., + // "llm_provider", ..)` so the domain tag is consistent), but the broader + // point is: any future caller that re-uses the same tag set for a + // different domain must NOT be silently dropped by this filter. + #[test] + fn transient_filter_keeps_events_with_no_domain_tag() { + let event = event_with_tags(&[("failure", "non_2xx"), ("status", "503")]); + assert!( + !is_transient_provider_http_failure(&event), + "missing domain tag means the event isn't provider-originated — must surface" + ); + } + + #[test] + fn transient_filter_keeps_events_from_other_domains() { + let event = event_with_tags(&[ + ("domain", "scheduler"), + ("failure", "non_2xx"), + ("status", "503"), + ]); + assert!( + !is_transient_provider_http_failure(&event), + "non-provider domain must surface even if failure/status tags collide" + ); + } + #[test] fn report_error_or_expected_does_not_panic() { report_error_or_expected( diff --git a/src/main.rs b/src/main.rs index decf62eb3..e5fb20f13 100644 --- a/src/main.rs +++ b/src/main.rs @@ -46,6 +46,19 @@ fn main() { environment: Some(std::borrow::Cow::Owned(resolve_environment())), send_default_pii: false, before_send: Some(std::sync::Arc::new(|mut event| { + // Defense-in-depth: drop transient-upstream provider failures that + // slipped past the call-site classifier. The reliable-provider + // layer already retries 429/408/502/503/504 with backoff + + // fallback, and the aggregate "all providers exhausted" event + // still fires for genuine outages. Per-attempt reports flood + // Sentry — see OPENHUMAN-TAURI-2E (~1393 events), -84 (~1050), + // -T (~871). The primary fix lives in + // `openhuman::providers::ops::should_report_provider_http_failure` + // (transient codes excluded). This filter catches any future call + // site that bypasses it. + if openhuman_core::core::observability::is_transient_provider_http_failure(&event) { + return None; + } // Strip server_name (hostname) to avoid leaking machine identity event.server_name = None; // Attach the cached account uid so Sentry can count unique users diff --git a/src/openhuman/agent/harness/tool_loop.rs b/src/openhuman/agent/harness/tool_loop.rs index 84bce73db..4da38fe6d 100644 --- a/src/openhuman/agent/harness/tool_loop.rs +++ b/src/openhuman/agent/harness/tool_loop.rs @@ -408,16 +408,39 @@ pub(crate) async fn run_tool_call_loop( ) } Err(e) => { - crate::core::observability::report_error_or_expected( - &e, - "agent", - "provider_chat", - &[ - ("provider", provider_name), - ("model", model), - ("iteration", &(iteration + 1).to_string()), - ], - ); + // Transient upstream failures (rate-limit, gateway 5xx, "no + // healthy upstream", etc.) are already classified + retried + // by reliable.rs and produce an aggregate Sentry event only + // when every provider/model is exhausted. Reporting each + // per-iteration provider_chat error here duplicates the + // signal and floods Sentry — see OPENHUMAN-TAURI-3Y/3Z + // (~46 events combined) and the underlying TAURI-2E/84/T + // (~3300 events from raw per-attempt 429/503/504 reports). + let transient = crate::openhuman::providers::reliable::is_rate_limited(&e) + || crate::openhuman::providers::reliable::is_upstream_unhealthy(&e); + if transient { + tracing::warn!( + domain = "agent", + operation = "provider_chat", + provider = provider_name, + model = model, + iteration = iteration + 1, + error = %format!("{e:#}"), + "[agent] transient provider_chat failure — retried upstream; \ + aggregated all-providers-exhausted will report if applicable" + ); + } else { + crate::core::observability::report_error_or_expected( + &e, + "agent", + "provider_chat", + &[ + ("provider", provider_name), + ("model", model), + ("iteration", &(iteration + 1).to_string()), + ], + ); + } return Err(e); } }; diff --git a/src/openhuman/providers/ops.rs b/src/openhuman/providers/ops.rs index bd466a624..ec33590c0 100644 --- a/src/openhuman/providers/ops.rs +++ b/src/openhuman/providers/ops.rs @@ -114,15 +114,18 @@ pub fn format_anyhow_chain(err: &anyhow::Error) -> String { /// Whether a non-2xx provider response is worth reporting to Sentry. /// -/// 429 Too Many Requests is a transient, caller-side throttling signal — the -/// reliable-provider layer already retries with backoff and falls back across -/// providers/models, and the aggregate "all providers exhausted" event still -/// fires if every attempt fails. Reporting each individual 429 floods Sentry -/// (see OPENHUMAN-TAURI-6Y: ~8K events/day from one user being rate-limited -/// by an upstream model). Callers should still propagate the error so retry -/// and fallback logic runs unchanged; this only gates the Sentry report. +/// Transient upstream statuses — 429 Too Many Requests, 408 Request Timeout, +/// and 502/503/504 gateway-layer failures — are caller-side throttling or +/// upstream-capacity signals. The reliable-provider layer already retries +/// with backoff and falls back across providers/models, and the aggregate +/// "all providers exhausted" event still fires if every attempt fails. +/// Reporting each individual transient failure floods Sentry (see +/// OPENHUMAN-TAURI-6Y / 2E / 84 / T: thousands of events/day per user from +/// a single upstream rate-limit / outage window). Callers should still +/// propagate the error so retry and fallback logic runs unchanged; this +/// only gates the per-attempt Sentry report. pub fn should_report_provider_http_failure(status: reqwest::StatusCode) -> bool { - status != reqwest::StatusCode::TOO_MANY_REQUESTS + !crate::core::observability::TRANSIENT_PROVIDER_HTTP_STATUSES.contains(&status.as_u16()) } /// Whether a "Budget exceeded" error from `provider` at `status` should be @@ -430,26 +433,38 @@ mod tests { } #[test] - fn skips_sentry_report_for_429_only() { - // 429 is transient rate-limit — reliable.rs retries + falls back, and - // the aggregate "all providers exhausted" event still fires for genuine - // outages. Reporting each 429 individually floods Sentry. - assert!(!should_report_provider_http_failure( - reqwest::StatusCode::TOO_MANY_REQUESTS - )); - // Everything else (auth, server, gateway, etc.) is still worth a report. - assert!(should_report_provider_http_failure( - reqwest::StatusCode::UNAUTHORIZED - )); - assert!(should_report_provider_http_failure( - reqwest::StatusCode::INTERNAL_SERVER_ERROR - )); - assert!(should_report_provider_http_failure( - reqwest::StatusCode::BAD_GATEWAY - )); - assert!(should_report_provider_http_failure( - reqwest::StatusCode::SERVICE_UNAVAILABLE - )); + fn skips_sentry_report_for_transient_upstream_statuses() { + // Transient statuses — 429 rate-limit, 408 client timeout, and 502/503/504 + // gateway-layer failures — are retried by reliable.rs. The aggregate + // "all providers exhausted" event still fires for genuine outages. + // Reporting each attempt individually floods Sentry (OPENHUMAN-TAURI-2E + // ~1393 events, 84 ~1050 events, T ~871 events). + for transient in [ + reqwest::StatusCode::TOO_MANY_REQUESTS, + reqwest::StatusCode::REQUEST_TIMEOUT, + reqwest::StatusCode::BAD_GATEWAY, + reqwest::StatusCode::SERVICE_UNAVAILABLE, + reqwest::StatusCode::GATEWAY_TIMEOUT, + ] { + assert!( + !should_report_provider_http_failure(transient), + "transient status {transient} must not trigger per-attempt Sentry report" + ); + } + // Auth + permanent server faults remain reportable — those are + // misconfiguration or genuine bugs, not transient capacity issues. + for reportable in [ + reqwest::StatusCode::UNAUTHORIZED, + reqwest::StatusCode::FORBIDDEN, + reqwest::StatusCode::BAD_REQUEST, + reqwest::StatusCode::NOT_FOUND, + reqwest::StatusCode::INTERNAL_SERVER_ERROR, + ] { + assert!( + should_report_provider_http_failure(reportable), + "status {reportable} must still report to Sentry" + ); + } } // Confirm the Budget-exceeded suppression predicate is scoped correctly. diff --git a/src/openhuman/providers/reliable.rs b/src/openhuman/providers/reliable.rs index 6d675852b..8f8bab787 100644 --- a/src/openhuman/providers/reliable.rs +++ b/src/openhuman/providers/reliable.rs @@ -75,13 +75,51 @@ fn is_context_window_exceeded(err: &anyhow::Error) -> bool { hints.iter().any(|hint| lower.contains(hint)) } -/// Detect provider-side temporary capacity/outage errors that are often surfaced -/// as HTTP 5xx with text like "no healthy upstream". +/// Detect provider-side temporary capacity/outage errors. Covers: +/// +/// - HTTP `408 Request Timeout`, `502 Bad Gateway`, `503 Service Unavailable`, +/// `504 Gateway Timeout` — both via direct `reqwest::Error` downcast and via +/// the formatted `" API error (): …"` text emitted by +/// `ops::api_error` (the path that actually reaches `report_error`). +/// - Provider-agnostic text markers like `"no healthy upstream"` / +/// `"upstream unavailable"` that don't come with a typed status. +/// +/// Pairs with [`is_rate_limited`] which handles 429 separately. Together they +/// form the transient-classifier the tool-call loop uses before deciding +/// whether to push a per-attempt event to Sentry (see OPENHUMAN-TAURI-2E / +/// -84 / -T / -G classes — per-iteration noise from upstream throttling). +/// +/// **Status list maintenance note**: the codes matched below (408/502/503/504) +/// are a subset of +/// [`crate::core::observability::TRANSIENT_PROVIDER_HTTP_STATUSES`] — that +/// const is the single source of truth for the `before_send` filter and the +/// call-site classifier in `providers/ops.rs`. We don't reference the const +/// directly here because this function takes a different code path (anyhow +/// error downcast vs typed `reqwest::StatusCode`) and because 429 is split out +/// into `is_rate_limited` (with its own retry-after parsing). If a new +/// transient status is added to the const, **also add it to this `matches!` +/// arm and the text-pattern list below**. +/// +/// Note: 429 lives in `TRANSIENT_PROVIDER_HTTP_STATUSES` but is intentionally +/// absent here — `is_rate_limited` handles it separately because 429 responses +/// may carry a `Retry-After` header that `parse_retry_after_ms` uses to pick a +/// precise backoff rather than the default exponential schedule. pub(crate) fn is_upstream_unhealthy(err: &anyhow::Error) -> bool { + if let Some(reqwest_err) = err.downcast_ref::() { + if let Some(status) = reqwest_err.status() { + if matches!(status.as_u16(), 408 | 502 | 503 | 504) { + return true; + } + } + } let lower = err.to_string().to_lowercase(); lower.contains("no healthy upstream") || lower.contains("upstream unavailable") || lower.contains("service unavailable") + || lower.contains("503 service unavailable") + || lower.contains("408 request timeout") + || lower.contains("502 bad gateway") + || lower.contains("504 gateway timeout") } /// Check if an error is a rate-limit (429) error. diff --git a/src/openhuman/providers/reliable_tests.rs b/src/openhuman/providers/reliable_tests.rs index e1d992315..7bf14aae5 100644 --- a/src/openhuman/providers/reliable_tests.rs +++ b/src/openhuman/providers/reliable_tests.rs @@ -799,6 +799,35 @@ fn upstream_unhealthy_does_not_flag_generic_error() { assert!(!is_upstream_unhealthy(&err)); } +// 408/502/504 must also classify as transient — `ops::api_error` formats +// the upstream failure as " API error (): ", and the +// tool-call loop ORs is_rate_limited (429) with is_upstream_unhealthy. Before +// this fix only 503/text-pattern matched; 408/502/504 leaked per-iteration +// Sentry events (CodeRabbit review on #1529, OPENHUMAN-TAURI-T/-2E/-84). +#[test] +fn upstream_unhealthy_detects_408_request_timeout() { + let err = anyhow::anyhow!("OpenAI API error (408 Request Timeout): upstream took too long"); + assert!(is_upstream_unhealthy(&err)); +} + +#[test] +fn upstream_unhealthy_detects_502_bad_gateway() { + let err = anyhow::anyhow!("Anthropic API error (502 Bad Gateway): bad gateway"); + assert!(is_upstream_unhealthy(&err)); +} + +#[test] +fn upstream_unhealthy_detects_504_gateway_timeout() { + let err = anyhow::anyhow!("OpenAI API error (504 Gateway Timeout): upstream timed out"); + assert!(is_upstream_unhealthy(&err)); +} + +#[test] +fn upstream_unhealthy_detects_503_service_unavailable_with_provider_prefix() { + let err = anyhow::anyhow!("OpenAI API error (503 Service Unavailable): backend overloaded"); + assert!(is_upstream_unhealthy(&err)); +} + #[test] fn failure_reason_upstream_unhealthy_wins_over_rate_limited() { // Both rate_limited AND upstream_unhealthy — upstream_unhealthy must win. diff --git a/tests/observability_smoke.rs b/tests/observability_smoke.rs new file mode 100644 index 000000000..47294d49b --- /dev/null +++ b/tests/observability_smoke.rs @@ -0,0 +1,144 @@ +//! Runtime smoke for the Sentry `before_send` filter that drops per-attempt +//! transient-upstream provider failures (OPENHUMAN-TAURI-2E / 84 / T). +//! +//! Unit tests in `src/core/observability.rs` exercise the pure filter +//! function. This integration test wires the actual `sentry::init` → +//! `before_send` → transport chain so we have proof the runtime path +//! behaves as designed: transient events are dropped, permanent events +//! and aggregate `all_exhausted` events still surface. + +use openhuman_core::core::observability::is_transient_provider_http_failure; +use sentry::protocol::Event; +use std::collections::BTreeMap; +use std::sync::Arc; + +fn event_with_tags(tags: &[(&str, &str)]) -> Event<'static> { + let mut event = Event::default(); + let mut t: BTreeMap = BTreeMap::new(); + for (k, v) in tags { + t.insert((*k).to_string(), (*v).to_string()); + } + event.tags = t; + event +} + +/// Drive an envelope-capturing Sentry client through a sequence of events +/// and return how many made it past `before_send`. +/// +/// `sentry::init` mutates the process-global Sentry hub; Cargo runs integration +/// test functions in parallel threads by default, so two `count_captured` calls +/// would otherwise race on the global hub and one test's `capture_event` could +/// land in another test's transport. Serialize the critical section here rather +/// than imposing `--test-threads=1` on the whole binary. +fn count_captured(events: Vec>) -> usize { + static SENTRY_TEST_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(()); + let _guard = SENTRY_TEST_LOCK + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()); + + let transport = sentry::test::TestTransport::new(); + let transport_for_factory = transport.clone(); + let options = sentry::ClientOptions { + dsn: Some( + "https://public@sentry.example.com/1" + .parse() + .expect("dsn parses"), + ), + // Same filter shape the real binary installs in main.rs. + before_send: Some(Arc::new(|event| { + if is_transient_provider_http_failure(&event) { + None + } else { + Some(event) + } + })), + transport: Some(Arc::new(move |_opts: &sentry::ClientOptions| { + transport_for_factory.clone() as Arc + })), + sample_rate: 1.0, + ..sentry::ClientOptions::default() + }; + let _sentry_guard = sentry::init(options); + for event in events { + sentry::capture_event(event); + } + sentry::Hub::current() + .client() + .map(|c| c.flush(Some(std::time::Duration::from_secs(2)))); + transport.fetch_and_clear_envelopes().len() +} + +#[test] +fn drops_per_attempt_429_503_504_408_502() { + // Each of these matches the tag shape `ops::api_error` sets when a + // transient upstream status returns. With the filter installed in + // before_send, none should leak through to the transport. + let events = ["429", "503", "504", "408", "502"] + .into_iter() + .map(|status| { + event_with_tags(&[ + ("domain", "llm_provider"), + ("failure", "non_2xx"), + ("status", status), + ]) + }) + .collect(); + assert_eq!( + count_captured(events), + 0, + "transient per-attempt failures must be filtered in before_send" + ); +} + +#[test] +fn keeps_permanent_failures() { + // 4xx auth / not-found / etc. and 500 internal errors are actionable — + // they must reach Sentry exactly as before. + let events = ["400", "401", "403", "404", "500"] + .into_iter() + .map(|status| { + event_with_tags(&[ + ("domain", "llm_provider"), + ("failure", "non_2xx"), + ("status", status), + ]) + }) + .collect(); + assert_eq!( + count_captured(events), + 5, + "permanent provider failures must reach Sentry" + ); +} + +#[test] +fn keeps_aggregate_all_exhausted_event() { + // The reliable_chat layer fires a single aggregate + // `failure=all_exhausted` event when every provider/model has been + // tried. That's the cascade signal we want — only the per-attempt + // noise gets dropped. + let event = event_with_tags(&[ + ("domain", "llm_provider"), + ("failure", "all_exhausted"), + ("model", "claude-haiku-4-5-20251001"), + ("attempts", "12"), + ]); + assert_eq!( + count_captured(vec![event]), + 1, + "aggregate all_exhausted event must surface for genuine outages" + ); +} + +#[test] +fn keeps_event_missing_status_tag() { + // Belt-and-suspenders: an event with `failure=non_2xx` but no `status` + // tag (e.g. a future call site forgets to attach one) must NOT be + // silently dropped — we'd rather see it and fix the tag emission. + let event = event_with_tags(&[("domain", "llm_provider"), ("failure", "non_2xx")]); + assert_eq!( + count_captured(vec![event]), + 1, + "event without status tag must not be silently dropped" + ); +}