From 8089510f0030142410dc42fee0db148fdf2a49ed Mon Sep 17 00:00:00 2001 From: phil Date: Tue, 28 Jul 2026 12:23:27 -0400 Subject: [PATCH] track in-flight over response streaming --- Cargo.lock | 40 +++++++++++++++ hubble/Cargo.toml | 2 + hubble/src/metrics.rs | 2 +- hubble/src/serve/mod.rs | 109 ++++++++++++++++++++++++++++++++++++++-- 4 files changed, 148 insertions(+), 5 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index b3df0f4..72e37ed 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1360,6 +1360,12 @@ dependencies = [ "cfg-if", ] +[[package]] +name = "endian-type" +version = "0.1.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c34f04666d835ff5d62e058c3995147c06f42fe86ff053337632bca83e42702d" + [[package]] name = "enum_dispatch" version = "0.3.13" @@ -2088,6 +2094,7 @@ dependencies = [ "console-subscriber", "headers-accept", "http", + "http-body", "hubble-sync", "hubble-sync-axum", "hubble-sync-rocksdb", @@ -2099,6 +2106,7 @@ dependencies = [ "mediatype", "metrics", "metrics-exporter-prometheus", + "metrics-util", "miette", "mimalloc", "rlimit", @@ -3172,11 +3180,15 @@ version = "0.20.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "96f8722f8562635f92f8ed992f26df0532266eb03d5202607c20c0d7e9745e13" dependencies = [ + "aho-corasick", "crossbeam-epoch", "crossbeam-utils", "hashbrown 0.16.1", + "indexmap", "metrics", + "ordered-float", "quanta", + "radix_trie", "rand 0.9.4", "rand_xoshiro", "rapidhash", @@ -3327,6 +3339,15 @@ version = "1.0.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "650eef8c711430f1a879fdd01d4745a7deea475becfb90269c06775983bbf086" +[[package]] +name = "nibble_vec" +version = "0.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "77a5d83df9f36fe23f0c3648c6bbb8b0298bb5f1939c8f2704431371f4b84d43" +dependencies = [ + "smallvec", +] + [[package]] name = "nom" version = "7.1.3" @@ -3446,6 +3467,15 @@ version = "0.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7c87def4c32ab89d880effc9e097653c8da5d6ef28e6b539d313baaacfbafcbe" +[[package]] +name = "ordered-float" +version = "5.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b7d950ca161dc355eaf28f82b11345ed76c6e1f6eb1f4f4479e0323b9e2fbd0e" +dependencies = [ + "num-traits", +] + [[package]] name = "owo-colors" version = "4.3.0" @@ -3850,6 +3880,16 @@ version = "5.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "69cdb34c158ceb288df11e18b4bd39de994f6657d83847bdffdbd7f346754b0f" +[[package]] +name = "radix_trie" +version = "0.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c069c179fcdc6a2fe24d8d18305cf085fdbd4f922c041943e203685d6a1c58fd" +dependencies = [ + "endian-type", + "nibble_vec", +] + [[package]] name = "rand" version = "0.8.6" diff --git a/hubble/Cargo.toml b/hubble/Cargo.toml index 06f8e12..9143529 100644 --- a/hubble/Cargo.toml +++ b/hubble/Cargo.toml @@ -23,6 +23,7 @@ jacquard-common = { workspace = true } jacquard-derive = { workspace = true } jacquard-lexicon = { workspace = true } headers-accept = "0.3.0" +http-body = "1" hubble-sync = { workspace = true } hubble-sync-axum = { workspace = true } hubble-sync-rocksdb = { workspace = true } @@ -55,4 +56,5 @@ features = ["http-listener"] [dev-dependencies] jacquard-lexicon = { workspace = true, features = ["codegen"] } +metrics-util = "0.20" tempfile = "3" diff --git a/hubble/src/metrics.rs b/hubble/src/metrics.rs index e55865e..d77fb06 100644 --- a/hubble/src/metrics.rs +++ b/hubble/src/metrics.rs @@ -72,7 +72,7 @@ pub fn describe_metrics() { ); describe_gauge!( HTTP_REQUESTS_IN_FLIGHT, - "http requests currently being handled, by endpoint" + "http requests in flight incl. response body streaming, by endpoint" ); describe_counter!( GETREPO_RESPONSES_TOTAL, diff --git a/hubble/src/serve/mod.rs b/hubble/src/serve/mod.rs index 47edf46..1bfc953 100644 --- a/hubble/src/serve/mod.rs +++ b/hubble/src/serve/mod.rs @@ -5,13 +5,18 @@ mod get_stats; mod hello; use std::path::PathBuf; +use std::pin::Pin; use std::sync::Arc; +use std::task::{Context, Poll}; use std::time::{Instant, SystemTime}; +use axum::body::Body; use axum::extract::{MatchedPath, Request}; use axum::middleware::Next; use axum::response::Response; use axum::{Router, routing::get}; +use bytes::Bytes; +use http_body::{Frame, SizeHint}; use hubble_sync::{DaslCid, RepoView, SyncHandle, Tid}; use hubble_sync_rocksdb::RocksEngine; use jacquard_common::deps::chrono; @@ -87,10 +92,12 @@ pub fn route(state: Arc) -> Router { .layer(axum::middleware::from_fn(track_request)) } -/// request metrics: a duration histogram labelled by method, endpoint, status +/// request metrics: a duration histogram labelled by method, endpoint, status, +/// and an in-flight gauge /// /// duration is measured to response *headers* -- streaming bodies (eg. getRepo -/// archives) keep transferring beyond what's recorded here +/// archives) keep transferring beyond what's recorded there. the in-flight +/// gauge, by contrast, rides the response body and covers the full transfer. async fn track_request(req: Request, next: Next) -> Response { // every label value comes from a closed set (methods, the route table, // status codes) so hostile requests can't blow up metric cardinality @@ -113,7 +120,6 @@ async fn track_request(req: Request, next: Next) -> Response { let start = Instant::now(); let in_flight = InFlightGuard::start(endpoint.clone()); let res = next.run(req).await; - drop(in_flight); // Drop also fires if this future is cancelled mid-await histogram!( HTTP_REQUEST_DURATION_SECONDS, @@ -123,5 +129,100 @@ async fn track_request(req: Request, next: Next) -> Response { ) .record(start.elapsed().as_secs_f64()); - res + // the gauge guard rides the response body from here, releasing when the + // body completes (or drops on client disconnect) -- so the gauge counts + // *active transfers*, not just handler time. a request future cancelled + // before this point drops the guard all the same. + res.map(|inner| { + Body::new(TrackedBody { + inner, + guard: in_flight, + }) + }) +} + +/// response body wrapper holding the in-flight gauge guard for as long as the +/// body is alive -- the middleware future itself ends at response headers, +/// long before a streaming export finishes transferring +struct TrackedBody { + inner: Body, + #[expect(dead_code, reason = "gauge guard released on body drop")] + guard: InFlightGuard, +} + +impl http_body::Body for TrackedBody { + type Data = Bytes; + type Error = axum::Error; + + fn poll_frame( + self: Pin<&mut Self>, + cx: &mut Context<'_>, + ) -> Poll, Self::Error>>> { + Pin::new(&mut self.get_mut().inner).poll_frame(cx) + } + + fn is_end_stream(&self) -> bool { + self.inner.is_end_stream() + } + + fn size_hint(&self) -> SizeHint { + self.inner.size_hint() + } +} + +#[cfg(test)] +mod tests { + use super::*; + use metrics_util::debugging::{DebugValue, DebuggingRecorder, Snapshotter}; + + fn in_flight_value(snapshotter: &Snapshotter) -> f64 { + snapshotter + .snapshot() + .into_vec() + .into_iter() + .find(|(k, ..)| k.key().name() == crate::metrics::HTTP_REQUESTS_IN_FLIGHT) + .map(|(.., v)| match v { + DebugValue::Gauge(g) => g.into_inner(), + other => panic!("expected gauge, got {other:?}"), + }) + .unwrap_or(0.0) + } + + fn tracked(recorder: &DebuggingRecorder) -> Body { + metrics::with_local_recorder(recorder, || { + Body::new(TrackedBody { + inner: Body::from("hello"), + guard: InFlightGuard::start("/test".to_owned()), + }) + }) + } + + // the regression these guard: the gauge used to release at response + // headers, undercounting active transfers (64 concurrent streaming exports + // read as ~4 in-flight). one snapshot per test: the debugging recorder's + // snapshots don't observe gauges non-destructively. + + #[test] + fn in_flight_gauge_held_while_body_alive() { + let recorder = DebuggingRecorder::new(); + let snapshotter = recorder.snapshotter(); + let body = tracked(&recorder); + assert_eq!(in_flight_value(&snapshotter), 1.0); + drop(body); + } + + #[tokio::test] + async fn in_flight_gauge_released_when_body_consumed() { + let recorder = DebuggingRecorder::new(); + let snapshotter = recorder.snapshotter(); + + let body = tracked(&recorder); + let bytes = axum::body::to_bytes(body, 16).await.expect("collect body"); + assert_eq!(&bytes[..], b"hello", "frames forwarded unchanged"); + assert_eq!( + in_flight_value(&snapshotter), + 0.0, + "released once the body is consumed and dropped", + ); + } } -- 2.51.2