From 28235f36d9efd45b26105b280b9464e19f1ba0b1 Mon Sep 17 00:00:00 2001 From: Jeremy HERGAULT Date: Sun, 6 Sep 2026 17:55:49 +0200 Subject: [PATCH 1/5] feat: implement probe for processor/service Signed-off-by: Jeremy HERGAULT --- prosa/Cargo.toml | 3 +- prosa/examples/my_prosa_settings.yml | 4 +- prosa/src/core/main.rs | 264 +++++++- prosa_book/src/ch01-02-01-observability.md | 77 ++- prosa_utils/Cargo.toml | 3 +- prosa_utils/src/config/observability.rs | 693 ++++++++++++++++++--- 6 files changed, 905 insertions(+), 139 deletions(-) diff --git a/prosa/Cargo.toml b/prosa/Cargo.toml index fc29d02..e9bb869 100644 --- a/prosa/Cargo.toml +++ b/prosa/Cargo.toml @@ -15,7 +15,8 @@ system-metrics = ["dep:memory-stats"] http-proxy = ["dep:async-http-proxy"] openssl = ["dep:openssl", "dep:tokio-openssl", "prosa-utils/config-openssl"] openssl-vendored = ["openssl", "openssl/vendored", "prosa-utils/config-openssl-vendored"] -prometheus = ["dep:prometheus", "prosa-utils/config-observability-prometheus"] +observability-http = ["prosa-utils/config-observability-http"] +prometheus = ["observability-http", "dep:prometheus", "prosa-utils/config-observability-prometheus"] queue = ["prosa-utils/queue"] [[example]] diff --git a/prosa/examples/my_prosa_settings.yml b/prosa/examples/my_prosa_settings.yml index ff920bf..9664755 100644 --- a/prosa/examples/my_prosa_settings.yml +++ b/prosa/examples/my_prosa_settings.yml @@ -1,10 +1,8 @@ name: my-prosa observability: + endpoint: 0.0.0.0:9100 level: INFO - metrics: - prometheus: - endpoint: 0.0.0.0:9100 traces: stdout: level: debug diff --git a/prosa/src/core/main.rs b/prosa/src/core/main.rs index c9d3908..b562093 100644 --- a/prosa/src/core/main.rs +++ b/prosa/src/core/main.rs @@ -20,6 +20,7 @@ use crate::otel::metrics::{Meter, MeterProvider as _}; use crate::otel::trace::TracerProvider as _; use crate::otel::{InstrumentationScope, KeyValue}; use crate::tracing::{debug, info, warn}; +use prosa_utils::config::observability::{HealthCheckCfg, HealthGuard, HealthState}; use prosa_utils::hash::{BuildIntHasher, IntHashMap, IntHashSet}; use std::sync::{ Arc, @@ -57,6 +58,8 @@ where scope_attributes: Vec, #[cfg(feature = "prometheus")] prometheus_registry: prometheus::Registry, + health: Arc, + health_check: HealthCheckCfg, meter_provider: opentelemetry_sdk::metrics::SdkMeterProvider, tracer_provider: opentelemetry_sdk::trace::SdkTracerProvider, stop: Arc, @@ -85,36 +88,30 @@ where internal_tx_queue: mpsc::Sender>, settings: &S, ) -> Main { + let observability = settings.get_observability(); + let health = Arc::new(HealthState::default()); #[cfg(feature = "prometheus")] - { - let prometheus_registry = prometheus::Registry::new(); - let meter_provider = settings - .get_observability() - .build_meter_provider(&prometheus_registry); - - Main { - internal_tx_queue, - name: settings.get_prosa_name(), - scope_attributes: settings.get_observability().get_scope_attributes(), - prometheus_registry, - meter_provider, - tracer_provider: settings.get_observability().build_tracer_provider(), - stop: Arc::new(AtomicBool::new(false)), - } - } + let prometheus_registry = prometheus::Registry::new(); - #[cfg(not(feature = "prometheus"))] - { - let meter_provider = settings.get_observability().build_meter_provider(); - - Main { - internal_tx_queue, - name: settings.get_prosa_name(), - scope_attributes: settings.get_observability().get_scope_attributes(), - meter_provider, - tracer_provider: settings.get_observability().build_tracer_provider(), - stop: Arc::new(AtomicBool::new(false)), - } + #[cfg(feature = "prometheus")] + let meter_provider = + observability.build_meter_provider_with_health(&prometheus_registry, health.clone()); + #[cfg(all(feature = "observability-http", not(feature = "prometheus")))] + let meter_provider = observability.build_meter_provider_with_health(health.clone()); + #[cfg(not(feature = "observability-http"))] + let meter_provider = observability.build_meter_provider(); + + Main { + internal_tx_queue, + name: settings.get_prosa_name(), + scope_attributes: observability.get_scope_attributes(), + #[cfg(feature = "prometheus")] + prometheus_registry, + health, + health_check: observability.get_health_check().clone(), + meter_provider, + tracer_provider: observability.build_tracer_provider(), + stop: Arc::new(AtomicBool::new(false)), } } @@ -222,6 +219,7 @@ where /// Method to stop all processors pub async fn stop(&self, reason: String) -> Result<(), SendError>> { self.stop.store(true, Ordering::Relaxed); + self.health.set_ready(false); Ok(self .internal_tx_queue .send(InternalMainMsg::Shutdown(reason)) @@ -274,6 +272,8 @@ where internal_rx_queue: mpsc::Receiver>, meter: Meter, stop: Arc, + health: Arc, + health_check: HealthCheckCfg, } impl ProcBusParam for MainProc @@ -293,11 +293,31 @@ impl MainProc where M: Sized + Clone + Debug + Tvf + Default + 'static + std::marker::Send + std::marker::Sync, { + fn update_health(&self) { + let processors_ready = self.health_check.required_processors().iter().all(|name| { + name.is_empty() + || self + .processors + .values() + .flat_map(|queues| queues.values()) + .any(|processor| processor.name() == name) + }); + let services_ready = self + .health_check + .required_services() + .iter() + .all(|name| name.is_empty() || self.services.exist_proc_service(name)); + + self.health + .set_ready(!self.stop.load(Ordering::Relaxed) && processors_ready && services_ready); + } + async fn remove_proc(&mut self, proc_id: u32) -> Option> { if let Some(proc) = self.processors.remove(&proc_id) { let mut new_services = (*self.services).clone(); new_services.remove_proc_services(proc_id); self.services = Arc::new(new_services); + self.update_health(); Some(proc) } else { None @@ -313,6 +333,7 @@ where proc_queue.get_queue_id(), ); self.services = Arc::new(new_services); + self.update_health(); Some(proc_queue) } else { None @@ -452,6 +473,8 @@ where let name = main.name().clone(); let meter = main.meter("prosa_main_task_meter"); let stop = main.stop.clone(); + let health = main.health.clone(); + let health_check = main.health_check.clone(); ( main, MainProc { @@ -462,6 +485,8 @@ where internal_rx_queue, meter, stop, + health, + health_check, }, ) } @@ -479,6 +504,19 @@ where } async fn run(mut self) { + let _health_guard = HealthGuard::new(self.health.clone()); + self.update_health(); + + // Monitor readiness + let health = self.health.clone(); + self.meter + .u64_observable_gauge("prosa_ready") + .with_description("Whether ProSA is ready to serve requests") + .with_callback(move |observer| { + observer.observe(u64::from(health.is_ready()), &[]); + }) + .build(); + #[cfg(feature = "system-metrics")] { // Monitor RAM usage @@ -608,6 +646,7 @@ where } prosa_main_record_proc!(); + self.update_health(); }, InternalMainMsg::DeleteProc(proc_id, proc_err) => { if self.remove_proc(proc_id).await.is_some() { @@ -644,6 +683,7 @@ where } } self.services = Arc::new(new_services); + self.update_health(); let _ = service_update.send(self.services.clone()); self.notify_srv_proc().await; } @@ -655,6 +695,7 @@ where new_services.add_service(name, proc_queue.clone()); } self.services = Arc::new(new_services); + self.update_health(); let _ = service_update.send(self.services.clone()); self.notify_srv_proc().await; } @@ -665,6 +706,7 @@ where new_services.remove_service_proc(&name, proc_id); } self.services = Arc::new(new_services); + self.update_health(); let _ = service_update.send(self.services.clone()); self.notify_srv_proc().await; }, @@ -674,11 +716,33 @@ where new_services.remove_service(&name, proc_id, queue_id); } self.services = Arc::new(new_services); + self.update_health(); let _ = service_update.send(self.services.clone()); self.notify_srv_proc().await; }, InternalMainMsg::Config(config) => { info!("Reloading ProSA configuration"); + + let health_check = match config + .config() + .get::("observability.health") + { + Ok(health_check) => Some(health_check), + Err(config::ConfigError::NotFound(_)) => { + Some(HealthCheckCfg::default()) + } + Err(error) => { + warn!("Can't reload health configuration: {error}"); + None + } + }; + if let Some(health_check) = health_check + && self.health_check != health_check + { + self.health_check = health_check; + self.update_health(); + } + for error in self.notify_config_proc_queue(config).await { if let BusError::ProcComm(proc_id, queue_id, _) = error { if queue_id > 0 { @@ -692,6 +756,7 @@ where }, InternalMainMsg::Shutdown(reason) => { warn!("ProSA is stopping: {}", reason); + self.health.set_ready(false); self.stop().await; // The shutdown mecanism will be implemented later @@ -701,6 +766,7 @@ where }, _ = signal::ctrl_c() => { warn!("ProSA is stopping"); + self.health.set_ready(false); self.stop().await; // The shutdown mecanism will be implemented later @@ -710,3 +776,145 @@ where } } } + +#[cfg(test)] +mod tests { + use super::*; + use crate::core::{proc::ProcParam, service::ProcService, settings::Settings}; + use prosa_utils::config::observability::Observability; + use prosa_utils::msg::simple_string_tvf::SimpleStringTvf; + use serde::Serialize; + use std::time::Duration; + + #[derive(Serialize)] + struct HealthSettings { + observability: Observability, + } + + impl Settings for HealthSettings { + fn get_prosa_name(&self) -> String { + "health-test".to_string() + } + + fn set_prosa_name(&mut self, _name: String) {} + + fn get_observability(&self) -> &Observability { + &self.observability + } + } + + async fn wait_until(condition: impl Fn() -> bool) { + tokio::time::timeout(Duration::from_secs(1), async { + while !condition() { + tokio::task::yield_now().await; + } + }) + .await + .expect("health state should update before the timeout"); + } + + fn config_from_yaml(yaml: &str) -> Arc { + let config = config::Config::builder() + .add_source(config::File::from_str(yaml, config::FileFormat::Yaml)) + .build() + .expect("reload configuration should build"); + Arc::new(ProsaConfig::from_config(config).expect("reload configuration should be valid")) + } + + #[cfg(feature = "prometheus")] + fn gauge_value(registry: &prometheus::Registry, name: &str) -> Option { + registry + .gather() + .into_iter() + .find(|family| family.name() == name) + .and_then(|family| { + family + .get_metric() + .first() + .map(|metric| metric.get_gauge().get_value()) + }) + } + + #[tokio::test] + async fn readiness_tracks_required_processors_and_services() { + let observability = serde_yaml::from_str( + r#" +health: + required_processors: ["", required_processor] + required_services: ["", REQUIRED_SERVICE] +"#, + ) + .expect("health configuration should deserialize"); + let settings = HealthSettings { observability }; + let (bus, main) = MainProc::::create(&settings, Some(1)); + let health = bus.health.clone(); + #[cfg(feature = "prometheus")] + let registry = bus.get_prometheus_registry().clone(); + let main_task = tokio::spawn(main.run()); + + wait_until(|| health.is_live()).await; + assert!(!health.is_started()); + assert!(!health.is_ready()); + #[cfg(feature = "prometheus")] + assert_eq!(Some(0.0), gauge_value(®istry, "prosa_ready")); + + let (processor_queue, mut processor_receiver) = mpsc::channel(1); + let processor = ProcParam::new( + 1, + "required_processor".to_string(), + processor_queue, + bus.clone(), + ); + bus.add_proc_queue(ProcService::new_proc(&processor, 0)) + .await + .expect("processor should register"); + bus.add_service(vec!["REQUIRED_SERVICE".to_string()], 1, 0) + .await + .expect("service should register"); + + wait_until(|| health.is_ready()).await; + assert!(health.is_started()); + #[cfg(feature = "prometheus")] + assert_eq!(Some(1.0), gauge_value(®istry, "prosa_ready")); + + let processor_drain = + tokio::spawn(async move { while processor_receiver.recv().await.is_some() {} }); + + bus.remove_service(vec!["REQUIRED_SERVICE".to_string()], 1, 0) + .await + .expect("service should unregister"); + wait_until(|| !health.is_ready()).await; + assert!(health.is_started()); + assert!(health.is_live()); + + bus.add_service(vec!["REQUIRED_SERVICE".to_string()], 1, 0) + .await + .expect("service should register again"); + wait_until(|| health.is_ready()).await; + + bus.update_config(config_from_yaml( + r#" +observability: + health: + required_services: [MISSING_SERVICE] +"#, + )) + .await + .expect("health configuration should reload"); + wait_until(|| !health.is_ready()).await; + + bus.update_config(config_from_yaml("observability: {}")) + .await + .expect("missing health configuration should restore defaults"); + wait_until(|| health.is_ready()).await; + + bus.stop("health test complete".to_string()) + .await + .expect("main task should stop"); + assert!(!health.is_ready()); + main_task.await.expect("main task should finish"); + processor_drain.abort(); + assert!(!health.is_live()); + assert!(health.is_started()); + } +} diff --git a/prosa_book/src/ch01-02-01-observability.md b/prosa_book/src/ch01-02-01-observability.md index cff2e54..46c73d2 100644 --- a/prosa_book/src/ch01-02-01-observability.md +++ b/prosa_book/src/ch01-02-01-observability.md @@ -182,13 +182,80 @@ flowchart LR As such, you can't directly send metric to it. It's the role of Prometheus to gather metrics from your application. -To do this, you need to declare a server that exposes your ProSA metrics: +The observability HTTP server is shared by Prometheus and the ProSA health probes. Configure its +listening address at the top level and enable the Prometheus exporter to expose `/metrics`: ```yaml observability: + endpoint: "0.0.0.0:9090" level: debug - metrics: - prometheus: - endpoint: "0.0.0.0:9090" ``` -> You also need to enable the feature `prometheus` for ProSA. +> You also need to enable the `prometheus` feature for ProSA. No additional Prometheus +> configuration is required. + +Only `GET` and `HEAD` requests to `/metrics` return metrics. Other paths do not expose the +Prometheus registry. + +### Health and readiness + +ProSA tracks readiness as part of its base observability support. The `prosa_ready` gauge is +exported through every configured metrics exporter, including OTLP, stdout, and Prometheus. Its +value is `1` when ready and `0` otherwise. Neither readiness requirements nor this metric require +the HTTP feature. + +When the `observability-http` feature is enabled, the observability server also exposes three +health endpoints: + +- `/startup` succeeds permanently after ProSA first becomes ready. +- `/live` succeeds while the main ProSA task is running. +- `/ready` succeeds while ProSA is running, is not shutting down, and all configured health + requirements are available. + +Processor and service requirements are optional. When both are present, every named processor and +service is required. Empty entries are ignored: + +```yaml +observability: + endpoint: "0.0.0.0:9090" + health: + required_processors: + - api_processor + - database_processor + required_services: + - CUSTOMER_LOOKUP + - PAYMENT +``` + +A required processor is available when it has at least one queue registered with the main task. A +required service is available when it has at least one registered provider. Losing either makes +`/ready` return `503 Service Unavailable`, but does not change liveness or reset `/startup`. +Without requirements, ProSA becomes ready when the main task starts. Requirement changes are +applied during configuration reload and immediately update readiness. + +The health server can be used without Prometheus by enabling the `observability-http` feature +without the `prometheus` feature. + +For Kubernetes, configure startup separately so liveness and readiness checks do not interfere +with initialization: + +```yaml +startupProbe: + httpGet: + path: /startup + port: 9090 + periodSeconds: 2 + failureThreshold: 30 +livenessProbe: + httpGet: + path: /live + port: 9090 + periodSeconds: 10 +readinessProbe: + httpGet: + path: /ready + port: 9090 + periodSeconds: 5 +``` + +See the [Kubernetes probe documentation](https://kubernetes.io/docs/tasks/configure-pod-container/configure-liveness-readiness-startup-probes/) +for deployment-specific timing and failure thresholds. diff --git a/prosa_utils/Cargo.toml b/prosa_utils/Cargo.toml index 6cbcddd..6362e88 100644 --- a/prosa_utils/Cargo.toml +++ b/prosa_utils/Cargo.toml @@ -18,7 +18,8 @@ config-openssl = ["config", "dep:openssl"] config-openssl-vendored = ["config-openssl", "openssl/vendored"] config-observability = ["dep:log", "dep:tracing-core", "dep:tracing-subscriber", "dep:tracing-opentelemetry", "dep:opentelemetry", "dep:opentelemetry_sdk", "dep:opentelemetry-stdout", "dep:opentelemetry-otlp", "dep:opentelemetry-appender-tracing"] config-observability-gzip = ["dep:flate2", "opentelemetry-otlp/gzip-tonic", "opentelemetry-otlp/gzip-http"] -config-observability-prometheus = ["config-observability", "dep:prometheus", "dep:opentelemetry-prometheus", "dep:tokio", "dep:hyper", "dep:http-body-util", "dep:hyper-util"] +config-observability-http = ["config-observability", "dep:tokio", "dep:hyper", "dep:http-body-util", "dep:hyper-util"] +config-observability-prometheus = ["config-observability-http", "dep:prometheus", "dep:opentelemetry-prometheus"] queue = [] [package.metadata.docs.rs] diff --git a/prosa_utils/src/config/observability.rs b/prosa_utils/src/config/observability.rs index 7ff2956..f1c24ea 100644 --- a/prosa_utils/src/config/observability.rs +++ b/prosa_utils/src/config/observability.rs @@ -8,6 +8,10 @@ use opentelemetry_sdk::{ trace::{SdkTracerProvider, Tracer}, }; use serde::{Deserialize, Serialize}; +use std::sync::{ + Arc, + atomic::{AtomicBool, Ordering}, +}; use std::{collections::HashMap, fmt, time::Duration}; use tracing_subscriber::{filter, prelude::*}; use tracing_subscriber::{layer::SubscriberExt, util::TryInitError}; @@ -105,103 +109,295 @@ impl fmt::Debug for OTLPExporterCfg { } } -#[cfg(feature = "config-observability-prometheus")] -/// Configuration struct of a prometheus metric exporter -#[derive(Default, Debug, Deserialize, Serialize, Clone)] -pub struct PrometheusExporterCfg { - endpoint: Option, +/// Requirements used to determine whether ProSA is ready. +#[derive(Default, Debug, Deserialize, Serialize, Clone, PartialEq, Eq)] +pub struct HealthCheckCfg { + /// Processor names that must currently have at least one registered queue. + #[serde(default)] + required_processors: Box<[String]>, + /// Service names that must currently have at least one registered provider. + #[serde(default)] + required_services: Box<[String]>, +} + +impl HealthCheckCfg { + /// Required processor names. + pub fn required_processors(&self) -> &[String] { + &self.required_processors + } + + /// Required service names. + pub fn required_services(&self) -> &[String] { + &self.required_services + } +} + +/// Shared ProSA health state used by observability exporters. +#[derive(Debug, Default)] +pub struct HealthState { + live: AtomicBool, + started: AtomicBool, + ready: AtomicBool, +} + +impl HealthState { + #[cfg(feature = "config-observability-http")] + fn healthy() -> Self { + Self { + live: AtomicBool::new(true), + started: AtomicBool::new(true), + ready: AtomicBool::new(true), + } + } + + /// Update whether the main ProSA task is running. + pub fn set_live(&self, live: bool) { + self.live.store(live, Ordering::Release); + if !live { + self.ready.store(false, Ordering::Release); + } + } + + /// Update readiness, permanently marking startup complete on the first ready state. + pub fn set_ready(&self, ready: bool) { + self.ready.store(ready, Ordering::Release); + if ready { + self.started.store(true, Ordering::Release); + } + } + + /// Whether the main ProSA task is currently running. + pub fn is_live(&self) -> bool { + self.live.load(Ordering::Acquire) + } + + /// Whether ProSA has reached readiness at least once. + pub fn is_started(&self) -> bool { + self.started.load(Ordering::Acquire) + } + + /// Whether ProSA currently satisfies its readiness requirements. + pub fn is_ready(&self) -> bool { + self.ready.load(Ordering::Acquire) + } +} + +/// Guard marking a [`HealthState`] live for its lifetime. +#[derive(Debug)] +pub struct HealthGuard(Arc); + +impl HealthGuard { + /// Mark the health state live until the returned guard is dropped. + pub fn new(health: Arc) -> Self { + health.set_live(true); + Self(health) + } +} + +impl Drop for HealthGuard { + fn drop(&mut self) { + self.0.set_live(false); + } +} + +#[cfg(feature = "config-observability-http")] +type ObservabilityResponse = hyper::Response< + http_body_util::Either, http_body_util::Full>, +>; + +#[cfg(feature = "config-observability-http")] +fn text_response( + status: hyper::StatusCode, + body: &'static str, + head: bool, +) -> ObservabilityResponse { + let response_body = if head || body.is_empty() { + http_body_util::Either::Left(http_body_util::Empty::new()) + } else { + http_body_util::Either::Right(http_body_util::Full::new(bytes::Bytes::from_static( + body.as_bytes(), + ))) + }; + let mut response = hyper::Response::new(response_body); + *response.status_mut() = status; + response.headers_mut().insert( + hyper::header::SERVER, + hyper::header::HeaderValue::from_static(concat!("ProSA/", env!("CARGO_PKG_VERSION"))), + ); + response.headers_mut().insert( + hyper::header::CONTENT_TYPE, + hyper::header::HeaderValue::from_static("text/plain; charset=utf-8"), + ); + response } #[cfg(feature = "config-observability-prometheus")] -impl PrometheusExporterCfg { - /// Start an HTTP server to expose the metrics if needed - pub(crate) fn init_prometheus_server( - &self, - registry: &prometheus::Registry, - ) -> Result<(), ExporterBuildError> { - if let Some(endpoint) = self.endpoint.clone() { - let registry = registry.clone(); - tokio::task::spawn(async move { - match tokio::net::TcpListener::bind(endpoint).await { - Ok(listener) => { - loop { - if let Ok((stream, _)) = listener.accept().await { - let io = hyper_util::rt::TokioIo::new(stream); - let registry = registry.clone(); - tokio::task::spawn(async move { - if let Err(err) = hyper::server::conn::http1::Builder::new() - .serve_connection( - io, - #[allow(unused)] - hyper::service::service_fn(|req| { - let registry = registry.clone(); - async move { - let metric_families = registry.gather(); - let encoder = prometheus::TextEncoder::new(); - if let Ok(metric_data) = - encoder.encode_to_string(&metric_families) - { - let response = hyper::Response::builder() - .header(hyper::header::SERVER, concat!("ProSA/", env!("CARGO_PKG_VERSION"))) - .header( - hyper::header::CONTENT_TYPE, - "text/plain; version=1.0.0", - ); - - #[cfg(feature = "config-observability-gzip")] - if req.headers().get(hyper::header::ACCEPT_ENCODING).is_some_and(|a| a.to_str().is_ok_and(|v| v.contains("gzip"))) { - let mut gz_encoder = flate2::write::GzEncoder::new(Vec::with_capacity(2048), flate2::Compression::fast()); - if std::io::Write::write_all(&mut gz_encoder, metric_data.as_bytes()).is_ok() - && let Ok(compressed_data) = gz_encoder.finish() - { - return response - .header(hyper::header::CONTENT_ENCODING, "gzip") - .body(http_body_util::Full::new( - bytes::Bytes::from(compressed_data)), - ) - .map_err(|e| e.to_string()); - } - } - - response - .body(http_body_util::Full::new( - bytes::Bytes::from(metric_data), - )) - .map_err(|e| e.to_string()) - } else { - Err("Can't serialize metrics".to_string()) - } - } - }), - ) - .await - { - log::debug!(target: "prosa::observability::prometheus_server", "Error serving prometheus connection: {err:?}"); - } - }); - } - } - } - Err(e) => { - log::error!(target: "prosa::observability::prometheus_server", "Failed to bind Prometheus metrics server: {e}"); - } - } - }); +fn metrics_response( + _request: &hyper::Request, + registry: &prometheus::Registry, + head: bool, +) -> ObservabilityResponse { + if head { + let mut response = text_response(hyper::StatusCode::OK, "", true); + response.headers_mut().insert( + hyper::header::CONTENT_TYPE, + hyper::header::HeaderValue::from_static(prometheus::TEXT_FORMAT), + ); + return response; + } + + let metric_families = registry.gather(); + let encoder = prometheus::TextEncoder::new(); + let Ok(metric_data) = encoder.encode_to_string(&metric_families) else { + return text_response( + hyper::StatusCode::INTERNAL_SERVER_ERROR, + "can't serialize metrics\n", + false, + ); + }; + + #[cfg(feature = "config-observability-gzip")] + let (response_body, compressed) = { + let mut response_body = bytes::Bytes::from(metric_data); + let mut compressed = false; + if _request + .headers() + .get(hyper::header::ACCEPT_ENCODING) + .is_some_and(|encoding| encoding.to_str().is_ok_and(|value| value.contains("gzip"))) + { + let mut encoder = flate2::write::GzEncoder::new( + Vec::with_capacity(2048), + flate2::Compression::fast(), + ); + if std::io::Write::write_all(&mut encoder, &response_body).is_ok() + && let Ok(compressed_data) = encoder.finish() + { + response_body = bytes::Bytes::from(compressed_data); + compressed = true; + } } + (response_body, compressed) + }; + #[cfg(not(feature = "config-observability-gzip"))] + let (response_body, compressed) = (bytes::Bytes::from(metric_data), false); - Ok(()) + let mut response = hyper::Response::new(http_body_util::Either::Right( + http_body_util::Full::new(response_body), + )); + *response.status_mut() = hyper::StatusCode::OK; + response.headers_mut().insert( + hyper::header::SERVER, + hyper::header::HeaderValue::from_static(concat!("ProSA/", env!("CARGO_PKG_VERSION"))), + ); + response.headers_mut().insert( + hyper::header::CONTENT_TYPE, + hyper::header::HeaderValue::from_static(prometheus::TEXT_FORMAT), + ); + if compressed { + response.headers_mut().insert( + hyper::header::CONTENT_ENCODING, + hyper::header::HeaderValue::from_static("gzip"), + ); } + response +} - pub(crate) fn get_resource( - &self, - attr: Vec, - ) -> opentelemetry_sdk::resource::Resource { - opentelemetry_sdk::resource::Resource::builder() - .with_attributes(attr) - .build() +#[cfg(feature = "config-observability-http")] +async fn handle_observability_request( + request: hyper::Request, + health: Arc, + #[cfg(feature = "config-observability-prometheus")] registry: prometheus::Registry, +) -> Result { + let head = request.method() == hyper::Method::HEAD; + if request.method() == hyper::Method::GET || head { + let response = match request.uri().path() { + #[cfg(feature = "config-observability-prometheus")] + "/metrics" => metrics_response(&request, ®istry, head), + "/startup" => { + if health.is_started() { + text_response(hyper::StatusCode::OK, "ok\n", head) + } else { + text_response(hyper::StatusCode::SERVICE_UNAVAILABLE, "starting\n", head) + } + } + "/live" => { + if health.is_live() { + text_response(hyper::StatusCode::OK, "ok\n", head) + } else { + text_response(hyper::StatusCode::SERVICE_UNAVAILABLE, "not live\n", head) + } + } + "/ready" => { + if health.is_ready() { + text_response(hyper::StatusCode::OK, "ok\n", head) + } else { + text_response(hyper::StatusCode::SERVICE_UNAVAILABLE, "not ready\n", head) + } + } + _ => text_response(hyper::StatusCode::NOT_FOUND, "not found\n", head), + }; + + Ok(response) + } else { + let mut response = text_response( + hyper::StatusCode::METHOD_NOT_ALLOWED, + "method not allowed\n", + true, + ); + response.headers_mut().insert( + hyper::header::ALLOW, + hyper::header::HeaderValue::from_static("GET, HEAD"), + ); + Ok(response) } } +#[cfg(feature = "config-observability-http")] +fn init_observability_server( + endpoint: String, + health: Arc, + #[cfg(feature = "config-observability-prometheus")] registry: prometheus::Registry, +) { + tokio::task::spawn(async move { + match tokio::net::TcpListener::bind(&endpoint).await { + Ok(listener) => loop { + match listener.accept().await { + Ok((stream, _)) => { + let io = hyper_util::rt::TokioIo::new(stream); + let health = health.clone(); + #[cfg(feature = "config-observability-prometheus")] + let registry = registry.clone(); + tokio::task::spawn(async move { + if let Err(err) = hyper::server::conn::http1::Builder::new() + .serve_connection( + io, + hyper::service::service_fn(move |request| { + handle_observability_request( + request, + health.clone(), + #[cfg(feature = "config-observability-prometheus")] + registry.clone(), + ) + }), + ) + .await + { + log::debug!(target: "prosa::observability::http_server", "Error serving observability connection: {err:?}"); + } + }); + } + Err(err) => { + log::error!(target: "prosa::observability::http_server", "Failed to accept observability connection: {err}"); + } + } + }, + Err(err) => { + log::error!(target: "prosa::observability::http_server", "Failed to bind observability server on {endpoint}: {err}"); + } + } + }); +} + /// Configuration struct of an stdout exporter #[derive(Default, Debug, Deserialize, Serialize, Copy, Clone)] pub(crate) struct StdoutExporterCfg { @@ -213,8 +409,6 @@ pub(crate) struct StdoutExporterCfg { #[derive(Default, Debug, Deserialize, Serialize, Clone)] pub struct TelemetryMetrics { otlp: Option, - #[cfg(feature = "config-observability-prometheus")] - prometheus: Option, stdout: Option, } @@ -251,8 +445,7 @@ impl TelemetryMetrics { } #[cfg(feature = "config-observability-prometheus")] - if let Some(prom) = &self.prometheus { - // configure OpenTelemetry to use this registry + { let exporter = opentelemetry_prometheus::exporter() .with_registry(registry.clone()) .with_resource_selector(opentelemetry_prometheus::ResourceSelector::All) @@ -260,11 +453,12 @@ impl TelemetryMetrics { .build() .map_err(|e| ExporterBuildError::InternalFailure(e.to_string()))?; meter_provider = meter_provider - .with_resource(prom.get_resource(resource_attr)) + .with_resource( + opentelemetry_sdk::resource::Resource::builder() + .with_attributes(resource_attr) + .build(), + ) .with_reader(exporter); - - // Initialize the Prometheus server if needed - prom.init_prometheus_server(registry)?; } if self.stdout.is_some() { @@ -431,6 +625,12 @@ pub struct Observability { /// Global level for observability #[serde(default)] level: TelemetryLevel, + /// Shared HTTP endpoint for health probes and Prometheus metrics. + #[cfg(feature = "config-observability-http")] + endpoint: Option, + /// Readiness requirements. + #[serde(default)] + health: HealthCheckCfg, /// Metrics settings of a ProSA metrics: Option, /// Logs settings of a ProSA @@ -469,6 +669,9 @@ impl Observability { Observability { attributes: HashMap::new(), level, + #[cfg(feature = "config-observability-http")] + endpoint: None, + health: HealthCheckCfg::default(), metrics: Some(TelemetryMetrics::default()), logs: Some(TelemetryData::default()), traces: Some(TelemetryData::default()), @@ -545,20 +748,96 @@ impl Observability { self.level } + /// Get the configured shared HTTP endpoint. + #[cfg(feature = "config-observability-http")] + pub fn get_endpoint(&self) -> Option<&str> { + self.endpoint.as_deref() + } + + /// Get the configured readiness requirements. + pub fn get_health_check(&self) -> &HealthCheckCfg { + &self.health + } + + #[cfg(feature = "config-observability-http")] + fn start_http_server( + &self, + health: Arc, + #[cfg(feature = "config-observability-prometheus")] registry: &prometheus::Registry, + ) { + if let Some(endpoint) = self.get_endpoint() { + #[cfg(feature = "config-observability-prometheus")] + let registry = registry.clone(); + + init_observability_server( + endpoint.to_string(), + health, + #[cfg(feature = "config-observability-prometheus")] + registry, + ); + } + } + /// Meter provider builder #[cfg(feature = "config-observability-prometheus")] pub fn build_meter_provider(&self, registry: &prometheus::Registry) -> SdkMeterProvider { - if let Some(settings) = &self.metrics { - settings - .build_provider(self.get_scope_attributes(), registry) - .unwrap_or_default() + let default_settings = TelemetryMetrics::default(); + let settings = self.metrics.as_ref().unwrap_or(&default_settings); + let meter_provider = settings + .build_provider(self.get_scope_attributes(), registry) + .unwrap_or_default(); + self.start_http_server(Arc::new(HealthState::healthy()), registry); + meter_provider + } + + /// Meter provider builder + #[cfg(all( + feature = "config-observability-http", + not(feature = "config-observability-prometheus") + ))] + pub fn build_meter_provider(&self) -> SdkMeterProvider { + let meter_provider = if let Some(settings) = &self.metrics { + settings.build_provider().unwrap_or_default() } else { SdkMeterProvider::default() - } + }; + self.start_http_server(Arc::new(HealthState::healthy())); + meter_provider + } + + /// Build a meter provider and report managed ProSA health on the shared HTTP server. + #[cfg(feature = "config-observability-prometheus")] + pub fn build_meter_provider_with_health( + &self, + registry: &prometheus::Registry, + health: Arc, + ) -> SdkMeterProvider { + let default_settings = TelemetryMetrics::default(); + let settings = self.metrics.as_ref().unwrap_or(&default_settings); + let meter_provider = settings + .build_provider(self.get_scope_attributes(), registry) + .unwrap_or_default(); + self.start_http_server(health, registry); + meter_provider + } + + /// Build a meter provider and report managed ProSA health on the shared HTTP server. + #[cfg(all( + feature = "config-observability-http", + not(feature = "config-observability-prometheus") + ))] + pub fn build_meter_provider_with_health(&self, health: Arc) -> SdkMeterProvider { + let meter_provider = if let Some(settings) = &self.metrics { + settings.build_provider().unwrap_or_default() + } else { + SdkMeterProvider::default() + }; + self.start_http_server(health); + meter_provider } /// Meter provider builder - #[cfg(not(feature = "config-observability-prometheus"))] + #[cfg(not(feature = "config-observability-http"))] pub fn build_meter_provider(&self) -> SdkMeterProvider { if let Some(settings) = &self.metrics { settings.build_provider().unwrap_or_default() @@ -684,6 +963,9 @@ impl Default for Observability { Self { attributes: HashMap::new(), level: TelemetryLevel::default(), + #[cfg(feature = "config-observability-http")] + endpoint: None, + health: HealthCheckCfg::default(), metrics: Some(TelemetryMetrics::default()), logs: Some(TelemetryData { otlp: None, @@ -705,6 +987,27 @@ impl Default for Observability { mod tests { use super::*; + #[cfg(feature = "config-observability-http")] + async fn health_request( + method: hyper::Method, + path: &str, + health: Arc, + ) -> ObservabilityResponse { + let request = hyper::Request::builder() + .method(method) + .uri(path) + .body(()) + .expect("health request should be valid"); + handle_observability_request( + request, + health, + #[cfg(feature = "config-observability-prometheus")] + prometheus::Registry::new(), + ) + .await + .expect("observability handler should be infallible") + } + #[test] fn otlp_http_authorization_preserves_literal_percent_triplets() { let config = OTLPExporterCfg { @@ -737,4 +1040,192 @@ mod tests { assert!(!debug.contains("password")); assert!(!debug.contains("secret")); } + + #[test] + fn health_configuration_uses_plain_requirements() { + let config: Observability = serde_yaml::from_str( + r#" +health: + required_processors: [api, "", worker, api] + required_services: [PAYMENT, ""] +"#, + ) + .expect("observability configuration should deserialize"); + + assert_eq!( + &[ + "api".to_string(), + String::new(), + "worker".to_string(), + "api".to_string() + ], + config.get_health_check().required_processors() + ); + assert_eq!( + &["PAYMENT".to_string(), String::new()], + config.get_health_check().required_services() + ); + } + + #[test] + fn health_guard_tracks_liveness() { + let health = Arc::new(HealthState::default()); + { + let _guard = HealthGuard::new(health.clone()); + assert!(health.is_live()); + health.set_ready(true); + } + assert!(!health.is_live()); + assert!(!health.is_ready()); + assert!(health.is_started()); + } + + #[cfg(feature = "config-observability-http")] + #[test] + fn observability_endpoint_is_used() { + let config: Observability = serde_yaml::from_str( + r#" +endpoint: 127.0.0.1:8080 +"#, + ) + .expect("observability configuration should deserialize"); + + assert_eq!(Some("127.0.0.1:8080"), config.get_endpoint()); + } + + #[cfg(feature = "config-observability-http")] + #[tokio::test] + async fn health_routes_follow_lifecycle_state() { + use http_body_util::BodyExt as _; + + let health = Arc::new(HealthState::default()); + assert_eq!( + hyper::StatusCode::SERVICE_UNAVAILABLE, + health_request(hyper::Method::GET, "/startup", health.clone()) + .await + .status() + ); + assert_eq!( + hyper::StatusCode::SERVICE_UNAVAILABLE, + health_request(hyper::Method::GET, "/live", health.clone()) + .await + .status() + ); + + health.set_live(true); + assert_eq!( + hyper::StatusCode::OK, + health_request(hyper::Method::GET, "/live", health.clone()) + .await + .status() + ); + health.set_ready(true); + health.set_ready(false); + assert_eq!( + hyper::StatusCode::OK, + health_request(hyper::Method::GET, "/startup", health.clone()) + .await + .status() + ); + assert_eq!( + hyper::StatusCode::SERVICE_UNAVAILABLE, + health_request(hyper::Method::GET, "/ready", health.clone()) + .await + .status() + ); + + let response = health_request(hyper::Method::HEAD, "/startup", health.clone()).await; + assert_eq!(hyper::StatusCode::OK, response.status()); + assert!( + response + .into_body() + .collect() + .await + .expect("response body should be readable") + .to_bytes() + .is_empty() + ); + assert_eq!( + hyper::StatusCode::METHOD_NOT_ALLOWED, + health_request(hyper::Method::POST, "/ready", health.clone()) + .await + .status() + ); + assert_eq!( + hyper::StatusCode::NOT_FOUND, + health_request(hyper::Method::GET, "/", health) + .await + .status() + ); + #[cfg(not(feature = "config-observability-prometheus"))] + assert_eq!( + hyper::StatusCode::NOT_FOUND, + health_request( + hyper::Method::GET, + "/metrics", + Arc::new(HealthState::default()) + ) + .await + .status() + ); + } + + #[cfg(feature = "config-observability-prometheus")] + #[tokio::test] + async fn metrics_are_only_exposed_on_metrics_route_when_enabled() { + let health = Arc::new(HealthState::healthy()); + let registry = prometheus::Registry::new(); + let request = hyper::Request::builder() + .method(hyper::Method::GET) + .uri("/metrics") + .body(()) + .expect("metrics request should be valid"); + let response = handle_observability_request(request, health.clone(), registry) + .await + .expect("observability handler should be infallible"); + assert_eq!(hyper::StatusCode::OK, response.status()); + assert_eq!( + Some(prometheus::TEXT_FORMAT), + response + .headers() + .get(hyper::header::CONTENT_TYPE) + .and_then(|value| value.to_str().ok()) + ); + + let request = hyper::Request::builder() + .method(hyper::Method::HEAD) + .uri("/metrics") + .body(()) + .expect("metrics request should be valid"); + let response = + handle_observability_request(request, health.clone(), prometheus::Registry::new()) + .await + .expect("observability handler should be infallible"); + assert_eq!(hyper::StatusCode::OK, response.status()); + assert_eq!( + Some(prometheus::TEXT_FORMAT), + response + .headers() + .get(hyper::header::CONTENT_TYPE) + .and_then(|value| value.to_str().ok()) + ); + + let request = hyper::Request::builder() + .method(hyper::Method::POST) + .uri("/metrics") + .body(()) + .expect("metrics request should be valid"); + let response = + handle_observability_request(request, health.clone(), prometheus::Registry::new()) + .await + .expect("observability handler should be infallible"); + assert_eq!(hyper::StatusCode::METHOD_NOT_ALLOWED, response.status()); + + assert_eq!( + hyper::StatusCode::NOT_FOUND, + health_request(hyper::Method::GET, "/anything", health) + .await + .status() + ); + } } From ff3461076206ae16404787bcf2668f803eb193ce Mon Sep 17 00:00:00 2001 From: Jeremy HERGAULT Date: Mon, 21 Sep 2026 08:51:15 +0200 Subject: [PATCH 2/5] feat: bump prosa-utils version to 0.5.2 Signed-off-by: Jeremy HERGAULT --- Cargo.toml | 2 +- prosa/Cargo.toml | 2 +- prosa_utils/Cargo.toml | 2 +- 3 files changed, 3 insertions(+), 3 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index 07d3c66..654c257 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -23,7 +23,7 @@ include = [ ] [workspace.dependencies] -prosa-utils = { version = "0.5.1", path = "prosa_utils" } +prosa-utils = { version = "0.5.2", path = "prosa_utils" } prosa-macros = { version = "0.5.1", path = "prosa_macros" } thiserror = "2" simple-mermaid = "0.2" diff --git a/prosa/Cargo.toml b/prosa/Cargo.toml index e9bb869..6ddee63 100644 --- a/prosa/Cargo.toml +++ b/prosa/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "prosa" -version = "0.5.1" +version = "0.5.2" authors.workspace = true description = "ProSA core" homepage.workspace = true diff --git a/prosa_utils/Cargo.toml b/prosa_utils/Cargo.toml index 6362e88..7252830 100644 --- a/prosa_utils/Cargo.toml +++ b/prosa_utils/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "prosa-utils" -version = "0.5.1" +version = "0.5.2" authors.workspace = true description = "ProSA utils" homepage.workspace = true From 735ceeb3dea927e0dbe59d768951451edf6e4948 Mon Sep 17 00:00:00 2001 From: Jeremy HERGAULT Date: Tue, 29 Sep 2026 15:13:17 +0200 Subject: [PATCH 3/5] feat: declare required service at startup for metrics Signed-off-by: Jeremy HERGAULT --- prosa/examples/my_prosa_settings.yml | 5 +++++ prosa/src/core/main.rs | 13 ++++++++++++- prosa/src/core/service.rs | 9 +++++++++ 3 files changed, 26 insertions(+), 1 deletion(-) diff --git a/prosa/examples/my_prosa_settings.yml b/prosa/examples/my_prosa_settings.yml index 9664755..a4ef4c4 100644 --- a/prosa/examples/my_prosa_settings.yml +++ b/prosa/examples/my_prosa_settings.yml @@ -3,6 +3,11 @@ name: my-prosa observability: endpoint: 0.0.0.0:9100 level: INFO + health: + required_processors: + - stub_proc + required_services: + - STUB_TEST traces: stdout: level: debug diff --git a/prosa/src/core/main.rs b/prosa/src/core/main.rs index b562093..67a9923 100644 --- a/prosa/src/core/main.rs +++ b/prosa/src/core/main.rs @@ -475,12 +475,23 @@ where let stop = main.stop.clone(); let health = main.health.clone(); let health_check = main.health_check.clone(); + + // Declare required services so they are visible (without processor) until a processor serve them + let mut services = ServiceTable::default(); + for service_name in health_check + .required_services() + .iter() + .filter(|name| !name.is_empty()) + { + services.declare_service(service_name); + } + ( main, MainProc { name, processors, - services: Arc::new(ServiceTable::default()), + services: Arc::new(services), config: None, internal_rx_queue, meter, diff --git a/prosa/src/core/service.rs b/prosa/src/core/service.rs index d0b0c81..cb0a079 100644 --- a/prosa/src/core/service.rs +++ b/prosa/src/core/service.rs @@ -87,6 +87,15 @@ where }) } + /// Method to declare a service without any processor that serve it + /// + /// Can be call only by the main task to declare expected services before processors are started + pub fn declare_service(&mut self, name: &str) { + self.table + .entry(name.into()) + .or_insert_with(|| (Vec::new(), atomic::AtomicU64::new(0))); + } + /// Method to add a service to the table /// /// Can be call only by the main task to modify the service table From 47bd60c90878d39a5c288bdde17e901fcbebbccc Mon Sep 17 00:00:00 2001 From: Jeremy HERGAULT Date: Thu, 1 Oct 2026 15:33:25 +0200 Subject: [PATCH 4/5] feat: simplify according to code review Signed-off-by: Jeremy HERGAULT --- Cargo.toml | 4 +- prosa/Cargo.toml | 3 +- prosa/src/core/main.rs | 46 +- prosa_book/src/ch01-02-01-observability.md | 18 +- prosa_utils/Cargo.toml | 15 +- prosa_utils/src/config/observability.rs | 587 +++++---------------- 6 files changed, 169 insertions(+), 504 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index 654c257..d49fe11 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -36,8 +36,8 @@ toml = "0.9" # Config Observability log = "0.4" -opentelemetry = { version = "0.32", features = ["metrics", "trace", "logs"] } -opentelemetry_sdk = { version = "0.32", features = ["metrics", "trace", "logs", "rt-tokio"] } +opentelemetry = { version = "0.33", features = ["metrics", "trace", "logs"] } +opentelemetry_sdk = { version = "0.33", features = ["metrics", "trace", "logs", "rt-tokio"] } prometheus = { version = "0.14", default-features = false } # Config SSL diff --git a/prosa/Cargo.toml b/prosa/Cargo.toml index 6ddee63..16c17ac 100644 --- a/prosa/Cargo.toml +++ b/prosa/Cargo.toml @@ -15,8 +15,7 @@ system-metrics = ["dep:memory-stats"] http-proxy = ["dep:async-http-proxy"] openssl = ["dep:openssl", "dep:tokio-openssl", "prosa-utils/config-openssl"] openssl-vendored = ["openssl", "openssl/vendored", "prosa-utils/config-openssl-vendored"] -observability-http = ["prosa-utils/config-observability-http"] -prometheus = ["observability-http", "dep:prometheus", "prosa-utils/config-observability-prometheus"] +prometheus = ["dep:prometheus", "prosa-utils/config-observability-prometheus"] queue = ["prosa-utils/queue"] [[example]] diff --git a/prosa/src/core/main.rs b/prosa/src/core/main.rs index 67a9923..9551b9d 100644 --- a/prosa/src/core/main.rs +++ b/prosa/src/core/main.rs @@ -20,7 +20,7 @@ use crate::otel::metrics::{Meter, MeterProvider as _}; use crate::otel::trace::TracerProvider as _; use crate::otel::{InstrumentationScope, KeyValue}; use crate::tracing::{debug, info, warn}; -use prosa_utils::config::observability::{HealthCheckCfg, HealthGuard, HealthState}; +use prosa_utils::config::observability::{HealthCheckCfg, HealthState}; use prosa_utils::hash::{BuildIntHasher, IntHashMap, IntHashSet}; use std::sync::{ Arc, @@ -92,13 +92,11 @@ where let health = Arc::new(HealthState::default()); #[cfg(feature = "prometheus")] let prometheus_registry = prometheus::Registry::new(); - #[cfg(feature = "prometheus")] - let meter_provider = - observability.build_meter_provider_with_health(&prometheus_registry, health.clone()); - #[cfg(all(feature = "observability-http", not(feature = "prometheus")))] - let meter_provider = observability.build_meter_provider_with_health(health.clone()); - #[cfg(not(feature = "observability-http"))] + let meter_provider = observability.build_meter_provider(&prometheus_registry); + #[cfg(feature = "prometheus")] + observability.start_http_server(health.clone(), &prometheus_registry); + #[cfg(not(feature = "prometheus"))] let meter_provider = observability.build_meter_provider(); Main { @@ -219,7 +217,6 @@ where /// Method to stop all processors pub async fn stop(&self, reason: String) -> Result<(), SendError>> { self.stop.store(true, Ordering::Relaxed); - self.health.set_ready(false); Ok(self .internal_tx_queue .send(InternalMainMsg::Shutdown(reason)) @@ -308,8 +305,7 @@ where .iter() .all(|name| name.is_empty() || self.services.exist_proc_service(name)); - self.health - .set_ready(!self.stop.load(Ordering::Relaxed) && processors_ready && services_ready); + self.health.set_ready(processors_ready && services_ready); } async fn remove_proc(&mut self, proc_id: u32) -> Option> { @@ -317,7 +313,6 @@ where let mut new_services = (*self.services).clone(); new_services.remove_proc_services(proc_id); self.services = Arc::new(new_services); - self.update_health(); Some(proc) } else { None @@ -333,7 +328,6 @@ where proc_queue.get_queue_id(), ); self.services = Arc::new(new_services); - self.update_health(); Some(proc_queue) } else { None @@ -432,6 +426,7 @@ where /// Method to shutdown all processors (return `true` if all processor are off, `false` otherwise) async fn stop(&mut self) -> bool { self.stop.store(true, Ordering::Relaxed); + self.health.set_ready(false); let mut is_stopped = true; for proc in self.processors.values() { for proc_service in proc.values() { @@ -515,7 +510,6 @@ where } async fn run(mut self) { - let _health_guard = HealthGuard::new(self.health.clone()); self.update_health(); // Monitor readiness @@ -657,7 +651,6 @@ where } prosa_main_record_proc!(); - self.update_health(); }, InternalMainMsg::DeleteProc(proc_id, proc_err) => { if self.remove_proc(proc_id).await.is_some() { @@ -694,7 +687,6 @@ where } } self.services = Arc::new(new_services); - self.update_health(); let _ = service_update.send(self.services.clone()); self.notify_srv_proc().await; } @@ -706,7 +698,6 @@ where new_services.add_service(name, proc_queue.clone()); } self.services = Arc::new(new_services); - self.update_health(); let _ = service_update.send(self.services.clone()); self.notify_srv_proc().await; } @@ -717,7 +708,6 @@ where new_services.remove_service_proc(&name, proc_id); } self.services = Arc::new(new_services); - self.update_health(); let _ = service_update.send(self.services.clone()); self.notify_srv_proc().await; }, @@ -727,7 +717,6 @@ where new_services.remove_service(&name, proc_id, queue_id); } self.services = Arc::new(new_services); - self.update_health(); let _ = service_update.send(self.services.clone()); self.notify_srv_proc().await; }, @@ -747,11 +736,8 @@ where None } }; - if let Some(health_check) = health_check - && self.health_check != health_check - { + if let Some(health_check) = health_check { self.health_check = health_check; - self.update_health(); } for error in self.notify_config_proc_queue(config).await { @@ -767,7 +753,6 @@ where }, InternalMainMsg::Shutdown(reason) => { warn!("ProSA is stopping: {}", reason); - self.health.set_ready(false); self.stop().await; // The shutdown mecanism will be implemented later @@ -777,13 +762,14 @@ where }, _ = signal::ctrl_c() => { warn!("ProSA is stopping"); - self.health.set_ready(false); self.stop().await; // The shutdown mecanism will be implemented later return; }, } + + self.update_health(); } } } @@ -863,7 +849,8 @@ health: let registry = bus.get_prometheus_registry().clone(); let main_task = tokio::spawn(main.run()); - wait_until(|| health.is_live()).await; + #[cfg(feature = "prometheus")] + wait_until(|| gauge_value(®istry, "prosa_ready").is_some()).await; assert!(!health.is_started()); assert!(!health.is_ready()); #[cfg(feature = "prometheus")] @@ -876,6 +863,8 @@ health: processor_queue, bus.clone(), ); + let processor_drain = + tokio::spawn(async move { while processor_receiver.recv().await.is_some() {} }); bus.add_proc_queue(ProcService::new_proc(&processor, 0)) .await .expect("processor should register"); @@ -888,15 +877,11 @@ health: #[cfg(feature = "prometheus")] assert_eq!(Some(1.0), gauge_value(®istry, "prosa_ready")); - let processor_drain = - tokio::spawn(async move { while processor_receiver.recv().await.is_some() {} }); - bus.remove_service(vec!["REQUIRED_SERVICE".to_string()], 1, 0) .await .expect("service should unregister"); wait_until(|| !health.is_ready()).await; assert!(health.is_started()); - assert!(health.is_live()); bus.add_service(vec!["REQUIRED_SERVICE".to_string()], 1, 0) .await @@ -922,10 +907,9 @@ observability: bus.stop("health test complete".to_string()) .await .expect("main task should stop"); - assert!(!health.is_ready()); main_task.await.expect("main task should finish"); processor_drain.abort(); - assert!(!health.is_live()); + assert!(!health.is_ready()); assert!(health.is_started()); } } diff --git a/prosa_book/src/ch01-02-01-observability.md b/prosa_book/src/ch01-02-01-observability.md index 46c73d2..1dde6e2 100644 --- a/prosa_book/src/ch01-02-01-observability.md +++ b/prosa_book/src/ch01-02-01-observability.md @@ -193,8 +193,7 @@ observability: > You also need to enable the `prometheus` feature for ProSA. No additional Prometheus > configuration is required. -Only `GET` and `HEAD` requests to `/metrics` return metrics. Other paths do not expose the -Prometheus registry. +Only the `/metrics` path returns metrics. Other paths do not expose the Prometheus registry. ### Health and readiness @@ -203,13 +202,13 @@ exported through every configured metrics exporter, including OTLP, stdout, and value is `1` when ready and `0` otherwise. Neither readiness requirements nor this metric require the HTTP feature. -When the `observability-http` feature is enabled, the observability server also exposes three -health endpoints: +When the `prometheus` feature is enabled, the observability server also exposes three health +endpoints: - `/startup` succeeds permanently after ProSA first becomes ready. -- `/live` succeeds while the main ProSA task is running. -- `/ready` succeeds while ProSA is running, is not shutting down, and all configured health - requirements are available. +- `/live` always succeeds; answering proves the process is alive. +- `/ready` succeeds while ProSA is not shutting down and all configured health requirements are + available. Processor and service requirements are optional. When both are present, every named processor and service is required. Empty entries are ignored: @@ -228,13 +227,10 @@ observability: A required processor is available when it has at least one queue registered with the main task. A required service is available when it has at least one registered provider. Losing either makes -`/ready` return `503 Service Unavailable`, but does not change liveness or reset `/startup`. +`/ready` return `503 Service Unavailable`, but does not reset `/startup`. Without requirements, ProSA becomes ready when the main task starts. Requirement changes are applied during configuration reload and immediately update readiness. -The health server can be used without Prometheus by enabling the `observability-http` feature -without the `prometheus` feature. - For Kubernetes, configure startup separately so liveness and readiness checks do not interfere with initialization: diff --git a/prosa_utils/Cargo.toml b/prosa_utils/Cargo.toml index 7252830..0a43bd0 100644 --- a/prosa_utils/Cargo.toml +++ b/prosa_utils/Cargo.toml @@ -18,8 +18,7 @@ config-openssl = ["config", "dep:openssl"] config-openssl-vendored = ["config-openssl", "openssl/vendored"] config-observability = ["dep:log", "dep:tracing-core", "dep:tracing-subscriber", "dep:tracing-opentelemetry", "dep:opentelemetry", "dep:opentelemetry_sdk", "dep:opentelemetry-stdout", "dep:opentelemetry-otlp", "dep:opentelemetry-appender-tracing"] config-observability-gzip = ["dep:flate2", "opentelemetry-otlp/gzip-tonic", "opentelemetry-otlp/gzip-http"] -config-observability-http = ["config-observability", "dep:tokio", "dep:hyper", "dep:http-body-util", "dep:hyper-util"] -config-observability-prometheus = ["config-observability-http", "dep:prometheus", "dep:opentelemetry-prometheus"] +config-observability-prometheus = ["config-observability", "dep:prometheus", "dep:opentelemetry-prometheus", "dep:tokio", "dep:hyper", "dep:http-body-util", "dep:hyper-util"] queue = [] [package.metadata.docs.rs] @@ -43,7 +42,7 @@ hex = "0.4" glob = { version = "0.3", optional = true } serde = { workspace = true, optional = true } serde_yaml = { version = "0.9", optional = true } -base64 = { version = "0.22", optional = true } +base64 = { version = "0.23", optional = true } percent-encoding = { version = "2", optional = true } uuid = { version = "1", optional = true, features = ["v1", "v4"] } @@ -54,14 +53,14 @@ openssl = { workspace = true, optional = true } log = { workspace = true, optional = true } tracing-core = { version = "0.1", optional = true } tracing-subscriber = { version = ">=0.3.20, < 0.4", features = ["std", "env-filter"], optional = true } -tracing-opentelemetry = { version = "0.33", optional = true } +tracing-opentelemetry = { version = "0.34", optional = true } opentelemetry = { workspace = true, optional = true } opentelemetry_sdk = { workspace = true, optional = true } -opentelemetry-stdout = { version = "0.32", optional = true, features = ["metrics", "trace", "logs"]} -opentelemetry-otlp = { version = "0.32", optional = true, features = ["metrics", "trace", "logs", "reqwest-rustls", "grpc-tonic", "http-json"]} +opentelemetry-stdout = { version = "0.33", optional = true, features = ["metrics", "trace", "logs"]} +opentelemetry-otlp = { version = "0.33", optional = true, features = ["metrics", "trace", "logs", "reqwest-rustls", "grpc-tonic", "http-json"]} prometheus = { workspace = true, optional = true } -opentelemetry-prometheus = { version = "0.32", optional = true } -opentelemetry-appender-tracing = { version = "0.32", optional = true } +opentelemetry-prometheus = { version = "0.33", optional = true } +opentelemetry-appender-tracing = { version = "0.33", optional = true } # Web Observability tokio = { workspace = true, optional = true } diff --git a/prosa_utils/src/config/observability.rs b/prosa_utils/src/config/observability.rs index f1c24ea..7bd0689 100644 --- a/prosa_utils/src/config/observability.rs +++ b/prosa_utils/src/config/observability.rs @@ -8,10 +8,7 @@ use opentelemetry_sdk::{ trace::{SdkTracerProvider, Tracer}, }; use serde::{Deserialize, Serialize}; -use std::sync::{ - Arc, - atomic::{AtomicBool, Ordering}, -}; +use std::sync::atomic::{AtomicU8, Ordering}; use std::{collections::HashMap, fmt, time::Duration}; use tracing_subscriber::{filter, prelude::*}; use tracing_subscriber::{layer::SubscriberExt, util::TryInitError}; @@ -132,270 +129,109 @@ impl HealthCheckCfg { } } +const HEALTH_STARTED: u8 = 0b01; +const HEALTH_READY: u8 = 0b10; + /// Shared ProSA health state used by observability exporters. #[derive(Debug, Default)] -pub struct HealthState { - live: AtomicBool, - started: AtomicBool, - ready: AtomicBool, -} +pub struct HealthState(AtomicU8); impl HealthState { - #[cfg(feature = "config-observability-http")] - fn healthy() -> Self { - Self { - live: AtomicBool::new(true), - started: AtomicBool::new(true), - ready: AtomicBool::new(true), - } - } - - /// Update whether the main ProSA task is running. - pub fn set_live(&self, live: bool) { - self.live.store(live, Ordering::Release); - if !live { - self.ready.store(false, Ordering::Release); - } - } - /// Update readiness, permanently marking startup complete on the first ready state. pub fn set_ready(&self, ready: bool) { - self.ready.store(ready, Ordering::Release); if ready { - self.started.store(true, Ordering::Release); + self.0 + .store(HEALTH_STARTED | HEALTH_READY, Ordering::Relaxed); + } else { + self.0.fetch_and(!HEALTH_READY, Ordering::Relaxed); } } - /// Whether the main ProSA task is currently running. - pub fn is_live(&self) -> bool { - self.live.load(Ordering::Acquire) - } - /// Whether ProSA has reached readiness at least once. pub fn is_started(&self) -> bool { - self.started.load(Ordering::Acquire) + self.0.load(Ordering::Relaxed) & HEALTH_STARTED != 0 } /// Whether ProSA currently satisfies its readiness requirements. pub fn is_ready(&self) -> bool { - self.ready.load(Ordering::Acquire) + self.0.load(Ordering::Relaxed) & HEALTH_READY != 0 } } -/// Guard marking a [`HealthState`] live for its lifetime. -#[derive(Debug)] -pub struct HealthGuard(Arc); - -impl HealthGuard { - /// Mark the health state live until the returned guard is dropped. - pub fn new(health: Arc) -> Self { - health.set_live(true); - Self(health) - } -} +#[cfg(feature = "config-observability-prometheus")] +type ObservabilityResponse = hyper::Response>; -impl Drop for HealthGuard { - fn drop(&mut self) { - self.0.set_live(false); - } +#[cfg(feature = "config-observability-prometheus")] +fn response_builder() -> hyper::http::response::Builder { + hyper::Response::builder().header( + hyper::header::SERVER, + concat!("ProSA/", env!("CARGO_PKG_VERSION")), + ) } -#[cfg(feature = "config-observability-http")] -type ObservabilityResponse = hyper::Response< - http_body_util::Either, http_body_util::Full>, ->; - -#[cfg(feature = "config-observability-http")] -fn text_response( - status: hyper::StatusCode, - body: &'static str, - head: bool, -) -> ObservabilityResponse { - let response_body = if head || body.is_empty() { - http_body_util::Either::Left(http_body_util::Empty::new()) +#[cfg(feature = "config-observability-prometheus")] +fn probe_response( + ok: bool, + error: &'static str, +) -> Result { + if ok { + response_builder().body("ok\n".into()) } else { - http_body_util::Either::Right(http_body_util::Full::new(bytes::Bytes::from_static( - body.as_bytes(), - ))) - }; - let mut response = hyper::Response::new(response_body); - *response.status_mut() = status; - response.headers_mut().insert( - hyper::header::SERVER, - hyper::header::HeaderValue::from_static(concat!("ProSA/", env!("CARGO_PKG_VERSION"))), - ); - response.headers_mut().insert( - hyper::header::CONTENT_TYPE, - hyper::header::HeaderValue::from_static("text/plain; charset=utf-8"), - ); - response + response_builder() + .status(hyper::StatusCode::SERVICE_UNAVAILABLE) + .body(error.into()) + } } #[cfg(feature = "config-observability-prometheus")] fn metrics_response( _request: &hyper::Request, registry: &prometheus::Registry, - head: bool, -) -> ObservabilityResponse { - if head { - let mut response = text_response(hyper::StatusCode::OK, "", true); - response.headers_mut().insert( - hyper::header::CONTENT_TYPE, - hyper::header::HeaderValue::from_static(prometheus::TEXT_FORMAT), - ); - return response; - } - - let metric_families = registry.gather(); - let encoder = prometheus::TextEncoder::new(); - let Ok(metric_data) = encoder.encode_to_string(&metric_families) else { - return text_response( - hyper::StatusCode::INTERNAL_SERVER_ERROR, - "can't serialize metrics\n", - false, - ); +) -> Result { + let Ok(metric_data) = prometheus::TextEncoder::new().encode_to_string(®istry.gather()) + else { + return response_builder() + .status(hyper::StatusCode::INTERNAL_SERVER_ERROR) + .body("can't serialize metrics\n".into()); }; + let response = response_builder().header(hyper::header::CONTENT_TYPE, prometheus::TEXT_FORMAT); #[cfg(feature = "config-observability-gzip")] - let (response_body, compressed) = { - let mut response_body = bytes::Bytes::from(metric_data); - let mut compressed = false; - if _request - .headers() - .get(hyper::header::ACCEPT_ENCODING) - .is_some_and(|encoding| encoding.to_str().is_ok_and(|value| value.contains("gzip"))) + if _request + .headers() + .get(hyper::header::ACCEPT_ENCODING) + .is_some_and(|a| a.to_str().is_ok_and(|v| v.contains("gzip"))) + { + let mut gz_encoder = + flate2::write::GzEncoder::new(Vec::with_capacity(2048), flate2::Compression::fast()); + if std::io::Write::write_all(&mut gz_encoder, metric_data.as_bytes()).is_ok() + && let Ok(compressed_data) = gz_encoder.finish() { - let mut encoder = flate2::write::GzEncoder::new( - Vec::with_capacity(2048), - flate2::Compression::fast(), - ); - if std::io::Write::write_all(&mut encoder, &response_body).is_ok() - && let Ok(compressed_data) = encoder.finish() - { - response_body = bytes::Bytes::from(compressed_data); - compressed = true; - } + return response + .header(hyper::header::CONTENT_ENCODING, "gzip") + .body(compressed_data.into()); } - (response_body, compressed) - }; - #[cfg(not(feature = "config-observability-gzip"))] - let (response_body, compressed) = (bytes::Bytes::from(metric_data), false); - - let mut response = hyper::Response::new(http_body_util::Either::Right( - http_body_util::Full::new(response_body), - )); - *response.status_mut() = hyper::StatusCode::OK; - response.headers_mut().insert( - hyper::header::SERVER, - hyper::header::HeaderValue::from_static(concat!("ProSA/", env!("CARGO_PKG_VERSION"))), - ); - response.headers_mut().insert( - hyper::header::CONTENT_TYPE, - hyper::header::HeaderValue::from_static(prometheus::TEXT_FORMAT), - ); - if compressed { - response.headers_mut().insert( - hyper::header::CONTENT_ENCODING, - hyper::header::HeaderValue::from_static("gzip"), - ); } - response -} -#[cfg(feature = "config-observability-http")] -async fn handle_observability_request( - request: hyper::Request, - health: Arc, - #[cfg(feature = "config-observability-prometheus")] registry: prometheus::Registry, -) -> Result { - let head = request.method() == hyper::Method::HEAD; - if request.method() == hyper::Method::GET || head { - let response = match request.uri().path() { - #[cfg(feature = "config-observability-prometheus")] - "/metrics" => metrics_response(&request, ®istry, head), - "/startup" => { - if health.is_started() { - text_response(hyper::StatusCode::OK, "ok\n", head) - } else { - text_response(hyper::StatusCode::SERVICE_UNAVAILABLE, "starting\n", head) - } - } - "/live" => { - if health.is_live() { - text_response(hyper::StatusCode::OK, "ok\n", head) - } else { - text_response(hyper::StatusCode::SERVICE_UNAVAILABLE, "not live\n", head) - } - } - "/ready" => { - if health.is_ready() { - text_response(hyper::StatusCode::OK, "ok\n", head) - } else { - text_response(hyper::StatusCode::SERVICE_UNAVAILABLE, "not ready\n", head) - } - } - _ => text_response(hyper::StatusCode::NOT_FOUND, "not found\n", head), - }; - - Ok(response) - } else { - let mut response = text_response( - hyper::StatusCode::METHOD_NOT_ALLOWED, - "method not allowed\n", - true, - ); - response.headers_mut().insert( - hyper::header::ALLOW, - hyper::header::HeaderValue::from_static("GET, HEAD"), - ); - Ok(response) - } + response.body(metric_data.into()) } -#[cfg(feature = "config-observability-http")] -fn init_observability_server( - endpoint: String, - health: Arc, - #[cfg(feature = "config-observability-prometheus")] registry: prometheus::Registry, -) { - tokio::task::spawn(async move { - match tokio::net::TcpListener::bind(&endpoint).await { - Ok(listener) => loop { - match listener.accept().await { - Ok((stream, _)) => { - let io = hyper_util::rt::TokioIo::new(stream); - let health = health.clone(); - #[cfg(feature = "config-observability-prometheus")] - let registry = registry.clone(); - tokio::task::spawn(async move { - if let Err(err) = hyper::server::conn::http1::Builder::new() - .serve_connection( - io, - hyper::service::service_fn(move |request| { - handle_observability_request( - request, - health.clone(), - #[cfg(feature = "config-observability-prometheus")] - registry.clone(), - ) - }), - ) - .await - { - log::debug!(target: "prosa::observability::http_server", "Error serving observability connection: {err:?}"); - } - }); - } - Err(err) => { - log::error!(target: "prosa::observability::http_server", "Failed to accept observability connection: {err}"); - } - } - }, - Err(err) => { - log::error!(target: "prosa::observability::http_server", "Failed to bind observability server on {endpoint}: {err}"); - } - } - }); +#[cfg(feature = "config-observability-prometheus")] +fn handle_observability_request( + request: &hyper::Request, + health: &HealthState, + registry: &prometheus::Registry, +) -> Result { + match request.uri().path() { + "/metrics" => metrics_response(request, registry), + "/startup" => probe_response(health.is_started(), "starting\n"), + // Answering at all proves the process is alive + "/live" => probe_response(true, ""), + "/ready" => probe_response(health.is_ready(), "not ready\n"), + _ => response_builder() + .status(hyper::StatusCode::NOT_FOUND) + .body("not found\n".into()), + } } /// Configuration struct of an stdout exporter @@ -626,7 +462,7 @@ pub struct Observability { #[serde(default)] level: TelemetryLevel, /// Shared HTTP endpoint for health probes and Prometheus metrics. - #[cfg(feature = "config-observability-http")] + #[cfg(feature = "config-observability-prometheus")] endpoint: Option, /// Readiness requirements. #[serde(default)] @@ -669,7 +505,7 @@ impl Observability { Observability { attributes: HashMap::new(), level, - #[cfg(feature = "config-observability-http")] + #[cfg(feature = "config-observability-prometheus")] endpoint: None, health: HealthCheckCfg::default(), metrics: Some(TelemetryMetrics::default()), @@ -749,7 +585,7 @@ impl Observability { } /// Get the configured shared HTTP endpoint. - #[cfg(feature = "config-observability-http")] + #[cfg(feature = "config-observability-prometheus")] pub fn get_endpoint(&self) -> Option<&str> { self.endpoint.as_deref() } @@ -759,85 +595,62 @@ impl Observability { &self.health } - #[cfg(feature = "config-observability-http")] - fn start_http_server( + /// Start the observability HTTP server (metrics and health probes) if an endpoint is configured + #[cfg(feature = "config-observability-prometheus")] + pub fn start_http_server( &self, - health: Arc, - #[cfg(feature = "config-observability-prometheus")] registry: &prometheus::Registry, + health: std::sync::Arc, + registry: &prometheus::Registry, ) { - if let Some(endpoint) = self.get_endpoint() { - #[cfg(feature = "config-observability-prometheus")] + if let Some(endpoint) = self.endpoint.clone() { let registry = registry.clone(); - - init_observability_server( - endpoint.to_string(), - health, - #[cfg(feature = "config-observability-prometheus")] - registry, - ); + tokio::task::spawn(async move { + let listener = match tokio::net::TcpListener::bind(&endpoint).await { + Ok(listener) => listener, + Err(e) => { + log::error!(target: "prosa::observability::http_server", "Failed to bind observability server on {endpoint}: {e}"); + return; + } + }; + loop { + if let Ok((stream, _)) = listener.accept().await { + let io = hyper_util::rt::TokioIo::new(stream); + let health = health.clone(); + let registry = registry.clone(); + tokio::task::spawn(async move { + if let Err(err) = hyper::server::conn::http1::Builder::new() + .serve_connection( + io, + hyper::service::service_fn(|req| { + std::future::ready(handle_observability_request( + &req, &health, ®istry, + )) + }), + ) + .await + { + log::debug!(target: "prosa::observability::http_server", "Error serving observability connection: {err:?}"); + } + }); + } + } + }); } } /// Meter provider builder #[cfg(feature = "config-observability-prometheus")] pub fn build_meter_provider(&self, registry: &prometheus::Registry) -> SdkMeterProvider { - let default_settings = TelemetryMetrics::default(); - let settings = self.metrics.as_ref().unwrap_or(&default_settings); - let meter_provider = settings + // Prometheus exporter is always attached, even without metrics settings + self.metrics + .clone() + .unwrap_or_default() .build_provider(self.get_scope_attributes(), registry) - .unwrap_or_default(); - self.start_http_server(Arc::new(HealthState::healthy()), registry); - meter_provider + .unwrap_or_default() } /// Meter provider builder - #[cfg(all( - feature = "config-observability-http", - not(feature = "config-observability-prometheus") - ))] - pub fn build_meter_provider(&self) -> SdkMeterProvider { - let meter_provider = if let Some(settings) = &self.metrics { - settings.build_provider().unwrap_or_default() - } else { - SdkMeterProvider::default() - }; - self.start_http_server(Arc::new(HealthState::healthy())); - meter_provider - } - - /// Build a meter provider and report managed ProSA health on the shared HTTP server. - #[cfg(feature = "config-observability-prometheus")] - pub fn build_meter_provider_with_health( - &self, - registry: &prometheus::Registry, - health: Arc, - ) -> SdkMeterProvider { - let default_settings = TelemetryMetrics::default(); - let settings = self.metrics.as_ref().unwrap_or(&default_settings); - let meter_provider = settings - .build_provider(self.get_scope_attributes(), registry) - .unwrap_or_default(); - self.start_http_server(health, registry); - meter_provider - } - - /// Build a meter provider and report managed ProSA health on the shared HTTP server. - #[cfg(all( - feature = "config-observability-http", - not(feature = "config-observability-prometheus") - ))] - pub fn build_meter_provider_with_health(&self, health: Arc) -> SdkMeterProvider { - let meter_provider = if let Some(settings) = &self.metrics { - settings.build_provider().unwrap_or_default() - } else { - SdkMeterProvider::default() - }; - self.start_http_server(health); - meter_provider - } - - /// Meter provider builder - #[cfg(not(feature = "config-observability-http"))] + #[cfg(not(feature = "config-observability-prometheus"))] pub fn build_meter_provider(&self) -> SdkMeterProvider { if let Some(settings) = &self.metrics { settings.build_provider().unwrap_or_default() @@ -963,7 +776,7 @@ impl Default for Observability { Self { attributes: HashMap::new(), level: TelemetryLevel::default(), - #[cfg(feature = "config-observability-http")] + #[cfg(feature = "config-observability-prometheus")] endpoint: None, health: HealthCheckCfg::default(), metrics: Some(TelemetryMetrics::default()), @@ -987,25 +800,13 @@ impl Default for Observability { mod tests { use super::*; - #[cfg(feature = "config-observability-http")] - async fn health_request( - method: hyper::Method, - path: &str, - health: Arc, - ) -> ObservabilityResponse { - let request = hyper::Request::builder() - .method(method) - .uri(path) + #[cfg(feature = "config-observability-prometheus")] + fn request(path: &str, health: &HealthState) -> ObservabilityResponse { + let request = hyper::Request::get(path) .body(()) - .expect("health request should be valid"); - handle_observability_request( - request, - health, - #[cfg(feature = "config-observability-prometheus")] - prometheus::Registry::new(), - ) - .await - .expect("observability handler should be infallible") + .expect("request should be valid"); + handle_observability_request(&request, health, &prometheus::Registry::new()) + .expect("response should be valid") } #[test] @@ -1068,19 +869,16 @@ health: } #[test] - fn health_guard_tracks_liveness() { - let health = Arc::new(HealthState::default()); - { - let _guard = HealthGuard::new(health.clone()); - assert!(health.is_live()); - health.set_ready(true); - } - assert!(!health.is_live()); + fn health_state_keeps_startup_after_ready() { + let health = HealthState::default(); + assert!(!health.is_started()); + health.set_ready(true); + health.set_ready(false); assert!(!health.is_ready()); assert!(health.is_started()); } - #[cfg(feature = "config-observability-http")] + #[cfg(feature = "config-observability-prometheus")] #[test] fn observability_endpoint_is_used() { let config: Observability = serde_yaml::from_str( @@ -1093,96 +891,21 @@ endpoint: 127.0.0.1:8080 assert_eq!(Some("127.0.0.1:8080"), config.get_endpoint()); } - #[cfg(feature = "config-observability-http")] - #[tokio::test] - async fn health_routes_follow_lifecycle_state() { - use http_body_util::BodyExt as _; - - let health = Arc::new(HealthState::default()); - assert_eq!( - hyper::StatusCode::SERVICE_UNAVAILABLE, - health_request(hyper::Method::GET, "/startup", health.clone()) - .await - .status() - ); - assert_eq!( - hyper::StatusCode::SERVICE_UNAVAILABLE, - health_request(hyper::Method::GET, "/live", health.clone()) - .await - .status() - ); + #[cfg(feature = "config-observability-prometheus")] + #[test] + fn observability_routes() { + let health = HealthState::default(); + let status = |path| request(path, &health).status(); + assert_eq!(hyper::StatusCode::SERVICE_UNAVAILABLE, status("/startup")); + assert_eq!(hyper::StatusCode::OK, status("/live")); + assert_eq!(hyper::StatusCode::SERVICE_UNAVAILABLE, status("/ready")); + assert_eq!(hyper::StatusCode::NOT_FOUND, status("/")); - health.set_live(true); - assert_eq!( - hyper::StatusCode::OK, - health_request(hyper::Method::GET, "/live", health.clone()) - .await - .status() - ); health.set_ready(true); - health.set_ready(false); - assert_eq!( - hyper::StatusCode::OK, - health_request(hyper::Method::GET, "/startup", health.clone()) - .await - .status() - ); - assert_eq!( - hyper::StatusCode::SERVICE_UNAVAILABLE, - health_request(hyper::Method::GET, "/ready", health.clone()) - .await - .status() - ); - - let response = health_request(hyper::Method::HEAD, "/startup", health.clone()).await; - assert_eq!(hyper::StatusCode::OK, response.status()); - assert!( - response - .into_body() - .collect() - .await - .expect("response body should be readable") - .to_bytes() - .is_empty() - ); - assert_eq!( - hyper::StatusCode::METHOD_NOT_ALLOWED, - health_request(hyper::Method::POST, "/ready", health.clone()) - .await - .status() - ); - assert_eq!( - hyper::StatusCode::NOT_FOUND, - health_request(hyper::Method::GET, "/", health) - .await - .status() - ); - #[cfg(not(feature = "config-observability-prometheus"))] - assert_eq!( - hyper::StatusCode::NOT_FOUND, - health_request( - hyper::Method::GET, - "/metrics", - Arc::new(HealthState::default()) - ) - .await - .status() - ); - } + assert_eq!(hyper::StatusCode::OK, status("/startup")); + assert_eq!(hyper::StatusCode::OK, status("/ready")); - #[cfg(feature = "config-observability-prometheus")] - #[tokio::test] - async fn metrics_are_only_exposed_on_metrics_route_when_enabled() { - let health = Arc::new(HealthState::healthy()); - let registry = prometheus::Registry::new(); - let request = hyper::Request::builder() - .method(hyper::Method::GET) - .uri("/metrics") - .body(()) - .expect("metrics request should be valid"); - let response = handle_observability_request(request, health.clone(), registry) - .await - .expect("observability handler should be infallible"); + let response = request("/metrics", &health); assert_eq!(hyper::StatusCode::OK, response.status()); assert_eq!( Some(prometheus::TEXT_FORMAT), @@ -1191,41 +914,5 @@ endpoint: 127.0.0.1:8080 .get(hyper::header::CONTENT_TYPE) .and_then(|value| value.to_str().ok()) ); - - let request = hyper::Request::builder() - .method(hyper::Method::HEAD) - .uri("/metrics") - .body(()) - .expect("metrics request should be valid"); - let response = - handle_observability_request(request, health.clone(), prometheus::Registry::new()) - .await - .expect("observability handler should be infallible"); - assert_eq!(hyper::StatusCode::OK, response.status()); - assert_eq!( - Some(prometheus::TEXT_FORMAT), - response - .headers() - .get(hyper::header::CONTENT_TYPE) - .and_then(|value| value.to_str().ok()) - ); - - let request = hyper::Request::builder() - .method(hyper::Method::POST) - .uri("/metrics") - .body(()) - .expect("metrics request should be valid"); - let response = - handle_observability_request(request, health.clone(), prometheus::Registry::new()) - .await - .expect("observability handler should be infallible"); - assert_eq!(hyper::StatusCode::METHOD_NOT_ALLOWED, response.status()); - - assert_eq!( - hyper::StatusCode::NOT_FOUND, - health_request(hyper::Method::GET, "/anything", health) - .await - .status() - ); } } From ddb6c816b8bd2d831af0a1a5202e248e723f7ecd Mon Sep 17 00:00:00 2001 From: Jeremy HERGAULT Date: Thu, 1 Oct 2026 16:47:06 +0200 Subject: [PATCH 5/5] feat: fix CI and improve code Signed-off-by: Jeremy HERGAULT --- prosa/Cargo.toml | 2 +- prosa/src/core/main.rs | 20 +++++++++++--------- prosa/src/core/settings.rs | 2 +- prosa_utils/Cargo.toml | 4 ++-- prosa_utils/src/config/observability.rs | 4 ++-- prosa_utils/src/config/ssl.rs | 4 ++-- prosa_utils/src/queue/lockfree.rs | 4 ++-- 7 files changed, 21 insertions(+), 19 deletions(-) diff --git a/prosa/Cargo.toml b/prosa/Cargo.toml index 16c17ac..76e600c 100644 --- a/prosa/Cargo.toml +++ b/prosa/Cargo.toml @@ -63,7 +63,7 @@ config = { version = "0.15", default-features = false, features = ["toml", "yaml notify = "8" glob = { version = "0.3" } toml.workspace = true -serde_yaml = "0.9" +yaml_serde = "0.10" opentelemetry.workspace = true opentelemetry_sdk.workspace = true diff --git a/prosa/src/core/main.rs b/prosa/src/core/main.rs index 9551b9d..0158de1 100644 --- a/prosa/src/core/main.rs +++ b/prosa/src/core/main.rs @@ -90,14 +90,16 @@ where ) -> Main { let observability = settings.get_observability(); let health = Arc::new(HealthState::default()); - #[cfg(feature = "prometheus")] - let prometheus_registry = prometheus::Registry::new(); - #[cfg(feature = "prometheus")] - let meter_provider = observability.build_meter_provider(&prometheus_registry); - #[cfg(feature = "prometheus")] - observability.start_http_server(health.clone(), &prometheus_registry); - #[cfg(not(feature = "prometheus"))] - let meter_provider = observability.build_meter_provider(); + cfg_select! { + feature = "prometheus" => { + let prometheus_registry = prometheus::Registry::new(); + let meter_provider = observability.build_meter_provider(&prometheus_registry); + observability.start_http_server(health.clone(), &prometheus_registry); + }, + not(feature = "prometheus") => { + let meter_provider = observability.build_meter_provider(); + }, + } Main { internal_tx_queue, @@ -834,7 +836,7 @@ mod tests { #[tokio::test] async fn readiness_tracks_required_processors_and_services() { - let observability = serde_yaml::from_str( + let observability = yaml_serde::from_str( r#" health: required_processors: ["", required_processor] diff --git a/prosa/src/core/settings.rs b/prosa/src/core/settings.rs index 3d1312d..fca9db3 100644 --- a/prosa/src/core/settings.rs +++ b/prosa/src/core/settings.rs @@ -122,7 +122,7 @@ pub trait Settings: Serialize { writeln!( f, "{}", - serde_yaml::to_string(&self) + yaml_serde::to_string(&self) .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))? ) } diff --git a/prosa_utils/Cargo.toml b/prosa_utils/Cargo.toml index 0a43bd0..65b0524 100644 --- a/prosa_utils/Cargo.toml +++ b/prosa_utils/Cargo.toml @@ -13,7 +13,7 @@ include.workspace = true default = ["msg", "config", "config-openssl", "config-observability"] msg = [] msg-serde = ["msg", "dep:serde", "chrono/serde", "bytes/serde"] -config = ["dep:glob","dep:serde","dep:serde_yaml", "dep:base64", "dep:percent-encoding", "dep:uuid"] +config = ["dep:glob","dep:serde","dep:yaml_serde", "dep:base64", "dep:percent-encoding", "dep:uuid"] config-openssl = ["config", "dep:openssl"] config-openssl-vendored = ["config-openssl", "openssl/vendored"] config-observability = ["dep:log", "dep:tracing-core", "dep:tracing-subscriber", "dep:tracing-opentelemetry", "dep:opentelemetry", "dep:opentelemetry_sdk", "dep:opentelemetry-stdout", "dep:opentelemetry-otlp", "dep:opentelemetry-appender-tracing"] @@ -41,7 +41,7 @@ hex = "0.4" # Config glob = { version = "0.3", optional = true } serde = { workspace = true, optional = true } -serde_yaml = { version = "0.9", optional = true } +yaml_serde = { version = "0.10", optional = true } base64 = { version = "0.23", optional = true } percent-encoding = { version = "2", optional = true } uuid = { version = "1", optional = true, features = ["v1", "v4"] } diff --git a/prosa_utils/src/config/observability.rs b/prosa_utils/src/config/observability.rs index 7bd0689..27ebaef 100644 --- a/prosa_utils/src/config/observability.rs +++ b/prosa_utils/src/config/observability.rs @@ -844,7 +844,7 @@ mod tests { #[test] fn health_configuration_uses_plain_requirements() { - let config: Observability = serde_yaml::from_str( + let config: Observability = yaml_serde::from_str( r#" health: required_processors: [api, "", worker, api] @@ -881,7 +881,7 @@ health: #[cfg(feature = "config-observability-prometheus")] #[test] fn observability_endpoint_is_used() { - let config: Observability = serde_yaml::from_str( + let config: Observability = yaml_serde::from_str( r#" endpoint: 127.0.0.1:8080 "#, diff --git a/prosa_utils/src/config/ssl.rs b/prosa_utils/src/config/ssl.rs index f872058..e871c80 100644 --- a/prosa_utils/src/config/ssl.rs +++ b/prosa_utils/src/config/ssl.rs @@ -395,7 +395,7 @@ tL4ndQavEi51mI38AjEAi/V3bNTIZargCyzuFJ0nN6T5U6VR5CmD1/iQMVtCnwr1 }; assert!(format!("{inline_store_le_x1_x2}").contains("ISRG Root X")); - let config_store_le_x1_x2: Store = serde_yaml::from_str( + let config_store_le_x1_x2: Store = yaml_serde::from_str( "certs: - | -----BEGIN CERTIFICATE----- @@ -448,7 +448,7 @@ tL4ndQavEi51mI38AjEAi/V3bNTIZargCyzuFJ0nN6T5U6VR5CmD1/iQMVtCnwr1 .expect("SSL certificate configuration should be read"); assert!(format!("{config_store_le_x1_x2}").contains("ISRG Root X")); - let config_store_file: Store = serde_yaml::from_str("path: \"/opt\"") + let config_store_file: Store = yaml_serde::from_str("path: \"/opt\"") .expect("Certificate configuration path should be read"); assert_eq!( Store::File { diff --git a/prosa_utils/src/queue/lockfree.rs b/prosa_utils/src/queue/lockfree.rs index 8a20ff3..449e8f1 100644 --- a/prosa_utils/src/queue/lockfree.rs +++ b/prosa_utils/src/queue/lockfree.rs @@ -191,7 +191,7 @@ macro_rules! impl_consume_queue { if !self.is_empty() { Ok(self .head - .fetch_update( + .try_update( std::sync::atomic::Ordering::Relaxed, std::sync::atomic::Ordering::Relaxed, |head| Some((head + 1) % self.max_capacity()), @@ -210,7 +210,7 @@ macro_rules! impl_consume_queue { pub unsafe fn consume(&self) -> Result<$p, QueueError> { if !self.is_empty() { self.head - .fetch_update( + .try_update( std::sync::atomic::Ordering::Relaxed, std::sync::atomic::Ordering::Relaxed, |head| Some((head + 1) % self.max_capacity()),