Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 3 additions & 3 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -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
Expand Down
2 changes: 1 addition & 1 deletion prosa/Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[package]
name = "prosa"
version = "0.5.1"
version = "0.5.2"
authors.workspace = true
description = "ProSA core"
homepage.workspace = true
Expand Down
9 changes: 6 additions & 3 deletions prosa/examples/my_prosa_settings.yml
Original file line number Diff line number Diff line change
@@ -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
Expand Down
261 changes: 232 additions & 29 deletions prosa/src/core/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -57,6 +58,8 @@ where
scope_attributes: Vec<KeyValue>,
#[cfg(feature = "prometheus")]
prometheus_registry: prometheus::Registry,
health: Arc<HealthState>,
health_check: HealthCheckCfg,
meter_provider: opentelemetry_sdk::metrics::SdkMeterProvider,
tracer_provider: opentelemetry_sdk::trace::SdkTracerProvider,
stop: Arc<AtomicBool>,
Expand Down Expand Up @@ -85,36 +88,28 @@ where
internal_tx_queue: mpsc::Sender<InternalMainMsg<M>>,
settings: &S,
) -> Main<M> {
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(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 = 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)),
}
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)),
Comment on lines +95 to +112

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Would it be worth it to use the new cfg_select macro ?

}
}

Expand Down Expand Up @@ -274,6 +269,8 @@ where
internal_rx_queue: mpsc::Receiver<InternalMainMsg<M>>,
meter: Meter,
stop: Arc<AtomicBool>,
health: Arc<HealthState>,
health_check: HealthCheckCfg,
}

impl<M> ProcBusParam for MainProc<M>
Expand All @@ -293,6 +290,24 @@ impl<M> MainProc<M>
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<ProcQueueMap<M>> {
if let Some(proc) = self.processors.remove(&proc_id) {
let mut new_services = (*self.services).clone();
Expand Down Expand Up @@ -411,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() {
Expand Down Expand Up @@ -452,16 +468,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,
},
)
}
Expand All @@ -479,6 +510,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
Expand Down Expand Up @@ -679,6 +722,24 @@ where
},
InternalMainMsg::Config(config) => {
info!("Reloading ProSA configuration");

let health_check = match config
.config()
.get::<HealthCheckCfg>("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 {
Expand Down Expand Up @@ -707,6 +768,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<ProsaConfig> {
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<f64> {
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::<SimpleStringTvf>::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(&registry, "prosa_ready").is_some()).await;
assert!(!health.is_started());
assert!(!health.is_ready());
#[cfg(feature = "prometheus")]
assert_eq!(Some(0.0), gauge_value(&registry, "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(&registry, "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());
}
}
Loading
Loading