From f55ce4885ba726bfadcd35b64161431676076002 Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Mon, 14 Sep 2026 10:05:48 +0000 Subject: [PATCH] fix(metrics): isolate scrape rendering from async workers Each metrics scrape traverses repository directories and encodes the response synchronously. Slow filesystem metadata operations could occupy a worker that also serves WebSocket and Git requests. Render on a blocking worker. Acquire a shared permit before spawning it and keep that permit inside the worker even if the HTTP request is canceled, so concurrent scrapes cannot start overlapping directory scans. Report a worker failure as an HTTP 500 response rather than losing the request task. Preserve the synchronous rendering API, metric contents and monitoring configuration. This does not cache metrics or claim a measured production speedup; it prevents scrape work from running on the async executor. Validation: a current-thread test checks that a queued scrape yields, capacity is released and async rendering preserves repository counts. --- CHANGELOG.md | 3 +++ docs/explanation/architecture.md | 2 ++ src/http/mod.rs | 13 ++++++++- src/metrics/mod.rs | 45 ++++++++++++++++++++++++++++++-- 4 files changed, 60 insertions(+), 3 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 1a7652e..311a2a2 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Fixed +- Render metrics on a blocking worker and serialize scrapes so repository + counting does not occupy async workers serving relay and Git requests. + - Keep periodic recovery checkpoints off async workers, release snapshot locks before disk I/O, and finish active checkpoints before the shutdown snapshot. diff --git a/docs/explanation/architecture.md b/docs/explanation/architecture.md index 3058ca3..3186132 100644 --- a/docs/explanation/architecture.md +++ b/docs/explanation/architecture.md @@ -78,6 +78,8 @@ runtime: - Disable Nagle on accepted sockets so small EVENT/EOSE and streaming Git writes do not wait for the peer's delayed acknowledgements, including when the immediate peer is a local reverse proxy +- Render metrics on a blocking worker, serializing concurrent scrapes before + spawning filesystem scans so monitoring cannot occupy the network executor - Initialize Nostr relay builder with custom [`Nip34WritePolicy`](src/nostr/builder.rs:51) - Set up shared storage (LMDB or Memory), purgatory, sync manager, and background maintenance tasks diff --git a/src/http/mod.rs b/src/http/mod.rs index 1c2c2c6..7523073 100644 --- a/src/http/mod.rs +++ b/src/http/mod.rs @@ -923,7 +923,18 @@ impl Service> for HttpService { if let Some(ref metrics) = self.metrics { let metrics = metrics.clone(); return Box::pin(async move { - let output = metrics.render(); + let output = match metrics.render_async().await { + Ok(output) => output, + Err(error) => { + tracing::error!(%error, "Metrics render worker failed"); + return Ok(add_cors_headers( + Response::builder().header("server", "ngit-grasp"), + ) + .status(500) + .body(full_body("Metrics temporarily unavailable")) + .unwrap()); + } + }; Ok( add_cors_headers(Response::builder().header("server", "ngit-grasp")) .status(200) diff --git a/src/metrics/mod.rs b/src/metrics/mod.rs index bdd79f9..b9644a6 100644 --- a/src/metrics/mod.rs +++ b/src/metrics/mod.rs @@ -564,6 +564,7 @@ pub fn record_manual_ejection( #[derive(Clone)] pub struct Metrics { inner: Arc, + render_permit: Arc, } struct MetricsInner { @@ -628,6 +629,7 @@ impl Metrics { let inner = MetricsInner::new(abuse_threshold, git_data_path); Self { inner: Arc::new(inner), + render_permit: Arc::new(tokio::sync::Semaphore::new(1)), } } @@ -794,6 +796,24 @@ impl Metrics { // === Rendering === + /// Render a scrape without blocking the async executor on filesystem + /// traversal or encoding. Concurrent scrapes wait before spawning work; + /// cancelling an HTTP request cannot release a running worker's permit. + pub async fn render_async(&self) -> Result { + let permit = self + .render_permit + .clone() + .acquire_owned() + .await + .expect("metrics render semaphore is never closed"); + let metrics = self.clone(); + tokio::task::spawn_blocking(move || { + let _permit = permit; + metrics.render() + }) + .await + } + /// Render all metrics in Prometheus text format. /// /// This method: @@ -1104,6 +1124,21 @@ impl Drop for GitOperationTimer { mod tests { use super::*; + async fn concurrent_scrapes_wait_without_blocking_the_runtime(metrics: &Metrics) { + use futures_util::FutureExt; + let permit = metrics.render_permit.clone().acquire_owned().await.unwrap(); + let render = metrics.render_async(); + tokio::pin!(render); + assert!(render.as_mut().now_or_never().is_none()); + drop(permit); + let output = tokio::time::timeout(std::time::Duration::from_secs(5), render) + .await + .unwrap() + .unwrap(); + assert!(output.contains("ngit_uptime_seconds")); + assert_eq!(metrics.render_permit.available_permits(), 1); + } + #[test] fn test_count_repositories_on_disk() { use std::fs; @@ -1153,8 +1188,8 @@ mod tests { /// /// If additional Metrics tests are needed, they should either be added to this /// test or use a separate test-specific Prometheus registry. - #[test] - fn test_metrics_with_repository_counting() { + #[tokio::test] + async fn test_metrics_with_repository_counting() { use std::fs; use tempfile::TempDir; @@ -1211,5 +1246,11 @@ mod tests { // Render should count 3 repos let output = metrics.render(); assert!(output.contains("ngit_repositories_total 3")); + concurrent_scrapes_wait_without_blocking_the_runtime(&metrics).await; + assert!(metrics + .render_async() + .await + .unwrap() + .contains("ngit_repositories_total 3")); } }