-
Notifications
You must be signed in to change notification settings - Fork 1
feat(core): add shared request timing collector and middleware #389
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: main
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -1,4 +1,5 @@ | ||
| use std::future::Future; | ||
| use std::marker::PhantomData; | ||
| use std::sync::Arc; | ||
| use web_time::Instant; | ||
|
|
||
|
|
@@ -8,6 +9,7 @@ use crate::context::RequestContext; | |
| use crate::error::EdgeError; | ||
| use crate::handler::DynHandler; | ||
| use crate::http::Response; | ||
| use crate::request_timing::RequestTimings; | ||
|
|
||
| pub type BoxMiddleware = Arc<dyn Middleware>; | ||
|
|
||
|
|
@@ -71,6 +73,56 @@ impl<'mw> Next<'mw> { | |
| } | ||
| } | ||
|
|
||
| /// Opt-in, insert-if-absent attachment of a request-local timing handle. | ||
| /// | ||
| /// Register before timing consumers. This never finalizes timings or renders | ||
| /// headers. Routes are matched before middleware, so unmatched 404/405 requests | ||
| /// bypass attachment. Do not install collectors with `RouterBuilder::with_state`. | ||
| pub struct RequestTimingMiddleware<const N: usize, D = ()> { | ||
| data: PhantomData<fn() -> D>, | ||
| excluded_paths: &'static [&'static str], | ||
| } | ||
|
|
||
| impl<const N: usize, D> Default for RequestTimingMiddleware<N, D> { | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 🤔 pub const TIMING: RequestTimingMiddleware<1> = RequestTimingMiddleware::default();
// error[E0015]: cannot call non-const associated function `<RequestTimingMiddleware<1> as Default>::default` in constantsA impl<const N: usize, D> RequestTimingMiddleware<N, D> {
#[must_use]
#[inline]
pub const fn new() -> Self {
Self { data: PhantomData, excluded_paths: &[] }
}
#[must_use]
#[inline]
pub const fn with_excluded_paths(mut self, paths: &'static [&'static str]) -> Self {
self.excluded_paths = paths;
self
}
}
// my_app_core/src/timing.rs
pub const TIMING: RequestTimingMiddleware<4, AppData> =
RequestTimingMiddleware::new().with_excluded_paths(&["/health"]);middleware = ["my_app_core::timing::TIMING"] |
||
| #[inline] | ||
| fn default() -> Self { | ||
| Self { | ||
| data: PhantomData, | ||
| excluded_paths: &[], | ||
| } | ||
| } | ||
| } | ||
|
|
||
| impl<const N: usize, D> RequestTimingMiddleware<N, D> { | ||
| /// Excludes exact URI paths for every method, ignoring query strings. | ||
| /// Existing handles are preserved even on excluded paths. | ||
| #[must_use] | ||
| #[inline] | ||
| pub fn with_excluded_paths(mut self, paths: &'static [&'static str]) -> Self { | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. ⛏ / 🌱 Exclusion list semantics and shape:
|
||
| self.excluded_paths = paths; | ||
| self | ||
| } | ||
| } | ||
|
|
||
| #[async_trait(?Send)] | ||
| impl<const N: usize, D: Default + Send + 'static> Middleware for RequestTimingMiddleware<N, D> { | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 🌱 A middleware-created handle can't be reached after the router returns: a probe confirmed the response carries no |
||
| #[inline] | ||
| async fn handle(&self, mut ctx: RequestContext, next: Next<'_>) -> Result<Response, EdgeError> { | ||
| if ctx | ||
| .request() | ||
| .extensions() | ||
| .get::<RequestTimings<N, D>>() | ||
| .is_none() | ||
| && !self.excluded_paths.contains(&ctx.request().uri().path()) | ||
| { | ||
| ctx.request_mut() | ||
| .extensions_mut() | ||
| .insert(RequestTimings::<N, D>::new()); | ||
| } | ||
| next.run(ctx).await | ||
| } | ||
| } | ||
|
|
||
| pub struct RequestLogger; | ||
|
|
||
| #[async_trait(?Send)] | ||
|
|
@@ -121,6 +173,250 @@ where | |
| FnMiddleware::new(func) | ||
| } | ||
|
|
||
| #[cfg(test)] | ||
| mod request_timing_tests { | ||
| use super::*; | ||
| use crate::body::Body; | ||
| use crate::handler::IntoHandler as _; | ||
| use crate::http::{Method, Request, StatusCode, request_builder}; | ||
| use crate::params::PathParams; | ||
| use crate::response::response_with_body; | ||
| use crate::router::RouterService; | ||
| use futures::executor::block_on; | ||
| use std::sync::Mutex; | ||
| use std::time::Duration; | ||
| use tower_service::Service as _; | ||
|
|
||
| type Timings = RequestTimings<1>; | ||
|
|
||
| struct BeforeAttachment; | ||
| #[async_trait(?Send)] | ||
| impl Middleware for BeforeAttachment { | ||
| async fn handle(&self, ctx: RequestContext, next: Next<'_>) -> Result<Response, EdgeError> { | ||
| assert!(ctx.request().extensions().get::<Timings>().is_none()); | ||
| next.run(ctx).await | ||
| } | ||
| } | ||
|
|
||
| struct Capture { | ||
| handles: Arc<Mutex<Vec<Timings>>>, | ||
| short_circuit: bool, | ||
| } | ||
|
|
||
| #[async_trait(?Send)] | ||
| impl Middleware for Capture { | ||
| async fn handle(&self, ctx: RequestContext, next: Next<'_>) -> Result<Response, EdgeError> { | ||
| let timings = ctx.request().extensions().get::<Timings>().unwrap().clone(); | ||
| self.handles.lock().unwrap().push(timings); | ||
| if self.short_circuit { | ||
| response_with_body(StatusCode::UNAUTHORIZED, Body::text("stopped")) | ||
| } else { | ||
| next.run(ctx).await | ||
| } | ||
| } | ||
| } | ||
|
|
||
| #[crate::action] | ||
| async fn presence(ctx: RequestContext) -> Result<Response, EdgeError> { | ||
| let installed = ctx.request().extensions().get::<Timings>(); | ||
| if let Some(timings) = installed { | ||
| timings.record(0, Duration::from_millis(3)).unwrap(); | ||
| } | ||
| response_with_body( | ||
| StatusCode::OK, | ||
| Body::text(if installed.is_some() { "yes" } else { "no" }), | ||
| ) | ||
| } | ||
|
|
||
| #[crate::action] | ||
| async fn failure(_ctx: RequestContext) -> Result<Response, EdgeError> { | ||
| Err(EdgeError::bad_request("unchanged timing error")) | ||
| } | ||
|
|
||
| fn request(method: Method, uri: &str) -> Request { | ||
| request_builder() | ||
| .method(method) | ||
| .uri(uri) | ||
| .body(Body::empty()) | ||
| .unwrap() | ||
| } | ||
|
|
||
| #[test] | ||
| fn request_timing_exclusions_are_exact_paths_for_every_method_and_query() { | ||
| let router = RouterService::builder() | ||
| .middleware(RequestTimingMiddleware::<1>::default().with_excluded_paths(&["/health"])) | ||
| .get("/health", presence) | ||
| .post("/health", presence) | ||
| .get("/health/", presence) | ||
| .get("/healthy", presence) | ||
| .build(); | ||
| for method in [Method::GET, Method::POST] { | ||
| for uri in ["/health", "/health?check=1"] { | ||
| let response = block_on(router.oneshot(request(method.clone(), uri))).unwrap(); | ||
| assert_eq!(response.body().as_bytes().unwrap(), b"no"); | ||
| let timings = Timings::new(); | ||
| timings.record(0, Duration::from_millis(8)).unwrap(); | ||
| let mut preinstalled = request(method.clone(), uri); | ||
| preinstalled.extensions_mut().insert(timings.clone()); | ||
| let preserved = block_on(router.oneshot(preinstalled)).unwrap(); | ||
| assert_eq!(preserved.body().as_bytes().unwrap(), b"yes"); | ||
| assert_eq!( | ||
| timings.snapshot(|view| view.phases).unwrap(), | ||
| [Some(Duration::from_millis(11))] | ||
| ); | ||
| } | ||
| } | ||
| for uri in ["/health/", "/healthy"] { | ||
| let response = block_on(router.oneshot(request(Method::GET, uri))).unwrap(); | ||
| assert_eq!(response.body().as_bytes().unwrap(), b"yes"); | ||
| } | ||
| let default_router = RouterService::builder() | ||
| .middleware(RequestTimingMiddleware::<1>::default()) | ||
| .get("/health", presence) | ||
| .build(); | ||
| let response = block_on(default_router.oneshot(request(Method::GET, "/health"))).unwrap(); | ||
| assert_eq!(response.body().as_bytes().unwrap(), b"yes"); | ||
| } | ||
|
|
||
| #[test] | ||
| fn request_timing_duplicate_installation_preserves_state_and_requests_are_fresh() { | ||
| let handles = Arc::new(Mutex::new(Vec::<Timings>::new())); | ||
| let router = RouterService::builder() | ||
| .middleware(RequestTimingMiddleware::<1>::default()) | ||
| .middleware(Capture { | ||
| handles: Arc::clone(&handles), | ||
| short_circuit: false, | ||
| }) | ||
| .middleware(RequestTimingMiddleware::<1>::default()) | ||
| .get("/timed", presence) | ||
| .build(); | ||
| let timings = Timings::new(); | ||
| timings.record(0, Duration::from_millis(8)).unwrap(); | ||
| let mut preinstalled = request(Method::GET, "/timed"); | ||
| preinstalled.extensions_mut().insert(timings.clone()); | ||
| block_on(router.oneshot(preinstalled)).unwrap(); | ||
| assert_eq!( | ||
| timings.snapshot(|view| view.phases).unwrap(), | ||
| [Some(Duration::from_millis(11))] | ||
| ); | ||
| for _ in 0_u8..2 { | ||
| block_on(router.oneshot(request(Method::GET, "/timed"))).unwrap(); | ||
| } | ||
| let saved = handles.lock().unwrap(); | ||
| assert_eq!(saved.len(), 3); | ||
| saved[1].record(0, Duration::from_millis(5)).unwrap(); | ||
| assert_eq!( | ||
| saved[1].snapshot(|view| view.phases).unwrap(), | ||
| [Some(Duration::from_millis(8))] | ||
| ); | ||
| assert_eq!( | ||
| saved[2].snapshot(|view| view.phases).unwrap(), | ||
| [Some(Duration::from_millis(3))] | ||
| ); | ||
| for handle in saved.iter() { | ||
| assert_eq!( | ||
| handle | ||
| .snapshot(|view| (view.headers_ready_total, view.request_elapsed)) | ||
| .unwrap(), | ||
| (None, None) | ||
| ); | ||
| } | ||
| } | ||
|
|
||
| #[test] | ||
| fn request_timing_preserves_short_circuits_and_errors_without_finalization() { | ||
| let handles = Arc::new(Mutex::new(Vec::<Timings>::new())); | ||
| let short_router = RouterService::builder() | ||
| .middleware(RequestTimingMiddleware::<1>::default()) | ||
| .middleware(Capture { | ||
| handles: Arc::clone(&handles), | ||
| short_circuit: true, | ||
| }) | ||
| .get("/timed", presence) | ||
| .build(); | ||
| let response = block_on(short_router.oneshot(request(Method::GET, "/timed"))).unwrap(); | ||
| assert_eq!(response.status(), StatusCode::UNAUTHORIZED); | ||
| assert_eq!(response.body().as_bytes().unwrap(), b"stopped"); | ||
| assert_eq!( | ||
| handles.lock().unwrap()[0] | ||
| .snapshot(|view| (view.phases, view.request_elapsed)) | ||
| .unwrap(), | ||
| ([None], None) | ||
| ); | ||
|
|
||
| let mut router = RouterService::builder() | ||
| .middleware(RequestTimingMiddleware::<1>::default()) | ||
| .middleware(Capture { | ||
| handles: Arc::clone(&handles), | ||
| short_circuit: false, | ||
| }) | ||
| .get("/error", failure) | ||
| .build(); | ||
| let error = block_on(router.call(request(Method::GET, "/error"))).unwrap_err(); | ||
| assert_eq!(error.status(), StatusCode::BAD_REQUEST); | ||
| assert_eq!(error.message(), "unchanged timing error"); | ||
| let rendered = block_on(router.oneshot(request(Method::GET, "/error"))).unwrap(); | ||
| assert_eq!(rendered.status(), StatusCode::BAD_REQUEST); | ||
| assert!(!rendered.headers().contains_key("server-timing")); | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. ⛏ Assertions that can't fail:
|
||
| let saved = handles.lock().unwrap(); | ||
| assert_eq!(saved.len(), 3); | ||
| for handle in saved.iter() { | ||
| assert_eq!( | ||
| handle | ||
| .snapshot(|view| (view.headers_ready_total, view.request_elapsed)) | ||
| .unwrap(), | ||
| (None, None) | ||
| ); | ||
| } | ||
| } | ||
|
|
||
| #[test] | ||
| fn request_timing_runs_after_prior_middleware_and_unmatched_routes_bypass_it() { | ||
| let handles = Arc::new(Mutex::new(Vec::<Timings>::new())); | ||
| let router = RouterService::builder() | ||
| .middleware(RequestTimingMiddleware::<1>::default()) | ||
| .middleware(Capture { | ||
| handles: Arc::clone(&handles), | ||
| short_circuit: false, | ||
| }) | ||
| .get("/timed", presence) | ||
| .build(); | ||
| for (method, uri, status) in [ | ||
| (Method::GET, "/missing", StatusCode::NOT_FOUND), | ||
| (Method::POST, "/timed", StatusCode::METHOD_NOT_ALLOWED), | ||
| ] { | ||
| let response = block_on(router.oneshot(request(method, uri))).unwrap(); | ||
| assert_eq!(response.status(), status); | ||
| } | ||
| assert!(handles.lock().unwrap().is_empty()); | ||
|
|
||
| let ordered = RouterService::builder() | ||
| .middleware(BeforeAttachment) | ||
| .middleware(RequestTimingMiddleware::<1>::default()) | ||
| .get("/timed", presence) | ||
| .build(); | ||
| assert_eq!( | ||
| block_on(ordered.oneshot(request(Method::GET, "/timed"))) | ||
| .unwrap() | ||
| .body() | ||
| .as_bytes() | ||
| .unwrap(), | ||
| b"yes" | ||
| ); | ||
| } | ||
|
|
||
| #[test] | ||
| fn request_timing_next_preserves_handler_error() { | ||
| let handler = failure.into_handler(); | ||
| let ctx = RequestContext::new(request(Method::GET, "/error"), PathParams::default()); | ||
| let error = block_on( | ||
| RequestTimingMiddleware::<1>::default().handle(ctx, Next::new(&[], handler.as_ref())), | ||
| ) | ||
| .unwrap_err(); | ||
| assert_eq!(error.message(), "unchanged timing error"); | ||
| } | ||
| } | ||
|
|
||
| #[cfg(test)] | ||
| mod tests { | ||
| use super::*; | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
♻️ Cross-crate
#[path]include, and the shared scenario never runs natively: the same include appears inedgezero-adapter-fastly/tests/contract.rs:224andedgezero-adapter-spin/tests/contract.rs:547. Two concerns:edgezero-core/tests/. All three adapters already enableedgezero-core'stest-utilsfeature in dev-deps, so the scenario could live in core behind#[cfg(any(test, feature = "test-utils"))](astest_envdoes) and be called as a normal function.edgezero-core/tests/support/request_timing.rsisn't a test target, andrequest_timing_consumer.rsdoesn't include it. Socargo test -p edgezero-corenever runs this scenario, and regressions only show up in the WASM jobs. At minimum, add#[path = "support/request_timing.rs"] mod support;and a#[test]wrapper inrequest_timing_consumer.rs.