From 4c40907187bcfb900aadfec9d3bdc0f5b39b0bdb Mon Sep 17 00:00:00 2001 From: Shiv Rossi Date: Thu, 20 Aug 2026 09:47:02 -0600 Subject: [PATCH] feat: add Rite pipeline observability (COD-433) --- Cargo.lock | 1 + crates/rite-cli/Cargo.toml | 1 + crates/rite-cli/src/main.rs | 44 ++++++++++++++- crates/rite-server/src/dispatch.rs | 33 ++++++++++-- crates/rite-server/src/lib.rs | 86 ++++++++++++++++++++++++++++-- 5 files changed, 155 insertions(+), 10 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index bf9fe0e..a71588a 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1077,6 +1077,7 @@ dependencies = [ "clap", "rite-server", "tokio", + "tracing", "tracing-subscriber", ] diff --git a/crates/rite-cli/Cargo.toml b/crates/rite-cli/Cargo.toml index aea42e7..1f9ae2e 100644 --- a/crates/rite-cli/Cargo.toml +++ b/crates/rite-cli/Cargo.toml @@ -13,6 +13,7 @@ axum.workspace = true clap.workspace = true rite-server = { path = "../rite-server" } tokio.workspace = true +tracing.workspace = true tracing-subscriber.workspace = true [lints] diff --git a/crates/rite-cli/src/main.rs b/crates/rite-cli/src/main.rs index 998446e..83186e4 100644 --- a/crates/rite-cli/src/main.rs +++ b/crates/rite-cli/src/main.rs @@ -1,7 +1,7 @@ use std::{net::SocketAddr, path::PathBuf}; use clap::Parser; -use rite_server::{app, configured_state, load_config, start_iris_subscription}; +use rite_server::{RiteConfig, app, configured_state, load_config, start_iris_subscription}; #[derive(Parser)] #[command(name = "rite", about = "Minimal event-to-action runtime")] @@ -19,9 +19,49 @@ async fn main() -> Result<(), Box> { tracing_subscriber::fmt::init(); let cli = Cli::parse(); let config = load_config(&std::fs::read_to_string(cli.config)?)?; + let state = configured_state(&cli.github_webhook_secret, config.clone())?; let listener = tokio::net::TcpListener::bind(cli.listen).await?; - let state = configured_state(&cli.github_webhook_secret, config)?; + for line in startup_summary(listener.local_addr()?, &config) { + tracing::info!("{line}"); + } + start_iris_subscription(&state); axum::serve(listener, app(state)).await?; Ok(()) } + +/// Human-readable, secret-free startup inventory. +fn startup_summary(listen: SocketAddr, config: &RiteConfig) -> Vec { + let mut lines = vec![ + format!("Rite listening on {listen}"), + "source github enabled=true".into(), + ]; + if let Some(iris) = &config.sources.iris { + lines.push(format!("source iris enabled={}", iris.enabled)); + } + lines.extend(config.rites.iter().map(|handler| { + format!( + "handler name={} source={} action=http_post", + handler.name, handler.source + ) + })); + lines +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn startup_summary_is_secret_free() { + let config = load_config("[[rites]]\nname = \"forward\"\nsource = \"github\"\nmatch = { event_type = \"push\" }\naction = { type = \"http_post\", url = \"http://example.test\" }\n[sources.iris]\nenabled = true\nbase_url = \"http://iris.internal\"\n").expect("config"); + let lines = startup_summary("127.0.0.1:8080".parse().expect("address"), &config); + assert!( + lines + .iter() + .any(|line| line.contains("handler name=forward")) + ); + assert!(lines.iter().any(|line| line == "source iris enabled=true")); + assert!(!lines.iter().any(|line| line.contains("iris.internal"))); + } +} diff --git a/crates/rite-server/src/dispatch.rs b/crates/rite-server/src/dispatch.rs index 89b154d..e0ec17e 100644 --- a/crates/rite-server/src/dispatch.rs +++ b/crates/rite-server/src/dispatch.rs @@ -164,22 +164,49 @@ async fn receive_event(state: &AppState, input: RawOperationInput) -> axum::resp Ok(event) => event, Err(error) => return (StatusCode::BAD_REQUEST, error.to_string()).into_response(), }; + state + .metrics + .events_received + .fetch_add(1, std::sync::atomic::Ordering::Relaxed); + tracing::info!(source = %event.source, event_type = %event.event_type, action = ?event.action, "Webhook event received"); let matched = state .handlers .iter() .filter(|handler| handler.matches(&event)) .collect::>(); + if !matched.is_empty() { + state + .metrics + .events_matched + .fetch_add(1, std::sync::atomic::Ordering::Relaxed); + } let mut executed = 0_usize; for handler in &matched { + tracing::info!(handler = %handler.name, "Webhook event matched handler"); match &handler.action { RiteAction::HttpPost { url } => { match state.client.post(url.clone()).json(&event).send().await { - Ok(response) if response.status().is_success() => executed += 1, + Ok(response) if response.status().is_success() => { + executed += 1; + state + .metrics + .actions_succeeded + .fetch_add(1, std::sync::atomic::Ordering::Relaxed); + tracing::info!(handler = %handler.name, status = %response.status(), "rite action completed"); + } Ok(response) => { + state + .metrics + .actions_failed + .fetch_add(1, std::sync::atomic::Ordering::Relaxed); tracing::warn!(handler = %handler.name, status = %response.status(), "rite action returned failure status"); } - Err(error) => { - tracing::warn!(handler = %handler.name, %error, "rite action request failed"); + Err(_error) => { + state + .metrics + .actions_failed + .fetch_add(1, std::sync::atomic::Ordering::Relaxed); + tracing::warn!(handler = %handler.name, "rite action request failed"); } } } diff --git a/crates/rite-server/src/lib.rs b/crates/rite-server/src/lib.rs index 7f5f128..12f7700 100644 --- a/crates/rite-server/src/lib.rs +++ b/crates/rite-server/src/lib.rs @@ -2,9 +2,15 @@ pub mod dispatch; -use std::sync::Arc; +use std::{ + sync::{ + Arc, + atomic::{AtomicU64, Ordering}, + }, + time::Instant, +}; -use axum::{Router, routing::get}; +use axum::{Json, Router, routing::get}; use rite_core::{RiteAction, RiteHandler}; use rite_sources::{github::GitHubSource, iris::IrisSource}; use serde::Deserialize; @@ -40,6 +46,28 @@ pub struct AppState { pub github: Arc, pub client: reqwest::Client, pub iris: Option, + pub metrics: Arc, +} + +/// Process-local counters exposed by `/status`. +pub struct Metrics { + pub events_received: AtomicU64, + pub events_matched: AtomicU64, + pub actions_succeeded: AtomicU64, + pub actions_failed: AtomicU64, + started_at: Instant, +} + +impl Default for Metrics { + fn default() -> Self { + Self { + events_received: AtomicU64::new(0), + events_matched: AtomicU64::new(0), + actions_succeeded: AtomicU64::new(0), + actions_failed: AtomicU64::new(0), + started_at: Instant::now(), + } + } } /// Creates the Rite HTTP application. @@ -53,6 +81,7 @@ pub struct AppState { pub fn app(state: AppState) -> Router { Router::new() .route("/health", get(health)) + .route("/status", get(status)) .merge(dispatch::generated::generated_router()) .with_state(state) } @@ -66,6 +95,19 @@ async fn health() -> &'static str { "ok" } +async fn status( + axum::extract::State(state): axum::extract::State, +) -> Json { + Json(serde_json::json!({ + "events_received": state.metrics.events_received.load(Ordering::Relaxed), + "events_matched": state.metrics.events_matched.load(Ordering::Relaxed), + "actions_succeeded": state.metrics.actions_succeeded.load(Ordering::Relaxed), + "actions_failed": state.metrics.actions_failed.load(Ordering::Relaxed), + "uptime_seconds": state.metrics.started_at.elapsed().as_secs(), + "handlers_loaded": state.handlers.len(), + })) +} + /// Builds state with a GitHub source and TOML-configured handlers. pub fn configured_state(secret: &str, config: RiteConfig) -> rite_core::Result { let iris = config @@ -79,6 +121,7 @@ pub fn configured_state(secret: &str, config: RiteConfig) -> rite_core::Result>(); + if !matched.is_empty() { + metrics.events_matched.fetch_add(1, Ordering::Relaxed); + } + for handler in matched { + tracing::info!(handler = %handler.name, "Iris event matched handler"); let RiteAction::HttpPost { url } = &handler.action; match client.post(url.clone()).json(&event).send().await { Ok(response) if response.status().is_success() => { + metrics.actions_succeeded.fetch_add(1, Ordering::Relaxed); tracing::info!(handler = %handler.name, "Iris event action completed"); } Ok(response) => { + metrics.actions_failed.fetch_add(1, Ordering::Relaxed); tracing::warn!(handler = %handler.name, status = %response.status(), "Iris event action returned failure status"); } - Err(error) => { - tracing::warn!(handler = %handler.name, %error, "Iris event action request failed"); + Err(_error) => { + metrics.actions_failed.fetch_add(1, Ordering::Relaxed); + tracing::warn!(handler = %handler.name, "Iris event action request failed"); } } } @@ -146,6 +203,7 @@ mod tests { let body = br#"{"ref":"refs/heads/main","repository":{"name":"rite"}}"#; let response = app + .clone() .oneshot( Request::post("/event/github") .header(header::CONTENT_TYPE, "application/json") @@ -165,6 +223,24 @@ mod tests { .expect("utf8") .contains("GitHub push to rite") ); + + let status = app + .oneshot( + Request::get("/status") + .body(Body::empty()) + .expect("request"), + ) + .await + .expect("response"); + assert_eq!(status.status(), StatusCode::OK); + let status = to_bytes(status.into_body(), usize::MAX) + .await + .expect("body"); + let status: serde_json::Value = serde_json::from_slice(&status).expect("json"); + assert_eq!(status["events_received"], 1); + assert_eq!(status["events_matched"], 0); + assert_eq!(status["actions_succeeded"], 0); + assert_eq!(status["actions_failed"], 0); } #[tokio::test]