Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 9 additions & 0 deletions crates/edgezero-adapter-cloudflare/tests/contract.rs
Original file line number Diff line number Diff line change
Expand Up @@ -285,4 +285,13 @@ mod tests {
let body = response.text().await.expect("text");
assert_eq!(body, "no");
}

#[wasm_bindgen_test::wasm_bindgen_test]
fn request_timing_clock_handle_and_attachment_work_on_runtime() {
super::request_timing_tests::clock_handle_and_attachment_work_on_runtime();
}
}

#[cfg(test)]
#[path = "../../edgezero-core/tests/support/request_timing.rs"]

Copy link
Copy Markdown
Contributor

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 in edgezero-adapter-fastly/tests/contract.rs:224 and edgezero-adapter-spin/tests/contract.rs:547. Two concerns:

  1. The three adapter suites now depend on a file layout inside edgezero-core/tests/. All three adapters already enable edgezero-core's test-utils feature in dev-deps, so the scenario could live in core behind #[cfg(any(test, feature = "test-utils"))] (as test_env does) and be called as a normal function.
  2. edgezero-core/tests/support/request_timing.rs isn't a test target, and request_timing_consumer.rs doesn't include it. So cargo test -p edgezero-core never 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 in request_timing_consumer.rs.

mod request_timing_tests;
9 changes: 9 additions & 0 deletions crates/edgezero-adapter-fastly/tests/contract.rs
Original file line number Diff line number Diff line change
Expand Up @@ -213,4 +213,13 @@ mod tests {
assert_eq!(response.get_status(), FastlyStatus::OK);
assert_eq!(response.take_body_bytes(), b"hello from fastly test");
}

#[test]
fn request_timing_clock_handle_and_attachment_work_on_runtime() {
super::request_timing_tests::clock_handle_and_attachment_work_on_runtime();
}
}

#[cfg(test)]
#[path = "../../edgezero-core/tests/support/request_timing.rs"]
mod request_timing_tests;
9 changes: 9 additions & 0 deletions crates/edgezero-adapter-spin/tests/contract.rs
Original file line number Diff line number Diff line change
Expand Up @@ -536,4 +536,13 @@ mod tests {
"no secret handle yields the no-handle marker"
);
}

#[test]
fn request_timing_clock_handle_and_attachment_work_on_runtime() {
super::request_timing_tests::clock_handle_and_attachment_work_on_runtime();
}
}

#[cfg(test)]
#[path = "../../edgezero-core/tests/support/request_timing.rs"]
mod request_timing_tests;
1 change: 1 addition & 0 deletions crates/edgezero-core/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@ pub mod manifest;
pub mod middleware;
pub mod params;
pub mod proxy;
pub mod request_timing;
pub mod responder;
pub mod response;
pub mod router;
Expand Down
296 changes: 296 additions & 0 deletions crates/edgezero-core/src/middleware.rs
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;

Expand All @@ -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>;

Expand Down Expand Up @@ -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> {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🤔 edgezero.toml can't register this middleware: app! turns each [app] middleware entry into builder = builder.middleware(#path); from a syn::ExprPath (edgezero-macros/src/app.rs:126-137). The generated app has no other place to add middleware, because App has no post-build hook and the macro emits Hooks::configure itself. RequestTimingMiddleware has private fields and only a default() that can't run in a const, so a manifest entry has no value to name:

pub const TIMING: RequestTimingMiddleware<1> = RequestTimingMiddleware::default();
// error[E0015]: cannot call non-const associated function `<RequestTimingMiddleware<1> as Default>::default` in constants

A const constructor fixes it. I tested this change: registering the const the way app! does serves /timed with a handle, still excludes /health, and strict clippy (--all-features -D warnings) passes:

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
    }
}

Default can then delegate to Self::new(). The guide's "Via Manifest" section could show the pattern:

// 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 {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⛏ / 🌱 Exclusion list semantics and shape:

  1. A second with_excluded_paths call replaces the first. A probe confirmed .with_excluded_paths(&["/health"]).with_excluded_paths(&["/ready"]) leaves /health instrumented. Please say "replaces any previous list" in the rustdoc.
  2. (🌱) &'static [&'static str] rules out exclusions loaded from config (e.g. edgezero.toml) without Box::leak. The spec calls a static list sufficient for current consumers. If a config-driven need appears, impl IntoIterator<Item = impl Into<Cow<'static, str>>> stored as Arc<[Cow<'static, str>]> keeps it allocation-free for static input.

self.excluded_paths = paths;
self
}
}

#[async_trait(?Send)]
impl<const N: usize, D: Default + Send + 'static> Middleware for RequestTimingMiddleware<N, D> {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The 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 RequestTimings. So code outside the chain (an adapter entry or outer service) can only finalize, measure streaming, or call set_resp_bytes by preinstalling its own handle, at which point the middleware isn't needed. The PR scopes this out deliberately. If a later consumer needs it, inserting a clone into response.extensions_mut() on the Ok path would close the gap without changing attachment semantics.

#[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)]
Expand Down Expand Up @@ -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"));

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⛏ Assertions that can't fail:

  • Here: the request goes to /error, so next.run returns Err and oneshot renders the response after the chain has returned. A middleware that added Server-Timing to Ok responses could never trip this. Move the check onto a 200, for example in request_timing_duplicate_installation_preserves_state_and_requests_are_fresh.
  • request_timing.rs:284: both sides lock the same Arc, so the two t0 values are always equal. The record/snapshot lines after it are what prove sharing. Use Arc::ptr_eq(&timings.0, &clone.0) or delete it.
  • request_timing.rs:334: nothing can change t0 after construction, so this always holds. The [Some(11ms)] assertion on line 337 is what proves the handler saw the preinstalled handle.
  • request_timing.rs:558: >= before holds even if the cancelled span recorded nothing. The fresh block below it is what proves cancellation records.
  • request_timing_consumer.rs:63: AppTimings isn't Clone, so it can never be inserted into extensions, and this lookup can only return None.

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::*;
Expand Down
Loading
Loading