diff --git a/app/src-tauri/Cargo.lock b/app/src-tauri/Cargo.lock index 662d42c4f..8ca6d1ce6 100644 --- a/app/src-tauri/Cargo.lock +++ b/app/src-tauri/Cargo.lock @@ -220,7 +220,7 @@ dependencies = [ "objc2-foundation 0.3.2", "parking_lot", "percent-encoding", - "windows-sys 0.59.0", + "windows-sys 0.60.2", "x11rb", ] @@ -1678,7 +1678,7 @@ dependencies = [ "libc", "option-ext", "redox_users 0.5.2", - "windows-sys 0.59.0", + "windows-sys 0.61.2", ] [[package]] @@ -1989,7 +1989,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb" dependencies = [ "libc", - "windows-sys 0.59.0", + "windows-sys 0.61.2", ] [[package]] @@ -3974,7 +3974,7 @@ version = "0.50.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7957b9740744892f114936ab4a57b3f487491bbeafaf8083688b16841a4240e5" dependencies = [ - "windows-sys 0.59.0", + "windows-sys 0.61.2", ] [[package]] @@ -4108,6 +4108,17 @@ dependencies = [ "objc2-foundation 0.3.2", ] +[[package]] +name = "objc2-contacts" +version = "0.3.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b034b578389f89a85c055eacc8d8b368be5f04a6c1b07f672bf3aec21d0ef621" +dependencies = [ + "block2 0.6.2", + "objc2 0.6.4", + "objc2-foundation 0.3.2", +] + [[package]] name = "objc2-core-data" version = "0.2.2" @@ -4466,6 +4477,7 @@ dependencies = [ "async-trait", "axum", "base64 0.22.1", + "block2 0.6.2", "chacha20poly1305", "chrono", "chrono-tz", @@ -4483,6 +4495,7 @@ dependencies = [ "fs2", "futures", "futures-util", + "glob", "hex", "hmac 0.12.1", "hostname", @@ -4493,6 +4506,9 @@ dependencies = [ "log", "mail-parser", "nu-ansi-term 0.46.0", + "objc2 0.6.4", + "objc2-contacts", + "objc2-foundation 0.3.2", "once_cell", "opentelemetry", "opentelemetry-otlp", @@ -4538,6 +4554,7 @@ dependencies = [ "urlencoding", "uuid", "wait-timeout", + "walkdir", "webpki-roots 1.0.6", "whisper-rs", "xz2", @@ -5420,7 +5437,7 @@ dependencies = [ "once_cell", "socket2", "tracing", - "windows-sys 0.59.0", + "windows-sys 0.60.2", ] [[package]] @@ -5846,7 +5863,7 @@ dependencies = [ "errno", "libc", "linux-raw-sys", - "windows-sys 0.59.0", + "windows-sys 0.61.2", ] [[package]] @@ -5905,7 +5922,7 @@ dependencies = [ "security-framework", "security-framework-sys", "webpki-root-certs", - "windows-sys 0.59.0", + "windows-sys 0.61.2", ] [[package]] @@ -7236,7 +7253,7 @@ dependencies = [ "getrandom 0.4.2", "once_cell", "rustix", - "windows-sys 0.59.0", + "windows-sys 0.61.2", ] [[package]] @@ -8534,7 +8551,7 @@ version = "0.1.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22" dependencies = [ - "windows-sys 0.48.0", + "windows-sys 0.61.2", ] [[package]] diff --git a/src/api/rest.rs b/src/api/rest.rs index 54eef08c7..064ddb98a 100644 --- a/src/api/rest.rs +++ b/src/api/rest.rs @@ -382,14 +382,44 @@ impl BackendOAuthClient { request = request.json(&body); } - let response = request - .send() - .await - .with_context(|| format!("backend request {} {}", method.as_str(), url.path()))?; + let response = request.send().await.map_err(|e| { + crate::core::observability::report_error( + e.to_string().as_str(), + "backend_api", + "authed_json", + &[ + ("method", method.as_str()), + ("path", url.path()), + ("failure", "transport"), + ], + ); + anyhow::Error::new(e).context(format!( + "backend request {} {}", + method.as_str(), + url.path() + )) + })?; let status = response.status(); let text = response.text().await.unwrap_or_default(); if !status.is_success() { + let status_str = status.as_u16().to_string(); + crate::core::observability::report_error( + format!( + "{} {} failed ({status}): {text}", + method.as_str(), + url.path() + ) + .as_str(), + "backend_api", + "authed_json", + &[ + ("method", method.as_str()), + ("path", url.path()), + ("status", status_str.as_str()), + ("failure", "non_2xx"), + ], + ); anyhow::bail!( "{} {} failed ({status}): {text}", method.as_str(), diff --git a/src/core/jsonrpc.rs b/src/core/jsonrpc.rs index 1c903bcf5..8926d8606 100644 --- a/src/core/jsonrpc.rs +++ b/src/core/jsonrpc.rs @@ -55,7 +55,19 @@ pub async fn rpc_handler(State(state): State, Json(req): Json { - tracing::info!("[rpc] {} -> err ({}ms): {}", method, ms, message); + // Session-expired bubbles up as an "error" but is an expected + // boundary condition (auth handler clears the local token and the + // UI re-auths). Don't spam Sentry with it. + if !is_session_expired_error(&message) { + crate::core::observability::report_error( + message.as_str(), + "rpc", + "invoke_method", + &[("method", method.as_str()), ("elapsed_ms", &ms.to_string())], + ); + } else { + tracing::info!("[rpc] {} -> err ({}ms): {}", method, ms, message); + } ( StatusCode::OK, Json(RpcFailure { diff --git a/src/core/mod.rs b/src/core/mod.rs index bd299d9ec..b879d4522 100644 --- a/src/core/mod.rs +++ b/src/core/mod.rs @@ -16,6 +16,7 @@ pub mod event_bus; pub mod jsonrpc; pub mod logging; pub mod memory_cli; +pub mod observability; pub mod rpc_log; pub mod shutdown; pub mod socketio; diff --git a/src/core/observability.rs b/src/core/observability.rs new file mode 100644 index 000000000..7b3d6e88b --- /dev/null +++ b/src/core/observability.rs @@ -0,0 +1,85 @@ +//! Centralised error reporting for the core. +//! +//! Wraps `tracing::error!` (which the global subscriber forwards to Sentry via +//! `sentry-tracing`) inside a `sentry::with_scope` so each captured event +//! carries consistent tags identifying the failing domain/operation plus any +//! callsite-specific context (session id, request id, tool name, …). +//! +//! Why this helper exists: errors that bubble up as `Result::Err` without ever +//! being logged at error level never reach Sentry. The agent-turn path is the +//! canonical example — `run_single` used to publish a `DomainEvent::AgentError` +//! and return `Err(_)`, but Sentry never saw it. Funnel error sites through +//! `report_error` so they show up tagged and grep-friendly in Sentry. + +use std::fmt::Display; + +/// A `(key, value)` pair attached as a Sentry tag. Tags are short, indexed, +/// and filterable in the Sentry UI — prefer them over free-form fields for +/// anything you'd want to facet on (`error_kind`, `tool_name`, `method`). +pub type Tag<'a> = (&'a str, &'a str); + +/// Capture an error to Sentry with structured tags. +/// +/// `domain` and `operation` are required and become tags `domain:<…>` and +/// `operation:<…>`. `extra` is an optional list of extra tag pairs. The error +/// itself is rendered via `Display` and emitted as a `tracing::error!` event, +/// which the Sentry tracing layer turns into a Sentry event under the active +/// scope. +/// +/// Use stable, low-cardinality values for tag keys/values so Sentry can group +/// and aggregate. High-cardinality data (full IDs, payloads) belongs in the +/// error message body, not in tags. +pub fn report_error( + err: &E, + domain: &str, + operation: &str, + extra: &[Tag<'_>], +) { + let message = err.to_string(); + sentry::with_scope( + |scope| { + scope.set_tag("domain", domain); + scope.set_tag("operation", operation); + for (k, v) in extra { + scope.set_tag(*k, *v); + } + }, + || { + tracing::error!( + domain = domain, + operation = operation, + error = %message, + "[observability] {domain}.{operation} failed: {message}" + ); + }, + ); +} + +#[cfg(test)] +mod tests { + use super::*; + + /// Helper must accept `&anyhow::Error`, `&dyn std::error::Error`, and + /// plain `&str` — the three shapes that show up at error sites today. + #[test] + fn report_error_accepts_common_error_shapes() { + let anyhow_err = anyhow::anyhow!("boom"); + report_error(&anyhow_err, "test", "anyhow_shape", &[]); + + let io_err = std::io::Error::other("io failed"); + report_error(&io_err, "test", "io_shape", &[("kind", "io")]); + + report_error("plain message", "test", "str_shape", &[]); + } + + #[test] + fn report_error_does_not_panic_with_many_tags() { + let err = anyhow::anyhow!("multi-tag"); + report_error( + &err, + "test", + "multi_tag", + &[("a", "1"), ("b", "2"), ("c", "3"), ("d", "4")], + ); + } +} diff --git a/src/openhuman/agent/harness/session/runtime.rs b/src/openhuman/agent/harness/session/runtime.rs index 44e3f58c4..14f24e25a 100644 --- a/src/openhuman/agent/harness/session/runtime.rs +++ b/src/openhuman/agent/harness/session/runtime.rs @@ -461,6 +461,21 @@ impl Agent { "Prompt flagged for security review and was not processed." } }; + let action_tag = match guard.action { + PromptEnforcementAction::Allow => "allow", + PromptEnforcementAction::Blocked => "blocked", + PromptEnforcementAction::ReviewBlocked => "review_blocked", + }; + crate::core::observability::report_error( + user_message, + "agent", + "prompt_injection_blocked", + &[ + ("session_id", self.event_session_id()), + ("channel", self.event_channel()), + ("action", action_tag), + ], + ); publish_global(DomainEvent::AgentError { session_id: self.event_session_id().to_string(), message: user_message.to_string(), @@ -487,6 +502,16 @@ impl Agent { } Err(err) => { let sanitized_message = Self::sanitize_event_error_message(&err); + crate::core::observability::report_error( + &err, + "agent", + "run_single", + &[ + ("session_id", self.event_session_id()), + ("channel", self.event_channel()), + ("error_kind", sanitized_message.as_str()), + ], + ); publish_global(DomainEvent::AgentError { session_id: self.event_session_id().to_string(), message: sanitized_message, diff --git a/src/openhuman/agent/harness/tool_loop.rs b/src/openhuman/agent/harness/tool_loop.rs index 1d6cd0179..e16f9ede9 100644 --- a/src/openhuman/agent/harness/tool_loop.rs +++ b/src/openhuman/agent/harness/tool_loop.rs @@ -191,26 +191,42 @@ pub(crate) async fn run_tool_call_loop( utilization_pct, reason, } => { - tracing::error!( - iteration, - utilization_pct, - "[agent_loop] context exhausted, aborting: {reason}" + let msg = format!("Context window exhausted ({utilization_pct}% full): {reason}"); + crate::core::observability::report_error( + msg.as_str(), + "agent", + "context_exhausted", + &[ + ("provider", provider_name), + ("model", model), + ("utilization_pct", &utilization_pct.to_string()), + ], ); - anyhow::bail!("Context window exhausted ({utilization_pct}% full): {reason}"); + anyhow::bail!(msg); } } tracing::debug!(iteration, "[agent_loop] sending LLM request"); let image_marker_count = multimodal::count_image_markers(history); if image_marker_count > 0 && !provider.supports_vision() { - return Err(ProviderCapabilityError { + let cap_err = ProviderCapabilityError { provider: provider_name.to_string(), capability: "vision".to_string(), message: format!( "received {image_marker_count} image marker(s), but this provider does not support vision input" ), - } - .into()); + }; + crate::core::observability::report_error( + &cap_err, + "agent", + "provider_capability", + &[ + ("provider", provider_name), + ("capability", "vision"), + ("model", model), + ], + ); + return Err(cap_err.into()); } let prepared_messages = @@ -346,6 +362,16 @@ pub(crate) async fn run_tool_call_loop( ) } Err(e) => { + crate::core::observability::report_error( + &e, + "agent", + "provider_chat", + &[ + ("provider", provider_name), + ("model", model), + ("iteration", &(iteration + 1).to_string()), + ], + ); return Err(e); } }; @@ -626,19 +652,33 @@ pub(crate) async fn run_tool_call_loop( } } Ok(Err(e)) => { - tracing::error!( - iteration, - tool = call.name.as_str(), - "[agent_loop] tool execution failed: {e}" + crate::core::observability::report_error( + &e, + "tool", + "execute", + &[ + ("tool", call.name.as_str()), + ("outcome", "failed"), + ("iteration", &(iteration + 1).to_string()), + ], ); (format!("Error executing {}: {e}", call.name), false) } Err(_) => { - tracing::error!( - iteration, - tool = call.name.as_str(), - secs = timeout_secs, - "[agent_loop] tool execution timed out" + let msg = format!( + "tool '{}' timed out after {} seconds", + call.name, timeout_secs + ); + crate::core::observability::report_error( + msg.as_str(), + "tool", + "execute", + &[ + ("tool", call.name.as_str()), + ("outcome", "timeout"), + ("timeout_secs", &timeout_secs.to_string()), + ("iteration", &(iteration + 1).to_string()), + ], ); ( format!( diff --git a/src/openhuman/channels/runtime/dispatch.rs b/src/openhuman/channels/runtime/dispatch.rs index 9929cba4c..9ab09df75 100644 --- a/src/openhuman/channels/runtime/dispatch.rs +++ b/src/openhuman/channels/runtime/dispatch.rs @@ -765,6 +765,15 @@ pub(crate) async fn process_channel_message( let active_provider = match get_or_create_provider(ctx.as_ref(), &route.provider).await { Ok(provider) => provider, Err(err) => { + crate::core::observability::report_error( + &err, + "channels", + "provider_init", + &[ + ("channel", msg.channel.as_str()), + ("provider", route.provider.as_str()), + ], + ); let safe_err = providers::sanitize_api_error(&err.to_string()); let message = format!( "⚠️ Failed to initialize provider `{}`. Please run `/models` to choose another provider.\nDetails: {safe_err}", @@ -1137,6 +1146,15 @@ pub(crate) async fn process_channel_message( " ❌ LLM error after {}ms: {e}", started_at.elapsed().as_millis() ); + crate::core::observability::report_error( + &e, + "channels", + "dispatch_llm_error", + &[ + ("channel", msg.channel.as_str()), + ("provider", route.provider.as_str()), + ], + ); if let Some(channel) = target_channel.as_ref() { if let Some(ref draft_id) = draft_message_id { let _ = channel @@ -1165,6 +1183,15 @@ pub(crate) async fn process_channel_message( timeout_msg, started_at.elapsed().as_millis() ); + crate::core::observability::report_error( + timeout_msg.as_str(), + "channels", + "dispatch_llm_timeout", + &[ + ("channel", msg.channel.as_str()), + ("timeout_secs", &ctx.message_timeout_secs.to_string()), + ], + ); let error_text = "⚠️ Request timed out while waiting for the model. Please try again.".to_string(); if let Some(channel) = target_channel.as_ref() { diff --git a/src/openhuman/composio/client.rs b/src/openhuman/composio/client.rs index 9e1dd95ab..af30503c7 100644 --- a/src/openhuman/composio/client.rs +++ b/src/openhuman/composio/client.rs @@ -323,6 +323,17 @@ impl ComposioClient { status, &body_text[..body_text.len().min(300)] ); + let status_str = status.as_u16().to_string(); + crate::core::observability::report_error( + format!("Backend returned {} for DELETE {}", status, url).as_str(), + "composio", + "delete", + &[ + ("path", path), + ("status", status_str.as_str()), + ("failure", "non_2xx"), + ], + ); anyhow::bail!("Backend returned {} for DELETE {}", status, url); } @@ -331,6 +342,12 @@ impl ComposioClient { let msg = envelope .error .unwrap_or_else(|| "unknown backend error".into()); + crate::core::observability::report_error( + msg.as_str(), + "composio", + "delete", + &[("path", path), ("failure", "envelope_error")], + ); anyhow::bail!("Backend error for DELETE {}: {}", url, msg); } envelope.data.ok_or_else(|| { diff --git a/src/openhuman/cron/bus.rs b/src/openhuman/cron/bus.rs index 4a525038f..afc49bac2 100644 --- a/src/openhuman/cron/bus.rs +++ b/src/openhuman/cron/bus.rs @@ -64,20 +64,32 @@ impl EventHandler for CronDeliverySubscriber { ); } Err(e) => { - tracing::warn!( - job_id = %job_id, - channel = %channel_lower, - error = %e, - "[cron] delivery failed" + crate::core::observability::report_error( + &e, + "cron", + "delivery", + &[ + ("job_id", job_id.as_str()), + ("channel", channel_lower.as_str()), + ("failure", "channel_send"), + ], ); } } } else { - tracing::warn!( - job_id = %job_id, - channel = %channel_lower, - available = ?self.channels_by_name.keys().collect::>(), - "[cron] no matching channel found for delivery" + let msg = format!( + "no matching channel '{}' found for cron delivery (job {})", + channel_lower, job_id + ); + crate::core::observability::report_error( + msg.as_str(), + "cron", + "delivery", + &[ + ("job_id", job_id.as_str()), + ("channel", channel_lower.as_str()), + ("failure", "channel_missing"), + ], ); } } diff --git a/src/openhuman/embeddings/ollama.rs b/src/openhuman/embeddings/ollama.rs index 1bc89cf80..5d5cf4e65 100644 --- a/src/openhuman/embeddings/ollama.rs +++ b/src/openhuman/embeddings/ollama.rs @@ -166,17 +166,25 @@ impl EmbeddingProvider for OllamaEmbedding { .send() .await .map_err(|e| { - anyhow::anyhow!( + let message = format!( "ollama embed request failed (is Ollama running at {}?): {e}", self.base_url - ) + ); + crate::core::observability::report_error( + message.as_str(), + "embeddings", + "ollama_embed", + &[("model", self.model.as_str()), ("failure", "transport")], + ); + anyhow::anyhow!(message) })?; if !resp.status().is_success() { let status = resp.status(); + let status_str = status.as_u16().to_string(); let body = resp.text().await.unwrap_or_default(); let detail = body.trim(); - anyhow::bail!( + let message = format!( "ollama embed failed with status {status}{}", if detail.is_empty() { String::new() @@ -184,6 +192,17 @@ impl EmbeddingProvider for OllamaEmbedding { format!(": {detail}") } ); + crate::core::observability::report_error( + message.as_str(), + "embeddings", + "ollama_embed", + &[ + ("model", self.model.as_str()), + ("status", status_str.as_str()), + ("failure", "non_2xx"), + ], + ); + anyhow::bail!(message); } let payload: OllamaEmbedResponse = resp diff --git a/src/openhuman/embeddings/openai.rs b/src/openhuman/embeddings/openai.rs index 5f8471680..61e0ec590 100644 --- a/src/openhuman/embeddings/openai.rs +++ b/src/openhuman/embeddings/openai.rs @@ -118,12 +118,24 @@ impl EmbeddingProvider for OpenAiEmbedding { if !resp.status().is_success() { let status = resp.status(); + let status_str = status.as_u16().to_string(); let text = resp.text().await.unwrap_or_default(); tracing::debug!( target: "openai::embed", "[openai] embed error: status={status}, body={text}" ); - anyhow::bail!("Embedding API error {status}: {text}"); + let message = format!("Embedding API error {status}: {text}"); + crate::core::observability::report_error( + message.as_str(), + "embeddings", + "openai_embed", + &[ + ("model", self.model.as_str()), + ("status", status_str.as_str()), + ("failure", "non_2xx"), + ], + ); + anyhow::bail!(message); } let json: serde_json::Value = resp.json().await?; diff --git a/src/openhuman/integrations/client.rs b/src/openhuman/integrations/client.rs index 583126807..c31a1c557 100644 --- a/src/openhuman/integrations/client.rs +++ b/src/openhuman/integrations/client.rs @@ -69,17 +69,28 @@ impl IntegrationClient { chain.push_str(&s.to_string()); src = s.source(); } - tracing::warn!("[integrations] POST {} failed: {}", url, chain); + crate::core::observability::report_error( + chain.as_str(), + "integrations", + "post", + &[("path", path), ("failure", "transport")], + ); anyhow::anyhow!("POST {} failed: {}", url, chain) })?; let status = resp.status(); if !status.is_success() { let _body_text = resp.text().await.unwrap_or_default(); - tracing::debug!( - "[integrations] POST {} → {} ", - url, - status + let status_str = status.as_u16().to_string(); + crate::core::observability::report_error( + format!("Backend returned {} for POST {}", status, url).as_str(), + "integrations", + "post", + &[ + ("path", path), + ("status", status_str.as_str()), + ("failure", "non_2xx"), + ], ); anyhow::bail!("Backend returned {} for POST {}", status, url); } @@ -89,6 +100,12 @@ impl IntegrationClient { let msg = envelope .error .unwrap_or_else(|| "unknown backend error".into()); + crate::core::observability::report_error( + msg.as_str(), + "integrations", + "post", + &[("path", path), ("failure", "envelope_error")], + ); anyhow::bail!("Backend error for POST {}: {}", url, msg); } envelope @@ -115,17 +132,28 @@ impl IntegrationClient { chain.push_str(&s.to_string()); src = s.source(); } - tracing::warn!("[integrations] GET {} failed: {}", url, chain); + crate::core::observability::report_error( + chain.as_str(), + "integrations", + "get", + &[("path", path), ("failure", "transport")], + ); anyhow::anyhow!("GET {} failed: {}", url, chain) })?; let status = resp.status(); if !status.is_success() { let _body_text = resp.text().await.unwrap_or_default(); - tracing::debug!( - "[integrations] GET {} → {} ", - url, - status + let status_str = status.as_u16().to_string(); + crate::core::observability::report_error( + format!("Backend returned {} for GET {}", status, url).as_str(), + "integrations", + "get", + &[ + ("path", path), + ("status", status_str.as_str()), + ("failure", "non_2xx"), + ], ); anyhow::bail!("Backend returned {} for GET {}", status, url); } @@ -135,6 +163,12 @@ impl IntegrationClient { let msg = envelope .error .unwrap_or_else(|| "unknown backend error".into()); + crate::core::observability::report_error( + msg.as_str(), + "integrations", + "get", + &[("path", path), ("failure", "envelope_error")], + ); anyhow::bail!("Backend error for GET {}: {}", url, msg); } envelope diff --git a/src/openhuman/memory/ingestion/queue.rs b/src/openhuman/memory/ingestion/queue.rs index dc0823d07..28a8bd8ec 100644 --- a/src/openhuman/memory/ingestion/queue.rs +++ b/src/openhuman/memory/ingestion/queue.rs @@ -171,9 +171,14 @@ async fn ingestion_worker( true } Err(e) => { - log::error!( - "[memory:ingestion_queue] extraction failed namespace={namespace} \ - doc_id={document_id} title={title}: {e}", + crate::core::observability::report_error( + &e, + "memory", + "ingestion_extract", + &[ + ("namespace", namespace.as_str()), + ("doc_id", document_id.as_str()), + ], ); false } diff --git a/src/openhuman/memory/tree/jobs/worker.rs b/src/openhuman/memory/tree/jobs/worker.rs index 6f563968d..5fa99ec9d 100644 --- a/src/openhuman/memory/tree/jobs/worker.rs +++ b/src/openhuman/memory/tree/jobs/worker.rs @@ -65,7 +65,12 @@ pub fn start(config: Config) { } } Err(err) => { - log::error!("[memory_tree::jobs] worker={} loop error: {:#}", idx, err); + crate::core::observability::report_error( + &err, + "memory", + "tree_jobs_worker", + &[("worker_idx", &idx.to_string())], + ); tokio::time::sleep(Duration::from_secs(1)).await; } } diff --git a/src/openhuman/providers/compatible.rs b/src/openhuman/providers/compatible.rs index ac3cb1e51..1e561a1e4 100644 --- a/src/openhuman/providers/compatible.rs +++ b/src/openhuman/providers/compatible.rs @@ -361,8 +361,23 @@ impl OpenAiCompatibleProvider { .await?; if !response.status().is_success() { + let status = response.status(); + let status_str = status.as_u16().to_string(); let error = response.text().await?; - anyhow::bail!("{} Responses API error: {error}", self.name); + let sanitized = super::sanitize_api_error(&error); + let message = format!("{} Responses API error: {sanitized}", self.name); + crate::core::observability::report_error( + message.as_str(), + "llm_provider", + "responses_api", + &[ + ("provider", self.name.as_str()), + ("model", model), + ("status", status_str.as_str()), + ("failure", "non_2xx"), + ], + ); + anyhow::bail!(message); } let body = response.text().await?; @@ -675,17 +690,28 @@ impl OpenAiCompatibleProvider { if !response.status().is_success() { let status = response.status(); + let status_str = status.as_u16().to_string(); let body = response.text().await.unwrap_or_default(); // Sanitize the upstream error body so we don't leak user // prompts, tool arguments, or credentials the backend // echoed back into the anyhow chain / logs. let sanitized = super::sanitize_api_error(&body); - anyhow::bail!( + let message = format!( "{} streaming API error ({}): {}", - self.name, - status, - sanitized + self.name, status, sanitized ); + crate::core::observability::report_error( + message.as_str(), + "llm_provider", + "streaming_chat", + &[ + ("provider", self.name.as_str()), + ("model", native_request.model.as_str()), + ("status", status_str.as_str()), + ("failure", "non_2xx"), + ], + ); + anyhow::bail!(message); } // Some OpenAI-compatible backends (and our e2e mock) accept @@ -1127,7 +1153,20 @@ impl Provider for OpenAiCompatibleProvider { }); } - anyhow::bail!("{} API error ({status}): {sanitized}", self.name); + let status_str = status.as_u16().to_string(); + let message = format!("{} API error ({status}): {sanitized}", self.name); + crate::core::observability::report_error( + message.as_str(), + "llm_provider", + "chat_completions", + &[ + ("provider", self.name.as_str()), + ("model", model), + ("status", status_str.as_str()), + ("failure", "non_2xx"), + ], + ); + anyhow::bail!(message); } let body = response.text().await?; @@ -1511,7 +1550,20 @@ impl Provider for OpenAiCompatibleProvider { }); } - anyhow::bail!("{} API error ({status}): {sanitized}", self.name); + let status_str = status.as_u16().to_string(); + let message = format!("{} API error ({status}): {sanitized}", self.name); + crate::core::observability::report_error( + message.as_str(), + "llm_provider", + "native_chat", + &[ + ("provider", self.name.as_str()), + ("model", model), + ("status", status_str.as_str()), + ("failure", "non_2xx"), + ], + ); + anyhow::bail!(message); } let response_bytes = response.bytes().await?; @@ -1575,6 +1627,8 @@ impl Provider for OpenAiCompatibleProvider { let url = self.chat_completions_url(); let client = self.http_client(); let auth_header = self.auth_header.clone(); + let provider_name = self.name.clone(); + let model_owned = model.to_string(); // Use a channel to bridge the async HTTP response to the stream let (tx, rx) = tokio::sync::mpsc::channel::>(100); @@ -1599,6 +1653,16 @@ impl Provider for OpenAiCompatibleProvider { let response = match req_builder.send().await { Ok(r) => r, Err(e) => { + crate::core::observability::report_error( + e.to_string().as_str(), + "llm_provider", + "stream_chat", + &[ + ("provider", provider_name.as_str()), + ("model", model_owned.as_str()), + ("failure", "transport"), + ], + ); let _ = tx.send(Err(StreamError::Http(e))).await; return; } @@ -1607,13 +1671,25 @@ impl Provider for OpenAiCompatibleProvider { // Check status if !response.status().is_success() { let status = response.status(); - let error = match response.text().await { + let status_str = status.as_u16().to_string(); + let raw_error = match response.text().await { Ok(e) => e, Err(_) => format!("HTTP error: {}", status), }; - let _ = tx - .send(Err(StreamError::Provider(format!("{}: {}", status, error)))) - .await; + let sanitized_error = super::sanitize_api_error(&raw_error); + let message = format!("{}: {}", status, sanitized_error); + crate::core::observability::report_error( + message.as_str(), + "llm_provider", + "stream_chat", + &[ + ("provider", provider_name.as_str()), + ("model", model_owned.as_str()), + ("status", status_str.as_str()), + ("failure", "non_2xx"), + ], + ); + let _ = tx.send(Err(StreamError::Provider(message))).await; return; } diff --git a/src/openhuman/providers/ops.rs b/src/openhuman/providers/ops.rs index 74ffde22c..8418df5e5 100644 --- a/src/openhuman/providers/ops.rs +++ b/src/openhuman/providers/ops.rs @@ -137,14 +137,30 @@ pub fn format_anyhow_chain(err: &anyhow::Error) -> String { } /// Build a sanitized provider error from a failed HTTP response. +/// +/// Also reports the failure to Sentry with `provider` and `status` tags so +/// upstream LLM errors are visible in observability without every call-site +/// having to remember to log. pub async fn api_error(provider: &str, response: reqwest::Response) -> anyhow::Error { let status = response.status(); + let status_str = status.as_u16().to_string(); let body = response .text() .await .unwrap_or_else(|_| "".to_string()); let sanitized = sanitize_api_error(&body); - anyhow::anyhow!("{provider} API error ({status}): {sanitized}") + let message = format!("{provider} API error ({status}): {sanitized}"); + crate::core::observability::report_error( + message.as_str(), + "llm_provider", + "api_error", + &[ + ("provider", provider), + ("status", status_str.as_str()), + ("failure", "non_2xx"), + ], + ); + anyhow::anyhow!(message) } /// Create the OpenHuman backend inference client (session JWT only). diff --git a/src/openhuman/providers/reliable.rs b/src/openhuman/providers/reliable.rs index 6163251ca..064ecfb0e 100644 --- a/src/openhuman/providers/reliable.rs +++ b/src/openhuman/providers/reliable.rs @@ -409,10 +409,21 @@ impl Provider for ReliableProvider { } } - anyhow::bail!( + let aggregate = format!( "All providers/models failed. Attempts:\n{}", failures.join("\n") - ) + ); + crate::core::observability::report_error( + aggregate.as_str(), + "llm_provider", + "reliable_chat_with_system", + &[ + ("model", model), + ("attempts", &failures.len().to_string()), + ("failure", "all_exhausted"), + ], + ); + anyhow::bail!(aggregate) } async fn chat_with_history( @@ -519,10 +530,21 @@ impl Provider for ReliableProvider { } } - anyhow::bail!( + let aggregate = format!( "All providers/models failed. Attempts:\n{}", failures.join("\n") - ) + ); + crate::core::observability::report_error( + aggregate.as_str(), + "llm_provider", + "reliable_chat_with_history", + &[ + ("model", model), + ("attempts", &failures.len().to_string()), + ("failure", "all_exhausted"), + ], + ); + anyhow::bail!(aggregate) } fn supports_native_tools(&self) -> bool { @@ -665,10 +687,21 @@ impl Provider for ReliableProvider { } } - anyhow::bail!( + let aggregate = format!( "All providers/models failed. Attempts:\n{}", failures.join("\n") - ) + ); + crate::core::observability::report_error( + aggregate.as_str(), + "llm_provider", + "reliable_chat", + &[ + ("model", model), + ("attempts", &failures.len().to_string()), + ("failure", "all_exhausted"), + ], + ); + anyhow::bail!(aggregate) } async fn chat_with_tools( @@ -776,10 +809,21 @@ impl Provider for ReliableProvider { } } - anyhow::bail!( + let aggregate = format!( "All providers/models failed. Attempts:\n{}", failures.join("\n") - ) + ); + crate::core::observability::report_error( + aggregate.as_str(), + "llm_provider", + "reliable_chat_with_tools", + &[ + ("model", model), + ("attempts", &failures.len().to_string()), + ("failure", "all_exhausted"), + ], + ); + anyhow::bail!(aggregate) } fn supports_streaming(&self) -> bool { diff --git a/src/openhuman/skills/ops_install.rs b/src/openhuman/skills/ops_install.rs index 9a2a171c2..76a5cf5a7 100644 --- a/src/openhuman/skills/ops_install.rs +++ b/src/openhuman/skills/ops_install.rs @@ -9,6 +9,24 @@ use super::ops_discover::{discover_skills_inner, is_workspace_trusted}; use super::ops_parse::parse_skill_md_str; use super::ops_types::{SkillFrontmatter, SkillScope, MAX_NAME_LEN, SKILL_MD}; +/// Strip userinfo, query, and fragment from a URL for safe inclusion in +/// observability tags. Returns `://[:]` on success, +/// or `""` on parse failure. Never returns the raw URL — even +/// validated install URLs may carry signed query params or embedded creds we +/// don't want flowing to Sentry. +fn redact_url(raw: &str) -> String { + match url::Url::parse(raw) { + Ok(u) => { + let scheme = u.scheme(); + let host = u.host_str().unwrap_or(""); + let port = u.port().map(|p| format!(":{p}")).unwrap_or_default(); + let path = u.path(); + format!("{scheme}://{host}{port}{path}") + } + Err(_) => "".to_string(), + } +} + /// Default wall-clock budget for the SKILL.md fetch. pub const DEFAULT_INSTALL_TIMEOUT_SECS: u64 = 60; /// Hard ceiling callers can request via `timeout_secs`. @@ -116,9 +134,12 @@ pub async fn install_skill_from_url( // tracked separately. validate_resolved_host(&fetch_url).await?; + let redacted_raw_url = redact_url(&raw_url); + let redacted_fetch_url = redact_url(&fetch_url); + tracing::debug!( - raw_url = %raw_url, - fetch_url = %fetch_url, + raw_url = %redacted_raw_url, + fetch_url = %redacted_fetch_url, workspace = %workspace_dir.display(), timeout_secs = timeout_secs, "[skills] install_skill_from_url: entry" @@ -138,26 +159,50 @@ pub async fn install_skill_from_url( .map_err(|e| format!("fetch failed: build http client: {e}"))?; tracing::info!( - fetch_url = %fetch_url, + fetch_url = %redacted_fetch_url, "[skills] install_skill_from_url: fetching SKILL.md" ); let response = match client.get(&fetch_url).send().await { Ok(resp) => resp, Err(e) => { - if e.is_timeout() { - return Err(format!("fetch timed out after {timeout_secs}s")); - } - return Err(format!("fetch failed: {e}")); + let (failure, msg) = if e.is_timeout() { + ("timeout", format!("fetch timed out after {timeout_secs}s")) + } else { + ("transport", format!("fetch failed: {e}")) + }; + crate::core::observability::report_error( + msg.as_str(), + "skills", + "install_fetch", + &[("url", redacted_fetch_url.as_str()), ("failure", failure)], + ); + return Err(msg); } }; let status = response.status(); if !status.is_success() { - return Err(format!( + let status_str = status.as_u16().to_string(); + let msg = format!( "fetch failed: {fetch_url} returned status {}", status.as_u16() - )); + ); + let report_msg = format!( + "fetch failed: {redacted_fetch_url} returned status {}", + status.as_u16() + ); + crate::core::observability::report_error( + report_msg.as_str(), + "skills", + "install_fetch", + &[ + ("url", redacted_fetch_url.as_str()), + ("status", status_str.as_str()), + ("failure", "non_2xx"), + ], + ); + return Err(msg); } if let Some(len) = response.content_length() { @@ -278,8 +323,8 @@ pub async fn install_skill_from_url( .collect(); tracing::info!( - raw_url = %raw_url, - fetch_url = %fetch_url, + raw_url = %redacted_raw_url, + fetch_url = %redacted_fetch_url, slug = %slug, bytes = content.len(), new_count = new_skills.len(), diff --git a/src/openhuman/update/core.rs b/src/openhuman/update/core.rs index 6dc68e166..f364e0815 100644 --- a/src/openhuman/update/core.rs +++ b/src/openhuman/update/core.rs @@ -99,17 +99,34 @@ pub async fn check_available() -> Result { .header("Accept", "application/vnd.github+json") .send() .await - .map_err(|e| format!("failed to fetch latest release: {e}"))?; + .map_err(|e| { + let msg = format!("failed to fetch latest release: {e}"); + crate::core::observability::report_error( + msg.as_str(), + "update", + "check_releases", + &[("failure", "transport")], + ); + msg + })?; if !response.status().is_success() { let status = response.status(); + let status_str = status.as_u16().to_string(); let body = response.text().await.unwrap_or_else(|_| "(no body)".into()); log::warn!( "[update] GitHub API returned {}: {}", status, &body[..body.len().min(200)] ); - return Err(format!("GitHub API error: {status}")); + let msg = format!("GitHub API error: {status}"); + crate::core::observability::report_error( + msg.as_str(), + "update", + "check_releases", + &[("status", status_str.as_str()), ("failure", "non_2xx")], + ); + return Err(msg); } let release: GitHubRelease = response @@ -180,14 +197,32 @@ pub async fn download_and_stage_with_version( .build() .map_err(|e| format!("failed to build HTTP client: {e}"))?; - let response = client - .get(download_url) - .send() - .await - .map_err(|e| format!("failed to download update: {e}"))?; + let response = client.get(download_url).send().await.map_err(|e| { + let msg = format!("failed to download update: {e}"); + crate::core::observability::report_error( + msg.as_str(), + "update", + "download", + &[("asset", asset_name), ("failure", "transport")], + ); + msg + })?; if !response.status().is_success() { - return Err(format!("download failed with status {}", response.status())); + let status = response.status(); + let status_str = status.as_u16().to_string(); + let msg = format!("download failed with status {}", status); + crate::core::observability::report_error( + msg.as_str(), + "update", + "download", + &[ + ("asset", asset_name), + ("status", status_str.as_str()), + ("failure", "non_2xx"), + ], + ); + return Err(msg); } let bytes = response diff --git a/src/openhuman/webhooks/bus.rs b/src/openhuman/webhooks/bus.rs index 50562b3e2..2a86e2a5f 100644 --- a/src/openhuman/webhooks/bus.rs +++ b/src/openhuman/webhooks/bus.rs @@ -98,7 +98,15 @@ impl EventHandler for WebhookRequestSubscriber { ); let decoded = decode_webhook_body(&request.body); if let Err(e) = &decoded { - tracing::error!("[webhook] rejecting — failed to decode body: {}", e); + crate::core::observability::report_error( + e.to_string().as_str(), + "webhooks", + "decode_body", + &[ + ("tunnel", tunnel_uuid.as_str()), + ("method", method.as_str()), + ], + ); let resp = WebhookResponseData { correlation_id: correlation_id.clone(), status_code: 400, @@ -126,14 +134,27 @@ impl EventHandler for WebhookRequestSubscriber { let (resp, err) = match result { Ok(Ok(output)) => (build_agent_response(&corr, 200, &output), None), Ok(Err(e)) => { - tracing::error!("[webhook] agent trigger failed: {}", e); + crate::core::observability::report_error( + e.as_str(), + "webhooks", + "agent_trigger", + &[ + ("correlation_id", corr.as_str()), + ("failure", "agent_error"), + ], + ); ( build_agent_response(&corr, 500, &format!("Agent error: {e}")), Some(e), ) } Err(_) => { - tracing::error!("[webhook] agent trigger timed out (60s)"); + crate::core::observability::report_error( + "agent triage timed out after 60s", + "webhooks", + "agent_trigger", + &[("correlation_id", corr.as_str()), ("failure", "timeout")], + ); ( build_agent_response(&corr, 504, "Agent triage timed out"), Some("timed out after 60s".to_string()),