use axum::http::{Method, StatusCode};
use loonfs::metrics::{
CounterHandle, DefaultMetricsRecorder, HistogramHandle, MetricEntry, MetricValue,
MetricsRecorder, MetricsSnapshot, LATENCY_SECONDS_BOUNDARIES,
};
use loonfs::RuntimeCacheStats;
use std::collections::HashMap;
use std::fmt::Write as _;
use std::sync::{Arc, Mutex, MutexGuard};
const UNMATCHED_ROUTE: &str = "unmatched";
const MAX_ROUTE_LABELS: usize = 128;
pub(super) struct ServerMetrics {
recorder: Arc<DefaultMetricsRecorder>,
routes: Mutex<RouteLabels>,
requests: Mutex<HashMap<&'static str, RequestInstruments>>,
busy_uploads: Arc<dyn CounterHandle>,
busy_downloads: Arc<dyn CounterHandle>,
}
struct RequestInstruments {
seconds: Arc<dyn HistogramHandle>,
served: HashMap<(&'static str, &'static str), Arc<dyn CounterHandle>>,
}
impl ServerMetrics {
pub(super) fn new() -> Arc<Self> {
let recorder = Arc::new(DefaultMetricsRecorder::new());
let busy = |kind: &'static str| {
recorder.register_counter(
"loonfs.server.busy_rejections",
"Requests refused at a concurrency limit",
&[("kind", kind)],
)
};
Arc::new(Self {
busy_uploads: busy("upload"),
busy_downloads: busy("download"),
routes: Mutex::new(RouteLabels::default()),
requests: Mutex::new(HashMap::new()),
recorder,
})
}
pub(super) fn recorder(&self) -> Arc<dyn MetricsRecorder> {
Arc::clone(&self.recorder) as Arc<dyn MetricsRecorder>
}
pub(super) fn request_served(
&self,
matched_route: Option<&str>,
method: &Method,
status: StatusCode,
elapsed_seconds: f64,
) {
let route = self.route_label(matched_route);
let (seconds, served) = {
let mut requests = lock(&self.requests);
let instruments = requests
.entry(route)
.or_insert_with(|| RequestInstruments::register(self.recorder.as_ref(), route));
instruments.served(
self.recorder.as_ref(),
route,
method_label(method),
status_class_label(status),
)
};
served.increment(1);
seconds.record(elapsed_seconds);
}
pub(super) fn upload_rejected_as_busy(&self) {
self.busy_uploads.increment(1);
}
pub(super) fn download_rejected_as_busy(&self) {
self.busy_downloads.increment(1);
}
pub(super) fn render(
&self,
cache: &RuntimeCacheStats,
upload_permits: usize,
download_permits: usize,
) -> String {
let mut rendered = render_snapshot(&self.recorder.snapshot());
render_scrape_gauges(&mut rendered, cache, upload_permits, download_permits);
rendered
}
fn route_label(&self, matched_route: Option<&str>) -> &'static str {
match matched_route {
Some(route) => lock(&self.routes).intern(route),
None => UNMATCHED_ROUTE,
}
}
}
impl RequestInstruments {
fn register(recorder: &dyn MetricsRecorder, route: &'static str) -> Self {
Self {
seconds: recorder.register_histogram(
"loonfs.server.request_seconds",
"Time to serve one request, by matched route",
&[("route", route)],
LATENCY_SECONDS_BOUNDARIES,
),
served: HashMap::new(),
}
}
fn served(
&mut self,
recorder: &dyn MetricsRecorder,
route: &'static str,
method: &'static str,
status_class: &'static str,
) -> (Arc<dyn HistogramHandle>, Arc<dyn CounterHandle>) {
let counter = self
.served
.entry((method, status_class))
.or_insert_with(|| {
recorder.register_counter(
"loonfs.server.requests",
"Requests served, by matched route, method, and status class",
&[
("route", route),
("method", method),
("status_class", status_class),
],
)
})
.clone();
(Arc::clone(&self.seconds), counter)
}
}
#[derive(Default)]
struct RouteLabels {
interned: HashMap<String, &'static str>,
}
impl RouteLabels {
fn intern(&mut self, route: &str) -> &'static str {
if let Some(label) = self.interned.get(route) {
return label;
}
if self.interned.len() >= MAX_ROUTE_LABELS {
return UNMATCHED_ROUTE;
}
let label: &'static str = Box::leak(route.to_owned().into_boxed_str());
self.interned.insert(route.to_owned(), label);
label
}
}
fn method_label(method: &Method) -> &'static str {
match *method {
Method::GET => "GET",
Method::HEAD => "HEAD",
Method::POST => "POST",
Method::PUT => "PUT",
Method::DELETE => "DELETE",
Method::OPTIONS => "OPTIONS",
Method::PATCH => "PATCH",
_ => "other",
}
}
fn status_class_label(status: StatusCode) -> &'static str {
match status.as_u16() / 100 {
1 => "1xx",
2 => "2xx",
3 => "3xx",
4 => "4xx",
5 => "5xx",
_ => "other",
}
}
fn render_snapshot(snapshot: &MetricsSnapshot) -> String {
let mut rendered = String::new();
let mut current_name = None;
for entry in snapshot.all() {
let name = prometheus_name(entry);
if current_name.as_deref() != Some(name.as_str()) {
write_metric_header(&mut rendered, &name, entry);
current_name = Some(name.clone());
}
write_entry(&mut rendered, &name, entry);
}
rendered
}
fn write_metric_header(rendered: &mut String, name: &str, entry: &MetricEntry) {
let kind = match entry.value {
MetricValue::Counter(_) => "counter",
MetricValue::Gauge(_) => "gauge",
MetricValue::Histogram { .. } => "histogram",
};
let _ = writeln!(rendered, "# HELP {name} {}", escape_help(entry.description));
let _ = writeln!(rendered, "# TYPE {name} {kind}");
}
fn write_entry(rendered: &mut String, name: &str, entry: &MetricEntry) {
match &entry.value {
MetricValue::Counter(value) => {
let _ = writeln!(rendered, "{name}{} {value}", labels(&entry.labels, None));
}
MetricValue::Gauge(value) => {
let _ = writeln!(rendered, "{name}{} {value}", labels(&entry.labels, None));
}
MetricValue::Histogram {
boundaries,
bucket_counts,
count,
sum,
} => {
let mut cumulative = 0u64;
for (boundary, filed) in boundaries.iter().zip(bucket_counts.iter()) {
cumulative += filed;
let _ = writeln!(
rendered,
"{name}_bucket{} {cumulative}",
labels(&entry.labels, Some(&format_float(*boundary)))
);
}
let _ = writeln!(
rendered,
"{name}_bucket{} {count}",
labels(&entry.labels, Some("+Inf"))
);
let _ = writeln!(
rendered,
"{name}_sum{} {}",
labels(&entry.labels, None),
format_float(*sum)
);
let _ = writeln!(
rendered,
"{name}_count{} {count}",
labels(&entry.labels, None)
);
}
}
}
fn render_scrape_gauges(
rendered: &mut String,
cache: &RuntimeCacheStats,
upload_permits: usize,
download_permits: usize,
) {
for (field, value) in cache_gauges(cache) {
let name = format!("loonfs_cache_{field}");
let _ = writeln!(rendered, "# HELP {name} Runtime cache counter `{field}`");
let _ = writeln!(rendered, "# TYPE {name} gauge");
let _ = writeln!(rendered, "{name} {value}");
}
for (name, description, value) in [
(
"loonfs_server_upload_permits_available",
"Proxied-upload slots free right now",
upload_permits,
),
(
"loonfs_server_download_permits_available",
"Proxied-content-read slots free right now",
download_permits,
),
] {
let _ = writeln!(rendered, "# HELP {name} {description}");
let _ = writeln!(rendered, "# TYPE {name} gauge");
let _ = writeln!(rendered, "{name} {value}");
}
}
fn cache_gauges(stats: &RuntimeCacheStats) -> [(&'static str, usize); 18] {
let RuntimeCacheStats {
latest_metadata_view_reads,
wal_tail_projection_cache_hits,
wal_tail_projection_cache_misses,
wal_tail_projection_cache_inserts,
wal_tail_projection_cache_evictions,
wal_tail_projection_cache_evicted_rows,
wal_tail_projection_cache_evicted_decoded_bytes,
wal_tail_projection_cache_uncacheable_count,
wal_tail_projection_cache_uncacheable_rows,
wal_tail_projection_cache_uncacheable_decoded_bytes,
wal_tail_projection_cache_cached_rows,
wal_tail_projection_cache_cached_decoded_bytes,
metadata_table_cache_hits,
metadata_table_cache_misses,
metadata_table_cache_inserts,
metadata_table_cache_evictions,
metadata_table_cache_filter_skips,
metadata_table_cache_filter_false_positives,
} = *stats;
[
("latest_metadata_view_reads", latest_metadata_view_reads),
(
"wal_tail_projection_cache_hits",
wal_tail_projection_cache_hits,
),
(
"wal_tail_projection_cache_misses",
wal_tail_projection_cache_misses,
),
(
"wal_tail_projection_cache_inserts",
wal_tail_projection_cache_inserts,
),
(
"wal_tail_projection_cache_evictions",
wal_tail_projection_cache_evictions,
),
(
"wal_tail_projection_cache_evicted_rows",
wal_tail_projection_cache_evicted_rows,
),
(
"wal_tail_projection_cache_evicted_decoded_bytes",
wal_tail_projection_cache_evicted_decoded_bytes,
),
(
"wal_tail_projection_cache_uncacheable_count",
wal_tail_projection_cache_uncacheable_count,
),
(
"wal_tail_projection_cache_uncacheable_rows",
wal_tail_projection_cache_uncacheable_rows,
),
(
"wal_tail_projection_cache_uncacheable_decoded_bytes",
wal_tail_projection_cache_uncacheable_decoded_bytes,
),
(
"wal_tail_projection_cache_cached_rows",
wal_tail_projection_cache_cached_rows,
),
(
"wal_tail_projection_cache_cached_decoded_bytes",
wal_tail_projection_cache_cached_decoded_bytes,
),
("metadata_table_cache_hits", metadata_table_cache_hits),
("metadata_table_cache_misses", metadata_table_cache_misses),
("metadata_table_cache_inserts", metadata_table_cache_inserts),
(
"metadata_table_cache_evictions",
metadata_table_cache_evictions,
),
(
"metadata_table_cache_filter_skips",
metadata_table_cache_filter_skips,
),
(
"metadata_table_cache_filter_false_positives",
metadata_table_cache_filter_false_positives,
),
]
}
fn prometheus_name(entry: &MetricEntry) -> String {
let mut name = entry.name.replace('.', "_");
if matches!(entry.value, MetricValue::Counter(_)) {
name.push_str("_total");
}
name
}
fn labels(labels: &[(&'static str, &'static str)], le: Option<&str>) -> String {
if labels.is_empty() && le.is_none() {
return String::new();
}
let mut rendered = String::from("{");
for (index, (key, value)) in labels.iter().enumerate() {
if index > 0 {
rendered.push(',');
}
let _ = write!(rendered, "{key}=\"{}\"", escape_label(value));
}
if let Some(le) = le {
if !labels.is_empty() {
rendered.push(',');
}
let _ = write!(rendered, "le=\"{le}\"");
}
rendered.push('}');
rendered
}
fn format_float(value: f64) -> String {
if value.fract() == 0.0 && value.abs() < 1e15 {
format!("{value:.1}")
} else {
format!("{value}")
}
}
fn escape_label(value: &str) -> String {
value
.replace('\\', "\\\\")
.replace('"', "\\\"")
.replace('\n', "\\n")
}
fn escape_help(value: &str) -> String {
value.replace('\\', "\\\\").replace('\n', "\\n")
}
fn lock<T>(table: &Mutex<T>) -> MutexGuard<'_, T> {
table
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
}
#[cfg(test)]
mod tests;