From 9c03cc59300363ae5339fc970dded08ede563e98 Mon Sep 17 00:00:00 2001 From: Christian Date: Mon, 28 Sep 2026 15:39:26 -0500 Subject: [PATCH] feat(core): add shared request timing collector and middleware --- .../tests/contract.rs | 9 + .../edgezero-adapter-fastly/tests/contract.rs | 9 + .../edgezero-adapter-spin/tests/contract.rs | 9 + crates/edgezero-core/src/lib.rs | 1 + crates/edgezero-core/src/middleware.rs | 296 +++++++++ crates/edgezero-core/src/request_timing.rs | 571 ++++++++++++++++++ .../tests/request_timing_consumer.rs | 107 ++++ .../tests/support/request_timing.rs | 71 +++ docs/guide/middleware.md | 111 +++- .../plans/2026-09-28-request-timing.md | 362 +++++++++++ .../specs/2026-09-28-request-timing-design.md | 449 ++++++++++++++ 11 files changed, 1973 insertions(+), 22 deletions(-) create mode 100644 crates/edgezero-core/src/request_timing.rs create mode 100644 crates/edgezero-core/tests/request_timing_consumer.rs create mode 100644 crates/edgezero-core/tests/support/request_timing.rs create mode 100644 docs/superpowers/plans/2026-09-28-request-timing.md create mode 100644 docs/superpowers/specs/2026-09-28-request-timing-design.md diff --git a/crates/edgezero-adapter-cloudflare/tests/contract.rs b/crates/edgezero-adapter-cloudflare/tests/contract.rs index 249ccefb..bfe6b8f2 100644 --- a/crates/edgezero-adapter-cloudflare/tests/contract.rs +++ b/crates/edgezero-adapter-cloudflare/tests/contract.rs @@ -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"] +mod request_timing_tests; diff --git a/crates/edgezero-adapter-fastly/tests/contract.rs b/crates/edgezero-adapter-fastly/tests/contract.rs index fb953b14..391e8a6b 100644 --- a/crates/edgezero-adapter-fastly/tests/contract.rs +++ b/crates/edgezero-adapter-fastly/tests/contract.rs @@ -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; diff --git a/crates/edgezero-adapter-spin/tests/contract.rs b/crates/edgezero-adapter-spin/tests/contract.rs index 7e4af0ee..2b1e2798 100644 --- a/crates/edgezero-adapter-spin/tests/contract.rs +++ b/crates/edgezero-adapter-spin/tests/contract.rs @@ -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; diff --git a/crates/edgezero-core/src/lib.rs b/crates/edgezero-core/src/lib.rs index 66aa8545..2a53d87e 100644 --- a/crates/edgezero-core/src/lib.rs +++ b/crates/edgezero-core/src/lib.rs @@ -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; diff --git a/crates/edgezero-core/src/middleware.rs b/crates/edgezero-core/src/middleware.rs index a9edaf2c..198e9e22 100644 --- a/crates/edgezero-core/src/middleware.rs +++ b/crates/edgezero-core/src/middleware.rs @@ -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; @@ -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 { + data: PhantomData D>, + excluded_paths: &'static [&'static str], +} + +impl Default for RequestTimingMiddleware { + #[inline] + fn default() -> Self { + Self { + data: PhantomData, + excluded_paths: &[], + } + } +} + +impl RequestTimingMiddleware { + /// 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 { + self.excluded_paths = paths; + self + } +} + +#[async_trait(?Send)] +impl Middleware for RequestTimingMiddleware { + #[inline] + async fn handle(&self, mut ctx: RequestContext, next: Next<'_>) -> Result { + if ctx + .request() + .extensions() + .get::>() + .is_none() + && !self.excluded_paths.contains(&ctx.request().uri().path()) + { + ctx.request_mut() + .extensions_mut() + .insert(RequestTimings::::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 { + assert!(ctx.request().extensions().get::().is_none()); + next.run(ctx).await + } + } + + struct Capture { + handles: Arc>>, + short_circuit: bool, + } + + #[async_trait(?Send)] + impl Middleware for Capture { + async fn handle(&self, ctx: RequestContext, next: Next<'_>) -> Result { + let timings = ctx.request().extensions().get::().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 { + let installed = ctx.request().extensions().get::(); + 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 { + 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::::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::::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")); + 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::::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::*; diff --git a/crates/edgezero-core/src/request_timing.rs b/crates/edgezero-core/src/request_timing.rs new file mode 100644 index 00000000..361c44b6 --- /dev/null +++ b/crates/edgezero-core/src/request_timing.rs @@ -0,0 +1,571 @@ +//! Bounded, best-effort per-request timing with application-owned data. +//! +//! All clones share one origin and mutex. Contention drops an entire operation; +//! it never blocks or substitutes zero for unavailable facts. Callbacks run under +//! the lock and must be short, synchronous and non-reentrant. Panics propagate, +//! without rollback: callbacks must leave their data valid even on unwind. +//! Poisoned locks recover best-effort, not with validation of application state. +//! Collection does not render headers or implicitly mark request completion. + +use std::sync::{Arc, Mutex, MutexGuard, TryLockError}; +use std::time::Duration; + +use web_time::Instant; + +/// Why a timing operation could not be performed. +#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)] +pub enum TimingError { + /// The slot is outside the collector's fixed array; nothing was changed. + #[error("timing slot is out of bounds")] + InvalidSlot, + /// Another operation holds the lock; no callback was invoked or state changed. + #[error("request timings are currently unavailable")] + Unavailable, +} + +struct Inner { + data: D, + headers_ready_total: Option, + phases: [Option; N], + request_elapsed: Option, + resp_bytes: Option, + t0: Instant, +} + +/// Consistent generic fields and borrowed application data, projected under lock. +/// +/// The callback may return owned values, but cannot retain `data` beyond the call. +/// Durations are unrounded; output conversion and exposure are application policy. +/// +/// ```compile_fail +/// use edgezero_core::request_timing::RequestTimings; +/// let timings = RequestTimings::<1, String>::new(); +/// let escaped = timings.snapshot(|view| view.data).unwrap(); +/// ``` +pub struct TimingSnapshot<'data, const N: usize, D> { + /// Application facts protected by the same lock as the generic fields. + pub data: &'data D, + /// Elapsed from the shared origin when this snapshot was taken. + pub elapsed: Duration, + /// First successful explicit headers-ready mark. + pub headers_ready_total: Option, + /// Accumulated durations; absent differs from a recorded zero. + pub phases: [Option; N], + /// First successful explicit request-completion mark. + pub request_elapsed: Option, + /// Last successfully recorded response byte count. + pub resp_bytes: Option, +} + +/// A cheap shared handle; neither cloning nor projection requires `D: Clone`. +/// +/// Create one per request at the chosen measurement boundary, not as router state. +/// Installed extension handles require `D: Send + 'static`, not `D: Sync`. +/// T0 is collector creation, not platform ingress or browser navigation start. +pub struct RequestTimings(Arc>>); + +impl Clone for RequestTimings { + #[inline] + fn clone(&self) -> Self { + Self(Arc::clone(&self.0)) + } + + #[inline] + fn clone_from(&mut self, source: &Self) { + self.0.clone_from(&source.0); + } +} + +impl Default for RequestTimings { + #[inline] + fn default() -> Self { + Self::with_data(D::default()) + } +} + +impl RequestTimings { + /// Starts a request clock with default application data. + #[must_use] + #[inline] + pub fn new() -> Self { + Self::default() + } +} + +impl RequestTimings { + /// Returns elapsed time from the immutable shared origin. + /// # Errors + /// Returns [`TimingError::Unavailable`] on contention. + #[inline] + pub fn elapsed(&self) -> Result { + Ok(self.try_inner()?.t0.elapsed()) + } + + /// Marks headers ready, first successful write wins (later calls succeed unchanged). + /// # Errors + /// Returns [`TimingError::Unavailable`] on contention. + #[inline] + pub fn mark_headers_ready(&self) -> Result<(), TimingError> { + let mut inner = self.try_inner()?; + if inner.headers_ready_total.is_none() { + inner.headers_ready_total = Some(inner.t0.elapsed()); + } + Ok(()) + } + + /// Marks request completion explicitly, first successful write wins. + /// Neither span drop nor middleware return calls this automatically. + /// # Errors + /// Returns [`TimingError::Unavailable`] on contention. + #[inline] + pub fn mark_request_elapsed(&self) -> Result<(), TimingError> { + let mut inner = self.try_inner()?; + if inner.request_elapsed.is_none() { + inner.request_elapsed = Some(inner.t0.elapsed()); + } + Ok(()) + } + + /// Saturating-adds a duration to a slot, preserving recorded zero. + /// # Errors + /// Returns [`TimingError::InvalidSlot`] before any mutation for invalid slots, + /// or [`TimingError::Unavailable`] on contention. + #[inline] + pub fn record(&self, slot: usize, duration: Duration) -> Result<(), TimingError> { + self.record_with(slot, duration, |_, _| ()) + } + + /// Accumulates a phase and updates application facts under one lock. + /// + /// The callback receives elapsed time from T0. It must be short, synchronous, + /// non-reentrant and leave data valid on unwind. A panic propagates and can + /// leave both the phase and application data partially changed; no rollback + /// is attempted. Poison recovery follows the module's best-effort contract. + /// # Errors + /// Invalid slots return [`TimingError::InvalidSlot`] before the callback or + /// any mutation. Contention returns [`TimingError::Unavailable`] likewise. + #[inline] + pub fn record_with R>( + &self, + slot: usize, + duration: Duration, + update: F, + ) -> Result { + Self::validate_slot(slot)?; + let mut inner = self.try_inner()?; + // Bounds were checked before acquiring the lock or invoking user code. + let phase = inner.phases.get_mut(slot).ok_or(TimingError::InvalidSlot)?; + *phase = Some(phase.unwrap_or(Duration::ZERO).saturating_add(duration)); + let elapsed = inner.t0.elapsed(); + Ok(update(&mut inner.data, elapsed)) + } + + /// Records the response byte count, last successful write wins. + /// # Errors + /// Returns [`TimingError::Unavailable`] on contention. + #[inline] + pub fn set_resp_bytes(&self, bytes: u64) -> Result<(), TimingError> { + self.try_inner()?.resp_bytes = Some(bytes); + Ok(()) + } + + /// Projects generic and application facts together, without cloning the payload. + /// The callback runs under lock and must be short, synchronous and non-reentrant. + /// Callback panics propagate; subsequent operations recover poison best-effort. + /// # Errors + /// Returns [`TimingError::Unavailable`] without invoking the callback on contention. + #[inline] + pub fn snapshot) -> R>( + &self, + project: F, + ) -> Result { + let inner = self.try_inner()?; + Ok(project(TimingSnapshot { + phases: inner.phases, + elapsed: inner.t0.elapsed(), + headers_ready_total: inner.headers_ready_total, + request_elapsed: inner.request_elapsed, + resp_bytes: inner.resp_bytes, + data: &inner.data, + })) + } + + /// Starts a phase span. Dropping it attempts to record work so far, including + /// cancellation, but never marks completion. Panic-abort may skip destructors. + /// # Errors + /// Returns [`TimingError::InvalidSlot`] for an out-of-bounds slot. + #[inline] + pub fn span(&self, slot: usize) -> Result, TimingError> { + Self::validate_slot(slot)?; + Ok(PhaseSpan { + timings: self.clone(), + slot, + started: Instant::now(), + }) + } + + fn try_inner(&self) -> Result>, TimingError> { + match self.0.try_lock() { + Ok(guard) => Ok(guard), + Err(TryLockError::Poisoned(poisoned)) => Ok(poisoned.into_inner()), + Err(TryLockError::WouldBlock) => Err(TimingError::Unavailable), + } + } + + /// Updates application facts using the shared origin, under the collector lock. + /// The callback has the same panic and reentrancy contract as [`Self::record_with`]. + /// # Errors + /// Returns [`TimingError::Unavailable`] without invoking the callback on contention. + #[inline] + pub fn update_data R>( + &self, + update: F, + ) -> Result { + let mut inner = self.try_inner()?; + let elapsed = inner.t0.elapsed(); + Ok(update(&mut inner.data, elapsed)) + } + + fn validate_slot(slot: usize) -> Result<(), TimingError> { + if slot < N { + Ok(()) + } else { + Err(TimingError::InvalidSlot) + } + } + + /// Starts a request clock with application data, without a `Default` bound. + #[must_use] + #[inline] + pub fn with_data(data: D) -> Self { + Self(Arc::new(Mutex::new(Inner { + t0: Instant::now(), + phases: [None; N], + headers_ready_total: None, + request_elapsed: None, + resp_bytes: None, + data, + }))) + } +} + +/// Records elapsed phase work on drop, discarding a contended sample. +#[must_use = "dropping a span records its elapsed duration"] +pub struct PhaseSpan { + slot: usize, + started: Instant, + timings: RequestTimings, +} + +impl Drop for PhaseSpan { + #[inline] + fn drop(&mut self) { + match self.timings.record(self.slot, self.started.elapsed()) { + Ok(()) | Err(TimingError::Unavailable | TimingError::InvalidSlot) => {} + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn aged(data: D) -> RequestTimings { + let timings = RequestTimings::with_data(data); + timings.0.lock().unwrap().t0 = Instant::now().checked_sub(Duration::from_secs(10)).unwrap(); + timings + } + + #[test] + fn clone_shares_aged_origin_and_facts_but_requests_are_independent() { + let timings = aged::<2, _>(()); + let clone = timings.clone(); + let origin = timings.0.lock().unwrap().t0; + assert_eq!(origin, clone.0.lock().unwrap().t0); + assert!(clone.elapsed().unwrap() >= Duration::from_secs(10)); + clone.record(0, Duration::from_millis(3)).unwrap(); + assert_eq!( + timings.snapshot(|view| view.phases).unwrap(), + [Some(Duration::from_millis(3)), None] + ); + let separate = RequestTimings::<2>::new(); + assert!(!Arc::ptr_eq(&timings.0, &separate.0)); + assert_eq!(separate.snapshot(|view| view.phases).unwrap(), [None; 2]); + } + + #[test] + fn middleware_and_handler_preserve_preaged_origin() { + use crate::body::Body; + use crate::context::RequestContext; + use crate::error::EdgeError; + use crate::http::{Response, StatusCode, request_builder}; + use crate::middleware::RequestTimingMiddleware; + use crate::response::response_with_body; + use crate::router::RouterService; + use futures::executor::block_on; + + #[crate::action] + async fn timed(ctx: RequestContext) -> Result { + let handle = ctx + .request() + .extensions() + .get::>() + .unwrap(); + assert!(handle.elapsed().unwrap() >= Duration::from_secs(10)); + handle.record(0, Duration::from_millis(3)).unwrap(); + handle.mark_headers_ready().unwrap(); + response_with_body(StatusCode::OK, Body::empty()) + } + + let timings = aged::<1, _>(()); + let origin = timings.0.lock().unwrap().t0; + timings.record(0, Duration::from_millis(8)).unwrap(); + let mut request = request_builder().uri("/timed").body(Body::empty()).unwrap(); + request.extensions_mut().insert(timings.clone()); + let router = RouterService::builder() + .middleware(RequestTimingMiddleware::<1>::default()) + .middleware(RequestTimingMiddleware::<1>::default()) + .get("/timed", timed) + .build(); + assert_eq!( + block_on(router.oneshot(request)).unwrap().status(), + StatusCode::OK + ); + assert_eq!(timings.0.lock().unwrap().t0, origin); + timings + .snapshot(|view| { + assert_eq!(view.phases, [Some(Duration::from_millis(11))]); + assert!(view.headers_ready_total.unwrap() >= Duration::from_secs(10)); + assert_eq!(view.request_elapsed, None); + }) + .unwrap(); + } + + #[test] + fn data_does_not_need_default_or_clone() { + struct Data(u8); + let timings = RequestTimings::<0, _>::with_data(Data(7)); + let clone = timings.clone(); + assert_eq!(clone.snapshot(|view| view.data.0).unwrap(), 7); + } + + #[test] + fn phases_preserve_zero_accumulate_saturate_and_reject_invalid_slots() { + let timings = RequestTimings::<2>::new(); + timings.record(0, Duration::ZERO).unwrap(); + assert_eq!( + timings.snapshot(|view| view.phases).unwrap(), + [Some(Duration::ZERO), None] + ); + timings.record(0, Duration::from_millis(2)).unwrap(); + timings.record(0, Duration::from_millis(3)).unwrap(); + assert_eq!( + timings.snapshot(|view| view.phases[0]).unwrap(), + Some(Duration::from_millis(5)) + ); + timings.record(0, Duration::MAX).unwrap(); + assert_eq!( + timings.snapshot(|view| view.phases[0]).unwrap(), + Some(Duration::MAX) + ); + assert_eq!( + timings.record(2, Duration::MAX), + Err(TimingError::InvalidSlot) + ); + assert_eq!( + timings.record(usize::MAX, Duration::ZERO), + Err(TimingError::InvalidSlot) + ); + assert!(matches!(timings.span(2), Err(TimingError::InvalidSlot))); + assert_eq!( + timings.snapshot(|view| view.phases).unwrap(), + [Some(Duration::MAX), None] + ); + assert_eq!( + RequestTimings::<0>::new().record(0, Duration::ZERO), + Err(TimingError::InvalidSlot) + ); + } + + #[test] + fn lifecycle_marks_are_explicit_first_write_and_bytes_are_last_write() { + let timings = aged::<1, _>(()); + timings + .snapshot(|view| { + assert_eq!(view.headers_ready_total, None); + assert_eq!(view.request_elapsed, None); + assert_eq!(view.resp_bytes, None); + }) + .unwrap(); + timings.mark_headers_ready().unwrap(); + timings.mark_request_elapsed().unwrap(); + let first = timings + .snapshot(|view| (view.headers_ready_total, view.request_elapsed)) + .unwrap(); + assert!(first.0.unwrap() >= Duration::from_secs(10)); + assert!(first.1.unwrap() >= first.0.unwrap()); + timings.mark_headers_ready().unwrap(); + timings.mark_request_elapsed().unwrap(); + assert_eq!( + timings + .snapshot(|view| (view.headers_ready_total, view.request_elapsed)) + .unwrap(), + first + ); + timings.set_resp_bytes(42).unwrap(); + timings.set_resp_bytes(0).unwrap(); + assert_eq!(timings.snapshot(|view| view.resp_bytes).unwrap(), Some(0)); + } + + #[test] + fn compound_updates_and_projection_share_one_lock_and_origin() { + let timings = aged::<1, _>((None, None)); + let returned = timings + .record_with(0, Duration::from_millis(9), |data, elapsed| { + data.0 = Some("before headers"); + data.1 = Some(elapsed); + 42_u32 + }) + .unwrap(); + assert_eq!(returned, 42); + timings + .snapshot(|view| { + assert_eq!(view.phases, [Some(Duration::from_millis(9))]); + assert_eq!(view.data.0, Some("before headers")); + assert!(view.data.1.unwrap() >= Duration::from_secs(10)); + assert!(view.elapsed >= view.data.1.unwrap()); + assert_eq!(timings.elapsed(), Err(TimingError::Unavailable)); + }) + .unwrap(); + timings + .update_data(|data, _| data.0 = Some("streaming")) + .unwrap(); + assert_eq!( + timings.snapshot(|view| view.data.0).unwrap(), + Some("streaming") + ); + assert_eq!( + timings.record_with(1, Duration::MAX, |_, _| panic!("invalid callback")), + Err(TimingError::InvalidSlot) + ); + assert_eq!( + timings.snapshot(|view| (view.phases, view.data.0)).unwrap(), + ([Some(Duration::from_millis(9))], Some("streaming")) + ); + } + + #[test] + fn contention_drops_whole_operations_without_invoking_callbacks() { + let timings = RequestTimings::<1, u32>::new(); + let guard = timings.0.lock().unwrap(); + assert_eq!( + timings.record_with(0, Duration::MAX, |_, _| panic!("contended update")), + Err(TimingError::Unavailable) + ); + assert_eq!( + timings.update_data(|_, _| panic!("contended payload")), + Err(TimingError::Unavailable) + ); + assert_eq!( + timings.snapshot(|_| panic!("contended snapshot")), + Err(TimingError::Unavailable) + ); + assert_eq!(timings.elapsed(), Err(TimingError::Unavailable)); + assert_eq!(timings.mark_headers_ready(), Err(TimingError::Unavailable)); + assert_eq!( + timings.mark_request_elapsed(), + Err(TimingError::Unavailable) + ); + assert_eq!(timings.set_resp_bytes(8), Err(TimingError::Unavailable)); + assert_eq!( + timings.record_with(1, Duration::MAX, |_, _| panic!("invalid update")), + Err(TimingError::InvalidSlot) + ); + drop(timings.span(0).unwrap()); + drop(guard); + timings + .snapshot(|view| { + assert_eq!(view.phases, [None]); + assert_eq!(*view.data, 0); + assert_eq!(view.headers_ready_total, None); + assert_eq!(view.request_elapsed, None); + assert_eq!(view.resp_bytes, None); + }) + .unwrap(); + timings.mark_headers_ready().unwrap(); + assert!( + timings + .snapshot(|view| view.headers_ready_total) + .unwrap() + .is_some() + ); + } + + #[cfg(not(target_arch = "wasm32"))] + #[test] + fn callback_panic_propagates_leaves_partial_state_and_recovers_poison() { + use std::panic::catch_unwind; + let timings = RequestTimings::<1, u32>::new(); + let outcome = catch_unwind(|| { + timings + .record_with(0, Duration::from_millis(4), |data, _| { + *data = 7; + panic!("callback interrupted"); + }) + .unwrap(); + }); + assert!(outcome.is_err()); + assert!(timings.0.is_poisoned()); + assert_eq!( + timings.snapshot(|view| (view.phases, *view.data)).unwrap(), + ([Some(Duration::from_millis(4))], 7) + ); + timings + .record_with(0, Duration::from_millis(2), |data, _| { + *data = data.saturating_add(1); + }) + .unwrap(); + assert_eq!( + timings.snapshot(|view| (view.phases, *view.data)).unwrap(), + ([Some(Duration::from_millis(6))], 8) + ); + timings.mark_request_elapsed().unwrap(); + assert!( + timings + .snapshot(|view| view.request_elapsed) + .unwrap() + .is_some() + ); + } + + #[test] + fn spans_record_drop_and_cancellation_without_completing_request() { + use futures::{FutureExt as _, future::pending}; + let timings = RequestTimings::<1>::new(); + let mut span = timings.span(0).unwrap(); + span.started = Instant::now().checked_sub(Duration::from_secs(2)).unwrap(); + drop(span); + let before = timings.snapshot(|view| view.phases[0]).unwrap().unwrap(); + assert!(before >= Duration::from_secs(2)); + let pending_work = async { + let _span = timings.span(0).unwrap(); + pending::<()>().await; + }; + // Poll once, then drop the pending future, as cancellation would. + assert!(pending_work.now_or_never().is_none()); + timings + .snapshot(|view| { + assert!(view.phases[0].unwrap() >= before); + assert_eq!(view.headers_ready_total, None); + assert_eq!(view.request_elapsed, None); + }) + .unwrap(); + let fresh = RequestTimings::<1>::new(); + let fresh_work = async { + let _span = fresh.span(0).unwrap(); + pending::<()>().await; + }; + assert!(fresh_work.now_or_never().is_none()); + assert!(fresh.snapshot(|view| view.phases[0]).unwrap().is_some()); + } +} diff --git a/crates/edgezero-core/tests/request_timing_consumer.rs b/crates/edgezero-core/tests/request_timing_consumer.rs new file mode 100644 index 00000000..9adeb339 --- /dev/null +++ b/crates/edgezero-core/tests/request_timing_consumer.rs @@ -0,0 +1,107 @@ +//! Public API fixture: application types stay outside the generic collector. +#[cfg(test)] +mod tests { + use std::cell::Cell; + use std::time::Duration; + + use edgezero_core::body::Body; + use edgezero_core::context::RequestContext; + use edgezero_core::error::EdgeError; + use edgezero_core::http::{Response, StatusCode, request_builder}; + use edgezero_core::middleware::RequestTimingMiddleware; + use edgezero_core::request_timing::{RequestTimings, TimingError}; + use edgezero_core::response::response_with_body; + use edgezero_core::router::RouterService; + use futures::executor::block_on; + + // Deliberately neither Clone nor Sync: only Send is needed behind the mutex. + #[derive(Default)] + struct AppData { + attempts: Cell, + started: Option, + } + + type Handle = RequestTimings<1, AppData>; + + enum Phase { + Fetch, + } + + impl Phase { + fn slot(self) -> usize { + match self { + Self::Fetch => 0, + } + } + } + + struct AppTimings(Handle); + + impl AppTimings { + fn from_context(ctx: &RequestContext) -> Option { + ctx.request() + .extensions() + .get::() + .cloned() + .map(Self) + } + + fn record(&self, phase: Phase, duration: Duration) -> Result<(), TimingError> { + self.0.record_with(phase.slot(), duration, |data, elapsed| { + data.attempts.set(data.attempts.get().saturating_add(1)); + data.started.get_or_insert(elapsed); + }) + } + } + + #[edgezero_core::action] + async fn timed_handler(ctx: RequestContext) -> Result { + let timings = AppTimings::from_context(&ctx).expect("installed generic handle"); + timings + .record(Phase::Fetch, Duration::from_millis(7)) + .expect("record"); + assert!(ctx.request().extensions().get::().is_none()); + response_with_body(StatusCode::OK, Body::empty()) + } + + #[test] + fn request_timing_public_facade_shares_extension_state_and_origin() { + let handle = Handle::new(); + let before = handle.elapsed().expect("elapsed"); + let mut request = request_builder() + .uri("/timed") + .body(Body::empty()) + .expect("request"); + request.extensions_mut().insert(handle.clone()); + let router = RouterService::builder() + .middleware(RequestTimingMiddleware::<1, AppData>::default()) + .get("/timed", timed_handler) + .build(); + let response = block_on(router.oneshot(request)).expect("response"); + assert_eq!(response.status(), StatusCode::OK); + handle + .snapshot(|view| { + assert_eq!(view.phases, [Some(Duration::from_millis(7))]); + assert_eq!(view.data.attempts.get(), 1); + assert!(view.data.started.expect("application mark") >= before); + assert!(view.elapsed >= view.data.started.expect("application mark")); + assert_eq!(view.headers_ready_total, None); + assert_eq!(view.request_elapsed, None); + }) + .expect("snapshot"); + let fresh = request_builder() + .uri("/timed") + .body(Body::empty()) + .expect("request"); + assert_eq!( + block_on(router.oneshot(fresh)).expect("response").status(), + StatusCode::OK + ); + assert_eq!( + handle + .snapshot(|view| view.data.attempts.get()) + .expect("snapshot"), + 1 + ); + } +} diff --git a/crates/edgezero-core/tests/support/request_timing.rs b/crates/edgezero-core/tests/support/request_timing.rs new file mode 100644 index 00000000..04d6c0f1 --- /dev/null +++ b/crates/edgezero-core/tests/support/request_timing.rs @@ -0,0 +1,71 @@ +//! Runtime-neutral timing scenario shared by the adapter WASM contract suites. + +use edgezero_core::body::Body; +use edgezero_core::context::RequestContext; +use edgezero_core::error::EdgeError; +use edgezero_core::http::{Response, StatusCode, request_builder}; +use edgezero_core::middleware::RequestTimingMiddleware; +use edgezero_core::request_timing::RequestTimings; +use edgezero_core::response::response_with_body; +use edgezero_core::router::RouterService; +use futures::executor::block_on; +use std::time::Duration; + +type Handle = RequestTimings<1, Option>; + +#[edgezero_core::action] +async fn timed(ctx: RequestContext) -> Result { + let handle = ctx.request().extensions().get::().unwrap(); + let span = handle.span(0).unwrap(); + handle + .update_data(|data, elapsed| *data = Some(elapsed)) + .unwrap(); + drop(span); + handle + .snapshot(|view| { + assert!(view.phases[0].is_some()); + assert!(view.elapsed >= view.data.unwrap()); + assert_eq!(view.headers_ready_total, None); + assert_eq!(view.request_elapsed, None); + }) + .unwrap(); + response_with_body(StatusCode::OK, Body::empty()) +} + +pub(crate) fn clock_handle_and_attachment_work_on_runtime() { + let handle = Handle::new(); + let before = handle.elapsed().unwrap(); + let mut request = request_builder().uri("/timed").body(Body::empty()).unwrap(); + request.extensions_mut().insert(handle.clone()); + let router = RouterService::builder() + .middleware(RequestTimingMiddleware::<1, Option>::default()) + .get("/timed", timed) + .build(); + assert_eq!( + block_on(router.oneshot(request)).unwrap().status(), + StatusCode::OK + ); + handle + .snapshot(|view| { + assert!(view.phases[0].is_some()); + assert!(view.data.unwrap() >= before); + assert!(view.elapsed >= view.data.unwrap()); + assert_eq!(view.headers_ready_total, None); + assert_eq!(view.request_elapsed, None); + }) + .unwrap(); + handle.mark_headers_ready().unwrap(); + assert!( + handle + .snapshot(|view| view.headers_ready_total.unwrap() >= before) + .unwrap() + ); + + // The normal attachment path must create a collector on this runtime too. + let uninstrumented = request_builder().uri("/timed").body(Body::empty()).unwrap(); + assert!(uninstrumented.extensions().get::().is_none()); + assert_eq!( + block_on(router.oneshot(uninstrumented)).unwrap().status(), + StatusCode::OK + ); +} diff --git a/docs/guide/middleware.md b/docs/guide/middleware.md index efefe3d3..7b099ec0 100644 --- a/docs/guide/middleware.md +++ b/docs/guide/middleware.md @@ -24,7 +24,7 @@ impl Middleware for RequestLogger { ) -> Result { let method = ctx.request().method().clone(); let path = ctx.request().uri().path().to_string(); - let start = std::time::Instant::now(); + let start = web_time::Instant::now(); let response = next.run(ctx).await?; @@ -58,7 +58,8 @@ middleware = [ ] ``` -Middleware are applied in order before routes are matched. +Routes are matched first. Middleware then run in registration order for a matched route; +unmatched 404/405 requests bypass the middleware chain. ### Programmatically @@ -164,31 +165,96 @@ impl Middleware for CorsMiddleware { ### Request Timing +`RequestTimingMiddleware` attaches a fresh `RequestTimings` only when +that exact type is absent. All clones share one portable clock and one mutex, +including the application-owned payload. Register it before timing consumers. +The default excludes no paths; static exclusions compare exact URI paths for every +method (including `/health?check=1`), without removing preinstalled handles. + +Applications own their phase enum and map it to bounded slots through a facade: + ```rust -pub struct TimingMiddleware; +use std::time::Duration; +use edgezero_core::context::RequestContext; +use edgezero_core::middleware::RequestTimingMiddleware; +use edgezero_core::request_timing::{RequestTimings, TimingError}; +use edgezero_core::router::RouterService; -#[async_trait(?Send)] -impl Middleware for TimingMiddleware { - async fn handle( - &self, - ctx: RequestContext, - next: Next<'_>, - ) -> Result { - let start = std::time::Instant::now(); +#[derive(Default)] +struct AppData { + attempts: u32, + started: Option, +} - let mut response = next.run(ctx).await?; +type Handle = RequestTimings<1, AppData>; - let duration = start.elapsed(); - response.headers_mut().insert( - "x-response-time", - format!("{}ms", duration.as_millis()).parse().unwrap(), - ); +enum Phase { + Fetch, +} - Ok(response) +impl Phase { + fn slot(self) -> usize { + match self { + Self::Fetch => 0, + } } } + +struct AppTimings(Handle); + +impl AppTimings { + fn from_context(ctx: &RequestContext) -> Option { + ctx.request().extensions().get::().cloned().map(Self) + } + + fn record(&self, phase: Phase, duration: Duration) -> Result<(), TimingError> { + self.0.record_with(phase.slot(), duration, |data, elapsed| { + data.attempts = data.attempts.saturating_add(1); + data.started.get_or_insert(elapsed); + }) + } +} + +let router = RouterService::builder() + .middleware(RequestTimingMiddleware::<1, AppData>::default() + .with_excluded_paths(&["/health"])) + .build(); ``` +Register application routes on that builder. The facade reads `Handle` from +extensions; it is **not** a second installed extension or clock. Do not use +`RouterBuilder::with_state` for collectors: router state is shared across requests +and overwrites same-type request extensions. The external `request_timing_consumer` +test compiles this pattern through a real handler, including a non-`Clone`, +non-`Sync` payload. Middleware payloads need only `Default + Send + 'static`; +manual `with_data(data)` construction needs no `Default`. + +- `record(slot, duration)` saturating-adds; absent and recorded zero stay distinct. + `span(slot)?` records work on drop, including cancelled work, but not completion. + Destructor execution is not guaranteed in panic-abort environments. +- `record_with` updates phase and payload together; `update_data` updates payload + alone using elapsed time from the same origin. `snapshot(|view| ...)` projects + generic fields and borrowed payload under that same lock, without cloning it. + Borrowed payload cannot escape the projection. +- Invalid slots return `TimingError::InvalidSlot` before mutation or callbacks. + Contention returns `TimingError::Unavailable`, dropping the whole operation + without invoking its callback. There is no blocking or retry. Span drop discards + an unavailable sample. Missing facts must not be interpreted as zero. +- Callbacks must be short, synchronous and non-reentrant: never call a lock-taking + collector operation inside one. Panics propagate and may leave partial state; + poisoned locks recover best-effort, **without rollback or payload validation**. + Application callbacks must keep data valid even on unwind. +- Call `mark_headers_ready()` and `mark_request_elapsed()` at explicit application + boundaries (first successful write wins); `set_resp_bytes(bytes)` is last-write. + T0 is collector creation at the chosen boundary, not socket ingress or browser + navigation. Wrapping or reattaching a handle never resets that origin. + +Attachment does not finalize a request, consume a streamed body, render +`Server-Timing`, or change cache/privacy policy. Header rendering and exposure +remain application-owned. Unmatched routes bypass middleware, and +`RouterService::oneshot` renders handler errors outside the chain. Use adapter or +outer-service hooks for terminal/error and streaming measurements as needed. + ## Early Returns Middleware can short-circuit the chain by not calling `next`: @@ -220,10 +286,11 @@ impl Middleware for RateLimiter { EdgeZero provides these middleware out of the box: -| Middleware | Purpose | -| --------------- | -------------------------------------------------------------------------------------------------------------------------------------------------------------------- | -| `RequestLogger` | Logs request method, path, and response status | -| `FnMiddleware` | Wraps a closure `Fn(RequestContext, Next<'_>) -> impl Future>`; build one with `middleware_fn(...)` or `FnMiddleware::new(...)` | +| Middleware | Purpose | +| ------------------------------- | -------------------------------------------------------------------------------------------------------------------------------------------------------------------- | +| `RequestTimingMiddleware` | Attaches a per-request timing collector if absent; configurable exact-path exclusions | +| `RequestLogger` | Logs request method, path, and response status | +| `FnMiddleware` | Wraps a closure `Fn(RequestContext, Next<'_>) -> impl Future>`; build one with `middleware_fn(...)` or `FnMiddleware::new(...)` | If you already hold an `Arc` (`BoxMiddleware`), register it with `.middleware_arc(...)` instead of `.middleware(...)`. diff --git a/docs/superpowers/plans/2026-09-28-request-timing.md b/docs/superpowers/plans/2026-09-28-request-timing.md new file mode 100644 index 00000000..b6b450e7 --- /dev/null +++ b/docs/superpowers/plans/2026-09-28-request-timing.md @@ -0,0 +1,362 @@ +# Shared request timing implementation plan + +Date: 2026-09-28 +Status: EdgeZero tasks 1–3 implemented, locally verified and independently reviewed; +draft-PR publication authorized. TS migration and joint validation remain pending. +Spec: [Shared request timing collector and middleware](../specs/2026-09-28-request-timing-design.md). + +## Outcome and approved scope + +Move generic timing collection and insert-if-absent middleware into EdgeZero. +Adopt it in Trusted Server without changing diagnostic wire output, rendering +policy, or adapter measurement boundaries. Ship the two changes together after +joint validation; do not land a temporary TS-only middleware extraction. + +The user approved typed application phase enums over bounded generic slots, +application-owned payload under the same clock and mutex, one concrete extension +handle, and TS-owned rendering. Exact public signatures and the callback failure +contract are the first checkpoint, not permission to redesign those decisions. + +Implementation of both sides and later draft-PR publication are now authorized. +Independent EdgeZero reviews cleared. This publication stage authorizes an +EdgeZero-only focused commit, normal push and one draft PR targeting main. +Merge, release/tag and final delivery still require separate approval. + +## Repositories and starting state + +| Repository | Worktree / branch | Reference | +| --- | --- | --- | +| EdgeZero | `~/worktrees/feat-shared-request-timing`, `feat/shared-request-timing` | `98930917d96cb665c7255f36ca8bbd61fc7539e1` | +| Trusted Server | `~/worktrees/feature-ts-console-improvements`, `feature/ts-console-improvements` | `f951955f537b89392dcb863fca7a01a5f6fa4846` | + +EdgeZero's main checkout was fast-forwarded to fetched `origin/main` before the +worktree was created. The spec moved into this worktree; the original checkout +is clean. Trusted Server currently pins six EdgeZero dependencies to `v0.0.7`. +Recheck both heads and worktrees before implementation rather than assuming they +have stayed unchanged. Keep one writer per worktree. + +## Delivery dependencies + +```mermaid +flowchart TD + A[Confirm API contract and baseline] --> B[Generic collector and consumer fixture] + B --> C[Shared attachment middleware] + C --> D[TS integration pinned to exact pushed candidate commit] + D --> E[Joint native and WASM validation] + E --> F[Approved EdgeZero merge and tag] + F --> G[Final TS pins without overrides] + G --> H[Revalidate and deliver TS] +``` + +A usable EdgeZero tag must precede the final TS dependency pin. Joint delivery +means both implementations are verified before releasing EdgeZero, then the +final tag-backed TS change is verified before delivery. It does not mean atomic +publication across repositories. + +## Task 1: Pin API details and capture the baseline + +Files: spec, this plan, existing EdgeZero core modules and CI definitions, +TS `request_timing.rs`, adapter timing integrations, and dependency manifests. +No production code changes in this task. + +- [x] Check repository instructions, status, exact heads, installed toolchain and + WASM targets. Do not update branches or discard unrelated work implicitly. +- [x] Record baseline `cargo test -p edgezero-core`, TS + `cargo test -p trusted-server-core request_timing`, and `cargo test-axum`. + Classify environment failures separately from code failures. +- [x] Pin the operation/result contract in the spec before implementing the API: + successful updates, unavailable reads/writes on contention, invalid slots + rejected without mutation, and first-write versus last-write semantics. +- [x] Define snapshot projection and payload updates under one lock. Mutable + access must not expose the origin or allow a request clock to be replaced. +- [x] Resolve callback failure behavior. Recommended contract to review: callbacks + are synchronous and non-reentrant; panics propagate normally; there is no + rollback; payloads must remain valid on unwind; poisoned mutexes recover + best-effort as TS does today. Do not advertise arbitrary callbacks as + infallible. If this contract is unsuitable, stop for an API decision. +- [x] Specify invalid-slot behavior for compound updates as well as plain records. + Validate slot inputs before invoking a payload callback so a rejected phase + cannot leave the payload changed. Do not introduce general transactions. +- [x] Record the concrete generic extension type and TS facade strategy. Agree on + names and minimum bounds; avoid requiring `D: Clone` for handle cloning. + +Exit: a short API contract in the spec, known baseline results, and no unresolved +safety semantics hidden in a TODO. This is a bounded API review, not another round +of choosing ownership or adding a generic renderer. + +Verification recorded for this implementation stage: + +- Base remains `98930917d96cb665c7255f36ca8bbd61fc7539e1`; no commit/push. + TS remains unchanged at `f951955f537b89392dcb863fca7a01a5f6fa4846`. +- Baseline EdgeZero core tests, TS core timing tests (13) and `cargo test-axum` + passed. The new consumer fixture initially failed to compile because the new + module did not exist; that is API red/green evidence, not a preexisting bug. +- Final core run: 484 unit tests, two integration tests, one compile-fail timing + lifetime doctest passed (13 existing doctests ignored). +- Workspace test, fmt, all-features clippy, feature check, generated-app build, + Fastly CLI/default gates, nested-config audits and excluded app-demo gates pass. +- Runtime contract suites pass: Cloudflare 8 in headless Chromium using lockfile- + matched wasm-bindgen 0.2.122, Fastly 7 in Viceroy 0.17.0, Spin 13 in Wasmtime + 44.0.1. Each includes the new portable clock/handle/router smoke test. Fastly's + WASM library suite passes all 88 tests. Adapter WASM check/clippy matrix passes, + including Fastly `fastly cli`. These are local runtime harnesses, not deployed + provider smoke tests or proof of TS migration. +- Rust 1.95.0; the missing `wasm32-wasip2` target was installed with approval. + Wasmtime and wasm-bindgen runners were installed run-locally, not globally. + Guide formatting, docs lint/format and VitePress build pass. +- Detailed command logs and report are outside the repository. No staged files. + Independent EdgeZero reviews subsequently cleared. Publication-stage core tests + (484 unit, two integration, one compile-fail doctest), fmt and all-features clippy + passed again. All TS migration/joint checks remain gates. + +## Task 2: Implement the generic collector and consumer fixture + +Files: + +- New `crates/edgezero-core/src/request_timing.rs` with colocated tests. +- `crates/edgezero-core/src/lib.rs` module declaration. +- New `crates/edgezero-core/tests/request_timing_consumer.rs` for public API usage. + +Work in small tested slices rather than implementing all methods at once. + +- [x] Add the bounded phase storage, application payload, single origin, and one + shared mutex. Use `web_time::Instant`; add no Tokio or UUID dependency. +- [x] Add construction, clone sharing, elapsed-from-origin access, saturating + accumulation, and consistent reads. Preserve absent versus recorded zero. +- [x] Add headers-ready/request-elapsed first-write marks and last-write bytes. + Do not stamp lifecycle marks automatically on middleware return or drop. +- [x] Add the narrow compound-update operation and its validated slot contract. + Test that contention prevents both phase and payload mutation. +- [x] Add drop-time phase spans. Cancellation records work so far, not request + completion; no destructor guarantee is made for panic-abort environments. +- [x] Test overflow, bounds, clone sharing, independent requests, first-write + marks, byte overwrite, compound snapshots, and the poison/panic behavior + agreed in task 1. Private test-only origin control is preferable to + wall-clock sleeps in core tests. +- [x] Compile an external consumer fixture with an application enum, a small + application payload, and a facade around the installed generic handle. + Prove extension insertion and handler access use one type and shared origin; + task 3 will add the middleware path. Use generic example phases, not auction + definitions, in EdgeZero tests. +- [x] Test public API use with a payload that is not `Clone`; verify the handle + still satisfies request-extension requirements with the necessary bounds. + +Check after each slice: `cargo test -p edgezero-core request_timing`. +Then run `cargo test -p edgezero-core`, including the consumer fixture. +Document that a missing API compilation failure is not runtime regression proof. +For existing behavior, use assertions that can detect incorrect values rather +than tests that merely check whether a field exists. + +Exit: public consumer fixture compiles, one-T0 and compound-state contracts pass, +and the TS facade can be expressed without a second clock or independently +installed wrapper extension. If not, return to the API checkpoint before adapters. + +## Task 3: Add shared attachment middleware and usage docs + +Files: + +- `crates/edgezero-core/src/middleware.rs` and its colocated tests. +- `crates/edgezero-core/tests/request_timing_consumer.rs`. +- `docs/guide/middleware.md`. + +- [x] Add generic attachment middleware using the collector from task 2. + Preserve an existing handle; create fresh state only for an eligible + request without one. Default to no exclusions. +- [x] Support the agreed static exact-path exclusion list. Compare URI paths, + not the full path/query string. Exclusion never removes an existing handle. +- [x] Verify default behavior, excluded paths, query strings, all methods, + duplicate registration, preinstalled handles, independent requests, + handler errors, and short-circuiting downstream middleware. +- [x] Exercise real router dispatch. Pin that unmatched 404/405 routes bypass + middleware and that an error rendered outside the chain remains outside + attachment middleware's response-finalization control. +- [x] Add an application-enum usage example based on the compiled fixture. + Explain opt-in installation, shared T0, payload ownership, and explicit + lifecycle marks. Add no generic header renderer. +- [x] Correct the guide's existing statement that middleware runs before route + matching when documenting this behavior. Use `web_time`, not a new + `std::time::Instant` example for cross-platform timing. Limit other guide + edits to the timing section and directly contradictory lifecycle wording. + +Checks: `cargo test -p edgezero-core`, then documentation formatting and lint for +the changed published guide. Internal spec/plan files are excluded from those +checks and need a direct whitespace/link/structure check. + +Exit: EdgeZero owns one reusable implementation with no TS dependency, and docs +show code checked by a real consumer fixture. + +## Task 4: Establish TS compatibility evidence before replacing mechanics + +Files in Trusted Server: + +- `crates/trusted-server-core/src/request_timing.rs`. +- `crates/trusted-server-adapter-axum/src/timing.rs`. +- Existing Cloudflare/Spin middleware tests and Fastly timing tests. + +- [ ] Capture current deterministic header and snapshot fixtures: order, precision, + missing versus zero, `u32` saturation, row-only phases, and auction fields. +- [ ] Add missing focused contract coverage for shared clones, compound wait and + placement, first-write marks/UUID, last-write bytes, and contention fallback. + These characterize intended current behavior and should pass before moving it. +- [ ] Add an Axum regression that inserts a pre-aged handle with a recorded phase + in an upstream service. Assert handler and private-response header retain + the same measurement. Run it before changing production code and capture + the failure caused by unconditional replacement. +- [ ] Pin health method differences: Cloudflare/Spin/Axum exclude every method on + `/health`; Fastly skips collection for `GET /health` only. Include query + strings and non-GET requests. Do not normalize these differences. + +Focused checks: `cargo test -p trusted-server-core request_timing`, +`cargo test-axum`, `cargo test-cloudflare`, `cargo test-spin`, and the nearest +Fastly health/timing filters through `cargo test-fastly`. + +Exit: current contracts are executable and the Axum reset regression has +failing-before evidence. Do not treat that deliberate failure as a baseline defect +to work around by weakening its assertions. + +## Task 5: Migrate TS against the candidate EdgeZero revision + +Files in Trusted Server: + +- `crates/trusted-server-core/src/request_timing.rs`. +- Cloudflare/Spin `src/middleware.rs` and `src/app.rs`. +- Axum `src/timing.rs` and Fastly `src/main.rs`, plus any existing attachment sites + in their app modules identified by searching `RequestTimings`. +- Core timing consumers only where generic extension lookup requires migration. +- Temporary exact pushed git-revision pins and regenerated lockfile; replace with + the eventual approved release tag before final TS delivery. + +- [ ] Record the exact candidate EdgeZero revision and any uncommitted diff used + for validation. After the reviewed EdgeZero draft PR is pushed, temporarily + pin all six TS dependencies to that exact git commit. Do not publish + developer-specific paths or local patches. +- [ ] Resolve the complete EdgeZero package dependency closure from the same + candidate git source, including macros/adapter registry. + Check `cargo tree`/metadata for duplicate EdgeZero core identities. Mixed + sources giving middleware and handlers different core types are invalid. +- [ ] Preserve the TS `Phase` enum and domain method names. Delegate generic + mechanics to EdgeZero, with auction state in its application payload. + Keep header rendering, privacy checks and access serialization local. +- [ ] Read the installed generic extension through the TS facade at every lookup. + Audit both initial-document and page-bids paths, telemetry, streaming, and + adapters. Do not leave an old-wrapper lookup that silently falls back to a + fresh collector or makes the page-bids missing-collector gate fire. +- [ ] Replace both Cloudflare/Spin middleware copies with shared registrations. + Keep sanitization first, preserve exclusions, and retain adapter integration + tests even though generic middleware unit tests now live upstream. +- [ ] Make Axum retrieve-or-create the same handle for non-excluded requests and + use it for terminal rendering. Run the task 4 regression and show it passes. +- [ ] Keep Fastly's clock creation and streaming/finalization boundaries unchanged. + Reuse its existing handle through every core-request conversion path. +- [ ] Preserve both diagnostic clock labels and browser payloads. Do not add the + separate initial-document missing-collector follow-up in this migration. +- [ ] Remove superseded collector mechanics and middleware tests only when their + replacement coverage exists. No parallel old/new collector at runtime. + +Checks after each adapter slice: its target-matched tests. Re-run core timing, +page-bids, initial-document diagnostic tests, and the Axum regression at the end. +Compare against task 4 fixtures; explain any intentional difference before +continuing. Only Axum's preinstalled-handle preservation is planned to improve. + +Exit: TS uses the candidate implementation on all four adapters with unchanged +wire/privacy behavior and no circular dependency or duplicate timing clock. + +## Task 6: Joint validation and independent review + +Run commands from their respective repository roots. Record exact revisions, +temporary git-revision pins, results, skipped checks, and environment blockers. A native pass +is not evidence of a WASM runtime pass. + +EdgeZero native and build gates: + +```sh +cargo test -p edgezero-core +cargo test --workspace --all-targets +cargo fmt --all -- --check +cargo clippy --workspace --all-targets --all-features -- -D warnings +cargo check --workspace --all-targets --features "fastly cloudflare spin" +cargo test -p edgezero-cli --test generated_project_builds -- --ignored +``` + +Run the current `.github/workflows/test.yml` and `format.yml` gates, including +feature-specific Fastly tests/clippy, the separately excluded app-demo workspace, +and WASM adapter contract tests/checks/clippy. Target mapping: + +| Adapter | EdgeZero target | Runtime runner | +| --- | --- | --- | +| Cloudflare | `wasm32-unknown-unknown` | Lockfile-matched `wasm-bindgen-test-runner` | +| Fastly | `wasm32-wasip1` | Pinned Viceroy | +| Spin | `wasm32-wasip2` | Pinned Wasmtime | + +For example, the Spin check is +`cargo check -p edgezero-adapter-spin --target wasm32-wasip2 --features spin`. +Adapter test/check command shapes and runner setup come from those workflows. +Add a clock/handle smoke assertion to the existing adapter contract tests if +needed to exercise the new collector on each runtime, not just compile it. +Keep native-only sleep/thread/poison tests out of unsupported WASM test paths. + +Trusted Server gates: + +```sh +cargo test-fastly +cargo test-axum +cargo test-cloudflare +cargo test-spin +./scripts/test-cli.sh +cargo test --manifest-path crates/trusted-server-integration-tests/Cargo.toml --test parity +cargo fmt --all -- --check +cargo clippy-fastly +cargo clippy-axum +cargo clippy-cloudflare +cargo clippy-cloudflare-wasm +cargo clippy-spin-native +cargo clippy-spin-wasm +``` + +TS Spin targets `wasm32-wasip1`, unlike EdgeZero's `wasm32-wasip2` matrix. Preserve +and validate both supported configurations; do not copy one repository's target +assumptions into the other. + +Run TS JavaScript diagnostics regressions and the full Vitest suite from +`crates/trusted-server-js/lib` with `npx vitest run`. Build browser artifacts and +bring up the existing Next.js integration environment using that suite's setup, +then run `npm run test:nextjs -- tests/nextjs/gpt-diagnostics.spec.ts` from +`crates/trusted-server-integration-tests/browser`. Run the applicable full browser +and integration CI jobs before PR handoff. No JS behavior change is planned. + +- [x] Fresh read-only review of collector API and synchronization contracts. +- [ ] Fresh read-only review of TS extension identity, timing origins, privacy, + and release pin compatibility. No concurrent writers during fixes. +- [ ] Fix accepted findings and re-run the affected checks. + +Exit: both repositories pass applicable gates against the same candidate; +remaining risks are explicit. Release is blocked on unverified consumer adoption. + +## Task 7: Publish and finish coordinated adoption + +Fresh independent EdgeZero reviews cleared and draft-PR publication is now +authorized. Keep the PR draft and do not request reviewers during this publication +stage. Merge/tag and final delivery remain gated on separate owner approval. The temporary TS git-revision pin is for joint review, not the final +release-backed delivery pin. + +- [ ] Record linked EdgeZero and TS PRs and their tested commits. Verify no new + changes invalidate joint evidence before approving the upstream release. +- [ ] Merge/tag EdgeZero through its normal release process. Select the tag then; + do not invent a tag or assume crates.io publication. +- [ ] Point all six TS EdgeZero dependencies at the released tag. Remove every + local override, regenerate `Cargo.lock`, and inspect the full dependency diff. +- [ ] Verify the released source matches the candidate tested jointly. If merge + changes alter it, re-run joint checks rather than relying on stale results. +- [ ] Re-run TS gates without overrides. From the TS root, run + `cargo metadata --locked --format-version 1` and inspect the resolved + EdgeZero package IDs and git sources against the released tag and commit. + Confirm the lockfile uses the released git source, no local paths, + and one intended EdgeZero core identity. +- [ ] Deliver TS only after final tag-backed validation and approval. Reply to + Aram with evidence, including the deliberate choice to keep rendering local. + Resolve review threads only with separate authorization. + +Completion: one shared EdgeZero middleware, generic collection used by all TS +adapters, no duplicate clocks, unchanged TS output/privacy contracts, and final +release-backed dependency pins. The other middleware duplication and the +initial-document missing-collector guard remain separate follow-ups. diff --git a/docs/superpowers/specs/2026-09-28-request-timing-design.md b/docs/superpowers/specs/2026-09-28-request-timing-design.md new file mode 100644 index 00000000..30b776f6 --- /dev/null +++ b/docs/superpowers/specs/2026-09-28-request-timing-design.md @@ -0,0 +1,449 @@ +# Shared request timing collector and middleware + +Date: 2026-09-28 +Status: Design and API checkpoint approved on 2026-09-28. EdgeZero implemented, locally verified and independently reviewed; draft-PR publication authorized. TS migration and joint validation pending. +Scope: `edgezero-core` and coordinated Trusted Server adoption, shipped together. +Plan: [Implementation plan](../plans/2026-09-28-request-timing.md). + +Implementation of both repositories and later draft-PR publication are authorized. +Independent EdgeZero reviews cleared; a focused commit, normal push and draft PR targeting main are authorized. TS remains unchanged during this publication stage. +Merge, release/tag, and final delivery still require separate approval. + +## Why this change + +Trusted Server PR [#1121](https://github.com/IABTechLab/trusted-server/pull/1121) +adds request-timing middleware in both its Cloudflare and Spin adapters. The two +implementations are identical. Aram's +[blocking review comment](https://github.com/IABTechLab/trusted-server/pull/1121#discussion_r4090027047) +asks for the generic collector and middleware to live in EdgeZero, while keeping +auction-specific state in Trusted Server. He also accepts an interim move of the +duplicated middleware into `trusted-server-core`. + +This is a cross-repository API change, not just moving a struct. EdgeZero cannot +depend on Trusted Server, and replacing the collector must not reset the clock +used by auction diagnostics or change telemetry and header output. + +Initial investigation used EdgeZero `567964158e4f8bd0d52321b9801de44966422e1b`, +locally tagged `v0.0.8`. Before creating the implementation worktree, local `main` +was fast-forwarded to fetched `origin/main` at +`98930917d96cb665c7255f36ca8bbd61fc7539e1`. The intervening changes leave the timing +integration points in core middleware, router, context, and Cargo dependencies +unchanged. The middleware guide changed and must be edited against this new base. + +The Trusted Server reference is `f951955f537b89392dcb863fca7a01a5f6fa4846`; it pins +six EdgeZero workspace dependencies to `v0.0.7`. This remains the unchanged TS baseline. EdgeZero tasks 1–3 are implemented and +locally verified against the base above; see the plan for evidence and remaining gates. + +## Approved decisions + +1. Applications use typed phase enums. EdgeZero stores bounded indexed slots; + Trusted Server's typed facade owns the enum-to-slot mapping. No string registry + or new public phase trait is needed for the first release. +2. A small application-defined payload shares the collector's clock and mutex. + Narrow synchronous update and snapshot operations preserve compound facts. +3. Install one concrete EdgeZero handle type in request extensions. Trusted Server + adds domain methods through a facade over that same handle, not a second clock. +4. Keep `Server-Timing` rendering and exposure policy in Trusted Server for the + first release. A generic renderer is out of scope. +5. Ship EdgeZero and Trusted Server adoption as one coordinated delivery. Do not + perform a temporary Trusted Server-only middleware move. + +## Goals + +- Provide one portable per-request clock and a reusable collector in EdgeZero. +- Attach a fresh collector through shared middleware only when one is absent. +- Preserve existing application marks, compound updates, and consistent snapshots. +- Keep collection independent from exposing timings in headers, logs, or bodies. +- Let Trusted Server remove its duplicate middleware without changing its public + diagnostic payloads or access-log schema. + +## Non-goals + +- Moving auction concepts, UUIDs, EC storage names, or Trusted Server privacy + policy into EdgeZero. +- Moving the other duplicated Trusted Server middlewares in this change. +- Automatically instrumenting every adapter, route, or streamed body. +- Changing router dispatch, error rendering, macros, manifests, or extractors. +- Adding tracing export, dynamic metric registration, retries, or a public clock + injection framework. +- Changing the browser's SSAT classification or comparing browser and server clocks. + +## Ownership + +| Concern | Owner after migration | +| --- | --- | +| Portable monotonic origin and cloneable request handle | EdgeZero | +| Bounded phase durations, saturating accumulation, drop-time spans | EdgeZero | +| Headers-ready and request-elapsed marks, response byte count | EdgeZero | +| Nonblocking collection and consistent snapshot mechanics | EdgeZero | +| Insert-if-absent middleware and configurable exclusion mechanism | EdgeZero | +| Eight-phase enum and mapping to generic storage | Trusted Server | +| Auction wait placement, UUID, dispatch/resolve/commit semantics | Trusted Server | +| `ts-*` header names, ordering, and existing serialization | Trusted Server | +| Header enablement, cache privacy checks, browser activation gates | Trusted Server | +| Where each adapter starts and finishes measurement | Application/adapter integration | + +The reviewer also identifies `server_timing_value` as movable. The approved scope +keeps Trusted Server's renderer in place for the first extraction: its names and +header-bearing phase list are application policy. A generic renderer would add +validation and ordering contracts that are unnecessary for this migration. +Explain this narrower scope in the review response; do not claim that the whole +of the reviewer's suggested extraction moved upstream. + +## Clock and lifecycle contract + +T0 is the instant the request collector is created at its configured application +or adapter boundary. It is not browser `navigationStart`, socket acceptance, or +an assumed uniform platform ingress timestamp. Existing adapters start after +different amounts of prologue work. The migration preserves those boundaries. + +Every clone, phase span, generic lifecycle mark, and application milestone must +refer to that same origin. Wrapping an existing handle must not allocate another +clock or an independent copy of its mutable state. + +```mermaid +flowchart TD + A[Adapter receives request] --> B{Collector already attached?} + B -->|Yes| C[Reuse existing handle and T0] + B -->|No| D[Configured attachment boundary] + D --> E{Request excluded?} + E -->|Yes| F[Continue without installing a collector] + E -->|No| G[Create request-local collector] + C --> H[Middleware and application record timings] + G --> H + H --> I[Application marks headers ready at its terminal boundary] + I --> J[Application emits permitted header timing] + J --> K[Application records body completion where observable] + K --> L[Application renders permitted access telemetry] +``` + +The diagram describes ownership, not a new automatic finalization pipeline. +Header rendering occurs at headers-ready; later body marks can appear only in +later snapshots such as access telemetry. + +### Middleware behavior + +The shared middleware is opt-in and attachment-only: + +1. Preserve a collector of the configured concrete type if already present. +2. Otherwise evaluate the application's exclusion policy. +3. For an eligible request, create and insert a fresh collector. +4. Pass the context to `next.run` and preserve its result unchanged. + +Default configuration excludes no paths. Trusted Server's Cloudflare and Spin +registrations configure an exact `/health` path exclusion for every method, +including requests with query strings on that path. Exclusion prevents +installation, not removal of a preinstalled collector. Duplicate middleware +installation must neither reset T0 nor erase recorded facts. + +Preserve the existing adapter differences: Axum also excludes `/health` for every +method, while Fastly short-circuits only `GET /health` before collector creation. +Non-GET Fastly health requests continue through normal timing setup. Normalizing +that policy across adapters would be a separate behavior change. + +Keep the exclusion API small. A static path list is sufficient for the current +Cloudflare/Spin consumers; a general predicate is an alternative only if a +concrete need emerges. Fastly keeps its existing method-aware entry-point check. +Neither collector nor application payload belongs in `RouterBuilder::with_state`: +that state is shared across requests and can overwrite same-type extensions. + +Register the middleware before consumers that need timing. Trusted Server's +client-IP sanitization must remain first; extraction must not change that order. + +### Boundaries the middleware cannot cover + +EdgeZero matches a route before running middleware. Unmatched 404/405 requests +bypass it. `RouterService::oneshot` also renders dispatch errors outside the +middleware chain. Returning a response does not mean a streaming body has been +consumed or sent. + +Consequently: + +- Axum retains its outer terminal timing service and error-response coverage. + For non-excluded requests, the proposed migration must retrieve the configured + concrete handle if present, creating and inserting one only when absent. Both + the handler and terminal renderer use that handle. This deliberately strengthens + the current wrapper, which unconditionally replaces a preinstalled collector; + it is not a claim about current behavior. The health bypass remains unchanged. +- Fastly retains its early collector creation, app-build span, header emission, + streaming duration, and successfully written-byte accounting. +- Cloudflare and Spin adopt the shared attachment middleware without gaining + `Server-Timing` emission incidentally. +- No middleware return or guard drop automatically stamps request completion. + +A cancelled phase span can record work up to its drop. That is not evidence that +the request, response body, or auction completed. Panic-abort environments do not +guarantee destructor execution. + +## Collector shape + +The ownership and state arrangement below are approved. Exact Rust signatures and +the payload failure contract need the API checkpoint in the implementation plan. + +Use a cloneable handle backed by one `Arc` and one mutex. The shared state holds +an immutable origin, generic lifecycle fields, a bounded phase array, and an +application-defined payload. A candidate name is +`RequestTimings`. + +The payload lets Trusted Server keep auction definitions in its own crate while +sharing synchronization with the generic fields. For example, recording auction +wait currently updates both a phase total and its placement under one lock. +Snapshots read these together. Splitting them across two locks would change that +contract even if both collectors shared T0. + +Generic operations should cover: + +- Construction, elapsed time from T0, and cheap handle cloning. +- Saturating phase accumulation and a span that records elapsed time on drop. +- First-successful-write headers-ready and request-elapsed marks. +- Last-write response byte count. +- A consistent read of generic fields and application payload. +- A short synchronous update of a phase and application facts under the same lock. + +Construction with `D::default()` is sufficient for middleware installation. Core +operations need not require every payload to implement `Default` or `Clone`. +Manual handle cloning must not accidentally impose `D: Clone`. Installed handles +must satisfy request-extension `Clone + Send + Sync + 'static` requirements; +middleware futures remain `?Send`. + +Trusted Server can keep its existing domain methods behind a facade over the +concrete EdgeZero handle. Middleware and consumers must use the same extension +type. A facade that reads the generic handle is acceptable; a wrapper type that +expects an independently installed extension is not. The first API implementation +slice must include a compiled usage fixture of this boundary before adapter +migration begins. + +### Phase representation alternatives + +| Option | Benefits | Costs | +| --- | --- | --- | +| Application enum implementing an EdgeZero phase trait | Typed generic calls and application-owned names | Must define count, unique mapping, and invalid mappings | +| Fixed indexed slots, with an application enum facade | Preserves bounded storage and current TS mapping with a small API | Raw indices need bounds handling | +| String-keyed phases | Convenient arbitrary labels | Requires allocation, cardinality limits, duplicate and overflow policy | + +Selected: fixed indexed slots inside EdgeZero, with typed phase enums at +application call sites. `N` is chosen at compile time, not from a request. Invalid +indices must be rejected without panicking or mutating state, and the caller must +be able to distinguish them from a successful write. Do not silently alias them +to another phase. The exact result type remains an API review decision. + +### Synchronization and failure behavior + +Preserve nonblocking `try_lock` behavior. A contended update drops the whole +sample; a contended read reports unavailable. Trusted Server maps unavailable +reads to its existing all-None snapshot or omitted header. Do not retry, block, +or substitute zero for missing facts. + +Any application update callback must be synchronous, short, and non-reentrant. +It cannot retain a guard across an await or call a lock-taking collector method +from inside itself. One lock makes successful compound updates and snapshots +consistent; it does not provide rollback if a callback panics partway through. + +The existing TS collector recovers poisoned locks because its state is simple +counters and optional facts. Arbitrary application payloads do not necessarily +have that property. Before accepting a public payload-update API, explicitly +choose its panic/poison contract. A candidate is documented best-effort recovery +with callbacks required to leave valid state even on unwind. An alternative is a +narrower mutation interface. Do not promise unconditional infallibility for +arbitrary user code, and do not change TS poison recovery without a deliberate +compatibility decision. + +## Trusted Server compatibility requirements + +These are migration acceptance criteria, not application features to add to EdgeZero. + +| Existing behavior | Required outcome | +| --- | --- | +| Repeated phase durations | Saturating addition; missing differs from recorded zero | +| Headers-ready, completion, auction milestones and UUID | First successful write wins | +| Response bytes and auction placement | Last successful write wins | +| Auction wait | Accumulate duration and update placement in one operation | +| Snapshot durations | Whole milliseconds, truncated and saturated to `u32::MAX` | +| Header durations | One decimal place; existing names, order and omission rules | +| Header emission | Mark headers-ready even when emission is disabled; append rather than overwrite | +| Cache policy | Retain existing private/no-store predicate; never change cacheability | +| Browser timing payload | Preserve field names, activation gates and omitted-value behavior | +| Access telemetry | Preserve field names, nullability, UUID join key and placement encoding | + +Auction UUID and dispatch marking remain independent. An attempted auction may +have an ID without a dispatch mark; an abandoned auction may have dispatch but no +resolution. Missing samples can also result from contention. Do not reinterpret +missing timing as proof that no auction occurred. + +The page-bids path already omits diagnostics when the adapter collector is +missing. The initial-document path lacks the equivalent guard. Aram identified +that as a nonblocking +[follow-up](https://github.com/IABTechLab/trusted-server/pull/1121#discussion_r4090027050), +not a currently reachable production failure. Track it separately with an +initial-document missing-collector regression test; no browser classification +change is required by that finding. + +## Delivery sequence + +```mermaid +flowchart TD + A[Confirm API details and authorize implementation] --> B[Implement EdgeZero collector and middleware] + B --> C[Integrate TS against the candidate EdgeZero revision] + C --> D[Validate both repositories before release] + D --> E[Approve EdgeZero merge and git tag] + E --> F[Pin all TS EdgeZero dependencies to that tag] + F --> G[Revalidate without local overrides] + G --> H[Approve coordinated TS delivery] +``` + +1. Resolve the implementation checkpoints below and review the linked plan before + starting code changes. +2. Implement in small tested steps, starting with + `crates/edgezero-core/src/request_timing.rs`, its declaration in `src/lib.rs`, + and shared middleware in `src/middleware.rs`. Colocate tests. Add usage docs in + `docs/guide/middleware.md`. Router changes are not expected. +3. Prepare Trusted Server adoption against the candidate EdgeZero revision using + temporary exact git-revision pins after the reviewed EdgeZero draft PR is + pushed. Pin all six TS dependencies to that same commit and verify one core + package identity, without developer-specific paths. Replace collector mechanics + and the Cloudflare/Spin middleware copies; preserve the TS domain facade and + Fastly/Axum boundaries. + Validate both repositories before considering the upstream release ready. +4. With separate publication approval, merge and tag EdgeZero through its normal + git-tag process. Do not assume a tag number or promise crates.io publication; + the inspected workspace has `publish = false`. +5. Update all six EdgeZero pins in TS `Cargo.toml` and review regenerated + `Cargo.lock`, including changes since `v0.0.7` unrelated to timing. Remove local + overrides and revalidate against the released tag before TS delivery. + +Ship together means coordinated readiness and dependent releases, not simultaneous +Git operations across repositories. EdgeZero's tag must exist before the final TS +pin can resolve. Temporary exact pushed git-revision pins support joint PR +validation only; replace them with the approved release tag before final TS +delivery. Neither side is considered complete without verified TS adoption. +No intermediate TS-only middleware move, unrelated lock-helper cleanup, or other +adapter middleware extraction is included. + +## Verification plan + +### Collector contracts + +Use deterministic durations and private test-only origin control where needed; +no public clock trait is required just for testing. Cover: + +- Clone sharing, independent requests, and preservation of a pre-aged T0. +- Absent versus recorded-zero phases, accumulation, saturation and invalid slots. +- First-write marks, last-write byte counts, and explicit lifecycle completion. +- Compound phase/payload updates and consistent snapshots. +- Contention dropping a whole operation and poison behavior under the chosen contract. +- Span drop and cancellation without fabricated request completion. + +### Middleware and consumer contracts + +Exercise the real router interface with `futures::executor::block_on`. Verify +insert-if-absent, repeated installation, exact health exclusion, preinstalled +handles on excluded paths, ordering, short circuits, and unchanged errors. +Pin the documented unmatched-route boundary rather than implying blanket coverage. +Test GET and non-GET `/health`, with and without query strings, against each +adapter's preserved method policy. For Axum, attach a pre-aged collector in an +upstream service and prove that its origin and recorded phases survive both +handler access and private-response header rendering. This regression should +fail against the current unconditional replacement before that wrapper changes. + +Before replacing TS mechanics, establish its existing output fixtures. Retain +private-header append/enable/cache tests, Axum terminal error coverage, Fastly +streaming byte counts, and both document and page-bids diagnostic gates. Prove +that a pre-aged collector yields nonzero pre-dispatch time without changing the +origin during extraction. Verify all four adapters still deliver the same handle +to their core consumers. + +### Commands for implementation + +EdgeZero: + +```sh +cargo test -p edgezero-core +cargo test --workspace --all-targets +cargo fmt --all -- --check +cargo clippy --workspace --all-targets --all-features -- -D warnings +cargo check --workspace --all-targets --features "fastly cloudflare spin" +cargo check -p edgezero-adapter-spin --target wasm32-wasip2 --features spin +``` + +Also run the repository's adapter WASM CI matrix, including Fastly +`wasm32-wasip1` and Cloudflare `wasm32-unknown-unknown`, plus existing generated-app +and demo checks relevant to a public core API change. Compile checks do not prove +runtime clock behavior; distinguish native execution, WASM compilation, and any +platform clock smoke tests in the implementation report. + +Trusted Server uses its target-matched `cargo test-fastly`, `cargo test-axum`, +`cargo test-cloudflare`, `cargo test-spin`, and `./scripts/test-cli.sh` commands. +Run `cargo fmt --all -- --check` and applicable adapter clippy aliases, including +WASM variants. Retain cross-adapter integration coverage and run JS/browser +regressions for auction diagnostics. Do not use bare workspace-wide Cargo tests +across its incompatible targets. + +For the internal spec/plan, inspect links, Markdown structure, and the final diff. +`docs/superpowers/**` is excluded from both the published site and the normal docs +Prettier check, so a passing site check would not validate this document. + +## Approved API checkpoint + +The supervisor approved these exact contracts before production edits: + +- `RequestTimings` holds one `Arc>>`. + `with_data(D)` needs neither `Default` nor `Clone`; `new`/`Default` need + `D: Default`. Manual handle cloning never needs `D: Clone`. Extension and + middleware payloads need `D: Send + 'static`, not `D: Sync`. +- `TimingError::{InvalidSlot, Unavailable}` distinguishes invalid slots from + contention. `record`, lifecycle marks and `set_resp_bytes` return + `Result<(), TimingError>`; `elapsed` returns `Result`. + Repeated first-write marks return `Ok(())` without replacing the value. +- `update_data(FnOnce(&mut D, Duration) -> R)` and + `record_with(slot, duration, FnOnce(&mut D, Duration) -> R)` return + `Result`. Duration is elapsed from the immutable origin. + `record_with` validates its slot before acquiring the lock or invoking the + callback, then accumulates the phase and invokes the callback under one lock. +- `snapshot(FnOnce(TimingSnapshot<'_, N, D>) -> R)` projects copied generic + fields/current elapsed and borrowed `&D` under that same lock. Neither the + borrowed payload nor the guard can escape; payload cloning is not required. +- Contention invokes no callback and changes no state. Callbacks are synchronous, + short and non-reentrant. Panics propagate, may leave phase/payload partially + changed, and poison the mutex. Poison recovery is best-effort, **not rollback** + or payload validation; callbacks must preserve payload validity on unwind. +- `span(slot)` returns `Result, TimingError>` after bounds + validation. Drop attempts a best-effort record, discarding contention; it never + marks request completion. +- `RequestTimingMiddleware` defaults to no exclusions and offers + `with_excluded_paths(&'static [&'static str])`. The TS facade reads exactly + `RequestTimings<8, TsPayload>` from extensions rather than installing itself. + +## Implementation checkpoints + +The five design decisions above are settled. These details remain to be pinned +before the corresponding implementation step, not reopened as scope choices: + +1. Specify update/read result types, invalid-slot handling, callback panic and + poison behavior, and reentrancy restrictions. Recovering a poisoned mutex must + not be described as rollback or validation of arbitrary payload state. +2. Compile a typed-phase, concrete-extension, domain-facade example against the + first API slice. Prove shared T0 and atomic phase/payload snapshots before + migrating adapters. +3. Record exact EdgeZero and TS revisions for joint validation. Choose a release + tag only through the normal approved publication process. + +If the API checkpoint requires a different state model or broadens the agreed +scope, return for a decision rather than implementing a silent redesign. + +## Source references + +EdgeZero files are relative to this repository at the revision named above: + +- `crates/edgezero-core/src/middleware.rs`: `Middleware`, `Next`, `RequestLogger`. +- `crates/edgezero-core/src/router.rs`: route matching, state injection, and + error-to-response conversion. +- `crates/edgezero-core/src/context.rs`: public request-extension access. +- `crates/edgezero-core/Cargo.toml`: existing `web-time` dependency. +- `.github/workflows/test.yml` and `.github/workflows/format.yml`: target matrices. + +Trusted Server sources at the reviewed revision: + +- [Collector and application rendering](https://github.com/IABTechLab/trusted-server/blob/f951955f537b89392dcb863fca7a01a5f6fa4846/crates/trusted-server-core/src/request_timing.rs). +- [Axum terminal timing service](https://github.com/IABTechLab/trusted-server/blob/f951955f537b89392dcb863fca7a01a5f6fa4846/crates/trusted-server-adapter-axum/src/timing.rs). +- [Fastly request and delivery boundaries](https://github.com/IABTechLab/trusted-server/blob/f951955f537b89392dcb863fca7a01a5f6fa4846/crates/trusted-server-adapter-fastly/src/main.rs). +- [Publisher diagnostics](https://github.com/IABTechLab/trusted-server/blob/f951955f537b89392dcb863fca7a01a5f6fa4846/crates/trusted-server-core/src/publisher.rs).