use std::time::Duration;
use camel_api::metrics::{AllocatorStat, MetricsCollector};
use prometheus::{CounterVec, HistogramVec, Opts, Registry};
mod families;
use families::StaticFamilies;
fn normalize_prom_name(name: &str) -> String {
let prefixed = if name.starts_with("camel_") {
name.to_string()
} else {
format!("camel_{name}")
};
let sanitized: String = prefixed
.chars()
.map(|c| {
if c.is_ascii_alphanumeric() || c == '_' || c == ':' {
c
} else {
'_'
}
})
.collect();
sanitized
}
fn sort_label_pairs<'a>(labels: &'a [(&'a str, &'a str)]) -> Vec<(&'a str, &'a str)> {
let mut pairs = labels.to_vec();
pairs.sort_by(|a, b| a.0.cmp(b.0));
pairs
}
fn counter_value_ok(v: f64) -> bool {
!v.is_nan() && v >= 0.0 && v.fract() == 0.0
}
struct DynCounter {
cv: CounterVec,
keys: Vec<String>,
}
struct DynHistogram {
hv: HistogramVec,
keys: Vec<String>,
}
fn keys_match(frozen: &[String], incoming: &[(&str, &str)]) -> bool {
if frozen.len() != incoming.len() {
return false;
}
frozen.iter().zip(incoming.iter()).all(|(f, (k, _))| f == k)
}
pub struct PrometheusMetrics {
registry: Registry,
families: StaticFamilies,
dyn_counters: dashmap::DashMap<String, Option<DynCounter>>,
dyn_histograms: dashmap::DashMap<String, Option<DynHistogram>>,
warned: dashmap::DashSet<String>,
max_dynamic_collectors: usize,
}
impl PrometheusMetrics {
pub fn new() -> Self {
let registry = Registry::new();
let families = families::register_static_families(®istry);
Self {
registry,
families,
dyn_counters: dashmap::DashMap::new(),
dyn_histograms: dashmap::DashMap::new(),
warned: dashmap::DashSet::new(),
max_dynamic_collectors: 1024,
}
}
pub fn max_dynamic_collectors(&self) -> usize {
self.max_dynamic_collectors
}
pub fn with_max_dynamic_collectors(mut self, n: usize) -> Self {
self.max_dynamic_collectors = n;
self
}
pub fn registry(&self) -> &Registry {
&self.registry
}
pub fn gather(&self) -> String {
families::render(&self.registry)
}
}
impl Default for PrometheusMetrics {
fn default() -> Self {
Self::new()
}
}
impl MetricsCollector for PrometheusMetrics {
fn record_exchange_duration(&self, route_id: &str, duration: Duration) {
let duration_secs = duration.as_secs_f64();
self.families
.exchange_duration_seconds
.with_label_values(&[route_id])
.observe(duration_secs);
}
fn increment_errors(&self, route_id: &str, error_type: &str) {
self.families
.errors_total
.with_label_values(&[route_id, error_type])
.inc();
}
fn increment_exchanges(&self, route_id: &str) {
self.families
.exchanges_total
.with_label_values(&[route_id])
.inc();
}
fn set_queue_depth(&self, queue: &str, depth: usize) {
self.families
.queue_depth
.with_label_values(&[queue])
.set(depth as f64);
}
fn set_pinned_client_cache_size(&self, component: &str, entries: u64) {
self.families
.pinned_client_cache_size
.with_label_values(&[component])
.set(entries as f64);
}
fn increment_pinned_client_cache_hit(&self, component: &str) {
self.families
.pinned_client_cache_hits_total
.with_label_values(&[component])
.inc();
}
fn increment_pinned_client_cache_miss(&self, component: &str) {
self.families
.pinned_client_cache_misses_total
.with_label_values(&[component])
.inc();
}
fn set_allocator_memory(&self, stat: AllocatorStat, bytes: u64) {
self.families
.allocator_memory_bytes
.with_label_values(&[stat.as_str()])
.set(bytes as f64);
}
fn record_circuit_breaker_change(&self, route_id: &str, _from: &str, to: &str) {
let state_value = |state: &str| -> f64 {
match state.to_lowercase().as_str() {
"closed" => 0.0,
"open" => 1.0,
"half_open" | "halfopen" => 2.0,
_ => -1.0, }
};
self.families
.circuit_breaker_state
.with_label_values(&[route_id])
.set(state_value(to));
}
fn increment_retry_attempt(&self, scheme: &str, operation: &str) {
self.families
.retry_attempts_total
.with_label_values(&[operation, scheme])
.inc();
}
fn increment_circuit_breaker_rejection(&self, route: &str) {
self.families
.circuit_breaker_rejections_total
.with_label_values(&[route])
.inc();
}
fn set_route_state(&self, route: &str, state: &str) {
self.families.route_state.set(route, state);
}
fn clear_route_state(&self, route: &str) {
self.families.route_state.remove(route);
}
fn record_build_info(&self, version: &str, git_sha: &str) {
self.families
.build_info
.with_label_values(&[git_sha, version])
.set(1);
}
fn record_uptime(&self, seconds: f64) {
self.families.uptime_seconds.set(seconds);
}
fn record_component_operation(&self, component: &str, operation: &str, outcome: &str) {
self.families
.component_operations_total
.with_label_values(&[component, operation, outcome])
.inc();
}
fn record_counter(&self, name: &str, value: f64, labels: &[(&str, &str)]) {
if !counter_value_ok(value) {
if self.warned.insert(name.to_string()) {
tracing::warn!(
name,
value,
"dynamic counter value rejected (NaN/negative/non-integer); \
further rejections for this name will be silent"
);
}
return;
}
let normalized = normalize_prom_name(name);
if normalized != name && self.warned.insert(format!("sanitize:{name}")) {
tracing::warn!(name, %normalized, "metric name sanitized for prometheus");
}
let sorted = sort_label_pairs(labels);
let values: Vec<&str> = sorted.iter().map(|(_, v)| *v).collect();
if self.dyn_counters.len() >= self.max_dynamic_collectors
&& !self.dyn_counters.contains_key(&normalized)
{
if self.warned.insert(format!("cap:{}", normalized)) {
tracing::warn!(
name,
cap = self.max_dynamic_collectors,
"dynamic counter cap exceeded; observation dropped"
);
}
return;
}
use dashmap::mapref::entry::Entry;
match self.dyn_counters.entry(normalized.clone()) {
Entry::Occupied(o) => match o.get() {
Some(dc) => {
if keys_match(&dc.keys, &sorted) {
dc.cv.with_label_values(&values).inc_by(value);
} else if self.warned.insert(name.to_string()) {
tracing::warn!(
name,
"dynamic counter label arity/key drift; observation dropped \
(further drift for this name will be silent)"
);
}
}
None => { }
},
Entry::Vacant(v) => {
let keys: Vec<String> = sorted.iter().map(|(k, _)| (*k).to_string()).collect();
let key_refs: Vec<&str> = keys.iter().map(|s| s.as_str()).collect();
let cv = match CounterVec::new(Opts::new(&normalized, "Dynamic counter"), &key_refs)
{
Ok(cv) => cv,
Err(_) => {
v.insert(None);
if self.warned.insert(name.to_string()) {
tracing::warn!(name, "dynamic counter creation failed; tombstoned");
}
return;
}
};
match self.registry.register(Box::new(cv.clone())) {
Ok(()) => {
cv.with_label_values(&values).inc_by(value);
v.insert(Some(DynCounter { cv, keys }));
}
Err(_) => {
v.insert(None);
if self.warned.insert(name.to_string()) {
tracing::warn!(
name,
"dynamic counter registration failed (possible name collision); \
tombstoned"
);
}
}
}
}
}
}
fn record_histogram(&self, name: &str, value: f64, labels: &[(&str, &str)]) {
if value.is_nan() {
if self.warned.insert(name.to_string()) {
tracing::warn!(
name,
"dynamic histogram value rejected (NaN); \
further NaN for this name will be silent"
);
}
return;
}
let normalized = normalize_prom_name(name);
if normalized != name && self.warned.insert(format!("sanitize:{name}")) {
tracing::warn!(name, %normalized, "metric name sanitized for prometheus");
}
let sorted = sort_label_pairs(labels);
let values: Vec<&str> = sorted.iter().map(|(_, v)| *v).collect();
if self.dyn_histograms.len() >= self.max_dynamic_collectors
&& !self.dyn_histograms.contains_key(&normalized)
{
if self.warned.insert(format!("cap:{}", normalized)) {
tracing::warn!(
name,
cap = self.max_dynamic_collectors,
"dynamic histogram cap exceeded; observation dropped"
);
}
return;
}
use dashmap::mapref::entry::Entry;
match self.dyn_histograms.entry(normalized.clone()) {
Entry::Occupied(o) => match o.get() {
Some(dh) => {
if keys_match(&dh.keys, &sorted) {
dh.hv.with_label_values(&values).observe(value);
} else if self.warned.insert(name.to_string()) {
tracing::warn!(
name,
"dynamic histogram label arity/key drift; observation dropped"
);
}
}
None => { }
},
Entry::Vacant(v) => {
let keys: Vec<String> = sorted.iter().map(|(k, _)| (*k).to_string()).collect();
let key_refs: Vec<&str> = keys.iter().map(|s| s.as_str()).collect();
let hv = match HistogramVec::new(
prometheus::HistogramOpts {
common_opts: Opts::new(&normalized, "Dynamic histogram"),
buckets: vec![
0.001, 0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1.0, 2.5, 5.0, 10.0,
],
},
&key_refs,
) {
Ok(hv) => hv,
Err(_) => {
v.insert(None);
if self.warned.insert(name.to_string()) {
tracing::warn!(name, "dynamic histogram creation failed; tombstoned");
}
return;
}
};
match self.registry.register(Box::new(hv.clone())) {
Ok(()) => {
hv.with_label_values(&values).observe(value);
v.insert(Some(DynHistogram { hv, keys }));
}
Err(_) => {
v.insert(None);
if self.warned.insert(name.to_string()) {
tracing::warn!(
name,
"dynamic histogram registration failed; tombstoned"
);
}
}
}
}
}
}
}
#[cfg(test)]
#[path = "tests.rs"]
mod tests;