diff --git a/Cargo.toml b/Cargo.toml index 07d3c66..d49fe11 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" @@ -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 fc29d02..76e600c 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 @@ -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/examples/my_prosa_settings.yml b/prosa/examples/my_prosa_settings.yml index ff920bf..a4ef4c4 100644 --- a/prosa/examples/my_prosa_settings.yml +++ b/prosa/examples/my_prosa_settings.yml @@ -1,10 +1,13 @@ name: my-prosa observability: + endpoint: 0.0.0.0:9100 level: INFO - metrics: - prometheus: - endpoint: 0.0.0.0:9100 + 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 c9d3908..0158de1 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, 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 { - #[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 observability = settings.get_observability(); + let health = Arc::new(HealthState::default()); + 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(); + }, } - #[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)), - } + 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)), } } @@ -274,6 +271,8 @@ where internal_rx_queue: mpsc::Receiver>, meter: Meter, stop: Arc, + health: Arc, + health_check: HealthCheckCfg, } impl ProcBusParam for MainProc @@ -293,6 +292,24 @@ 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(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(); @@ -411,6 +428,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() { @@ -452,16 +470,31 @@ 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(); + + // 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, stop, + health, + health_check, }, ) } @@ -479,6 +512,18 @@ where } async fn run(mut self) { + 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 @@ -679,6 +724,24 @@ where }, 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; + } + for error in self.notify_config_proc_queue(config).await { if let BusError::ProcComm(proc_id, queue_id, _) = error { if queue_id > 0 { @@ -707,6 +770,148 @@ where return; }, } + + self.update_health(); } } } + +#[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 = yaml_serde::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()); + + #[cfg(feature = "prometheus")] + wait_until(|| gauge_value(®istry, "prosa_ready").is_some()).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(), + ); + 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"); + 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")); + + 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()); + + 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"); + main_task.await.expect("main task should finish"); + processor_drain.abort(); + assert!(!health.is_ready()); + assert!(health.is_started()); + } +} 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 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_book/src/ch01-02-01-observability.md b/prosa_book/src/ch01-02-01-observability.md index cff2e54..1dde6e2 100644 --- a/prosa_book/src/ch01-02-01-observability.md +++ b/prosa_book/src/ch01-02-01-observability.md @@ -182,13 +182,76 @@ 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 the `/metrics` path returns 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 `prometheus` feature is enabled, the observability server also exposes three health +endpoints: + +- `/startup` succeeds permanently after ProSA first becomes ready. +- `/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: + +```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 reset `/startup`. +Without requirements, ProSA becomes ready when the main task starts. Requirement changes are +applied during configuration reload and immediately update readiness. + +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..65b0524 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 @@ -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,8 +41,8 @@ hex = "0.4" # Config glob = { version = "0.3", optional = true } serde = { workspace = true, optional = true } -serde_yaml = { version = "0.9", optional = true } -base64 = { version = "0.22", 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"] } @@ -53,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 7ff2956..27ebaef 100644 --- a/prosa_utils/src/config/observability.rs +++ b/prosa_utils/src/config/observability.rs @@ -8,6 +8,7 @@ use opentelemetry_sdk::{ trace::{SdkTracerProvider, Tracer}, }; use serde::{Deserialize, Serialize}; +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}; @@ -105,100 +106,131 @@ impl fmt::Debug for OTLPExporterCfg { } } +/// 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 + } +} + +const HEALTH_STARTED: u8 = 0b01; +const HEALTH_READY: u8 = 0b10; + +/// Shared ProSA health state used by observability exporters. +#[derive(Debug, Default)] +pub struct HealthState(AtomicU8); + +impl HealthState { + /// Update readiness, permanently marking startup complete on the first ready state. + pub fn set_ready(&self, ready: bool) { + if ready { + self.0 + .store(HEALTH_STARTED | HEALTH_READY, Ordering::Relaxed); + } else { + self.0.fetch_and(!HEALTH_READY, Ordering::Relaxed); + } + } + + /// Whether ProSA has reached readiness at least once. + pub fn is_started(&self) -> bool { + self.0.load(Ordering::Relaxed) & HEALTH_STARTED != 0 + } + + /// Whether ProSA currently satisfies its readiness requirements. + pub fn is_ready(&self) -> bool { + self.0.load(Ordering::Relaxed) & HEALTH_READY != 0 + } +} + #[cfg(feature = "config-observability-prometheus")] -/// Configuration struct of a prometheus metric exporter -#[derive(Default, Debug, Deserialize, Serialize, Clone)] -pub struct PrometheusExporterCfg { - endpoint: Option, +type ObservabilityResponse = hyper::Response>; + +#[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-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 probe_response( + ok: bool, + error: &'static str, +) -> Result { + if ok { + response_builder().body("ok\n".into()) + } else { + response_builder() + .status(hyper::StatusCode::SERVICE_UNAVAILABLE) + .body(error.into()) + } +} - Ok(()) +#[cfg(feature = "config-observability-prometheus")] +fn metrics_response( + _request: &hyper::Request, + registry: &prometheus::Registry, +) -> 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")] + 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() + { + return response + .header(hyper::header::CONTENT_ENCODING, "gzip") + .body(compressed_data.into()); + } } - pub(crate) fn get_resource( - &self, - attr: Vec, - ) -> opentelemetry_sdk::resource::Resource { - opentelemetry_sdk::resource::Resource::builder() - .with_attributes(attr) - .build() + response.body(metric_data.into()) +} + +#[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()), } } @@ -213,8 +245,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 +281,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 +289,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 +461,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-prometheus")] + endpoint: Option, + /// Readiness requirements. + #[serde(default)] + health: HealthCheckCfg, /// Metrics settings of a ProSA metrics: Option, /// Logs settings of a ProSA @@ -469,6 +505,9 @@ impl Observability { Observability { attributes: HashMap::new(), level, + #[cfg(feature = "config-observability-prometheus")] + endpoint: None, + health: HealthCheckCfg::default(), metrics: Some(TelemetryMetrics::default()), logs: Some(TelemetryData::default()), traces: Some(TelemetryData::default()), @@ -545,16 +584,69 @@ impl Observability { self.level } + /// Get the configured shared HTTP endpoint. + #[cfg(feature = "config-observability-prometheus")] + 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 + } + + /// 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: std::sync::Arc, + registry: &prometheus::Registry, + ) { + if let Some(endpoint) = self.endpoint.clone() { + let registry = registry.clone(); + 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 { - if let Some(settings) = &self.metrics { - settings - .build_provider(self.get_scope_attributes(), registry) - .unwrap_or_default() - } else { - SdkMeterProvider::default() - } + // 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() } /// Meter provider builder @@ -684,6 +776,9 @@ impl Default for Observability { Self { attributes: HashMap::new(), level: TelemetryLevel::default(), + #[cfg(feature = "config-observability-prometheus")] + endpoint: None, + health: HealthCheckCfg::default(), metrics: Some(TelemetryMetrics::default()), logs: Some(TelemetryData { otlp: None, @@ -705,6 +800,15 @@ impl Default for Observability { mod tests { use super::*; + #[cfg(feature = "config-observability-prometheus")] + fn request(path: &str, health: &HealthState) -> ObservabilityResponse { + let request = hyper::Request::get(path) + .body(()) + .expect("request should be valid"); + handle_observability_request(&request, health, &prometheus::Registry::new()) + .expect("response should be valid") + } + #[test] fn otlp_http_authorization_preserves_literal_percent_triplets() { let config = OTLPExporterCfg { @@ -737,4 +841,78 @@ mod tests { assert!(!debug.contains("password")); assert!(!debug.contains("secret")); } + + #[test] + fn health_configuration_uses_plain_requirements() { + let config: Observability = yaml_serde::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_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-prometheus")] + #[test] + fn observability_endpoint_is_used() { + let config: Observability = yaml_serde::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-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_ready(true); + assert_eq!(hyper::StatusCode::OK, status("/startup")); + assert_eq!(hyper::StatusCode::OK, status("/ready")); + + let response = request("/metrics", &health); + 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()) + ); + } } 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()),