diff --git a/.cargo/config.toml b/.cargo/config.toml index f8fcfadaf..aea0df2d5 100644 --- a/.cargo/config.toml +++ b/.cargo/config.toml @@ -30,6 +30,10 @@ build-fastly = "build -p trusted-server-core -p trusted-server-adapter-fastly -p check-fastly = "check -p trusted-server-core -p trusted-server-adapter-fastly -p trusted-server-js -p trusted-server-openrtb --target wasm32-wasip1" clippy-fastly = "clippy -p trusted-server-core -p trusted-server-adapter-fastly -p trusted-server-js -p trusted-server-openrtb --all-targets --all-features --target wasm32-wasip1 -- -D warnings" test-fastly = "test -p trusted-server-core -p trusted-server-adapter-fastly -p trusted-server-js -p trusted-server-openrtb --target wasm32-wasip1" +# Feature-on counterpart of `test-fastly`. `test-fastly` does NOT pass +# --all-features, so without this the reusable-sandbox code is linted by +# `clippy-fastly` but its tests never execute. +test-fastly-reuse = "test -p trusted-server-adapter-fastly --target wasm32-wasip1 --features reusable-sandbox" # --- Axum adapter (native dev server) --- build-axum = "build -p trusted-server-adapter-axum" diff --git a/.github/workflows/test.yml b/.github/workflows/test.yml index a7fc78d08..545a4d3a6 100644 --- a/.github/workflows/test.yml +++ b/.github/workflows/test.yml @@ -59,6 +59,9 @@ jobs: - name: Run tests run: cargo test-fastly + - name: Run tests (reusable-sandbox feature) + run: cargo test-fastly-reuse + - name: Run template cache ESI local harness run: BID_DELAY=3 ./scripts/template-cache-local-test.sh esi diff --git a/AGENTS.md b/AGENTS.md index 738fcab34..2780414f6 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -98,6 +98,7 @@ spin up --from crates/trusted-server-adapter-spin # See .cargo/config.toml; default-members = [fastly] so Viceroy can locate # the binary via `cargo run --bin`. cargo test-fastly # Fastly adapter + core (wasm32-wasip1 via Viceroy) +cargo test-fastly-reuse # Fastly adapter with the reusable-sandbox feature on cargo test-axum # Axum dev server adapter (native) cargo test-cloudflare # Cloudflare Workers adapter (native host) cargo test-spin # Spin adapter route tests (native host) @@ -346,7 +347,7 @@ Every PR must pass: 1. `cargo fmt --all -- --check` 2. `cargo clippy-fastly && cargo clippy-axum && cargo clippy-cloudflare && cargo clippy-cloudflare-wasm && cargo clippy-spin-native && cargo clippy-spin-wasm && cargo clippy-cli && cargo clippy-codegen` -3. `cargo test-fastly && cargo test-axum && cargo test-cloudflare && cargo test-spin` +3. `cargo test-fastly && cargo test-fastly-reuse && cargo test-axum && cargo test-cloudflare && cargo test-spin` 4. `cargo test --manifest-path crates/trusted-server-integration-tests/Cargo.toml --test parity` 5. JS build and test (`cd crates/trusted-server-js/lib && npx vitest run`) 6. JS format (`cd crates/trusted-server-js/lib && npm run format`) diff --git a/Cargo.lock b/Cargo.lock index 7a45462b1..4a1c06bcd 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1428,7 +1428,7 @@ dependencies = [ [[package]] name = "edgezero-adapter" version = "0.1.0" -source = "git+https://github.com/stackpop/edgezero?tag=v0.0.8#567964158e4f8bd0d52321b9801de44966422e1b" +source = "git+https://github.com/stackpop/edgezero?rev=35a72835322fe0127beb2ce998e4988e95974373#35a72835322fe0127beb2ce998e4988e95974373" dependencies = [ "toml", ] @@ -1436,7 +1436,7 @@ dependencies = [ [[package]] name = "edgezero-adapter-axum" version = "0.1.0" -source = "git+https://github.com/stackpop/edgezero?tag=v0.0.8#567964158e4f8bd0d52321b9801de44966422e1b" +source = "git+https://github.com/stackpop/edgezero?rev=35a72835322fe0127beb2ce998e4988e95974373#35a72835322fe0127beb2ce998e4988e95974373" dependencies = [ "anyhow", "async-trait", @@ -1464,7 +1464,7 @@ dependencies = [ [[package]] name = "edgezero-adapter-cloudflare" version = "0.1.0" -source = "git+https://github.com/stackpop/edgezero?tag=v0.0.8#567964158e4f8bd0d52321b9801de44966422e1b" +source = "git+https://github.com/stackpop/edgezero?rev=35a72835322fe0127beb2ce998e4988e95974373#35a72835322fe0127beb2ce998e4988e95974373" dependencies = [ "anyhow", "async-trait", @@ -1487,7 +1487,7 @@ dependencies = [ [[package]] name = "edgezero-adapter-fastly" version = "0.1.0" -source = "git+https://github.com/stackpop/edgezero?tag=v0.0.8#567964158e4f8bd0d52321b9801de44966422e1b" +source = "git+https://github.com/stackpop/edgezero?rev=35a72835322fe0127beb2ce998e4988e95974373#35a72835322fe0127beb2ce998e4988e95974373" dependencies = [ "anyhow", "async-stream", @@ -1516,7 +1516,7 @@ dependencies = [ [[package]] name = "edgezero-adapter-spin" version = "0.1.0" -source = "git+https://github.com/stackpop/edgezero?tag=v0.0.8#567964158e4f8bd0d52321b9801de44966422e1b" +source = "git+https://github.com/stackpop/edgezero?rev=35a72835322fe0127beb2ce998e4988e95974373#35a72835322fe0127beb2ce998e4988e95974373" dependencies = [ "anyhow", "async-trait", @@ -1543,7 +1543,7 @@ dependencies = [ [[package]] name = "edgezero-cli" version = "0.1.0" -source = "git+https://github.com/stackpop/edgezero?tag=v0.0.8#567964158e4f8bd0d52321b9801de44966422e1b" +source = "git+https://github.com/stackpop/edgezero?rev=35a72835322fe0127beb2ce998e4988e95974373#35a72835322fe0127beb2ce998e4988e95974373" dependencies = [ "chrono", "clap", @@ -1568,7 +1568,7 @@ dependencies = [ [[package]] name = "edgezero-core" version = "0.1.0" -source = "git+https://github.com/stackpop/edgezero?tag=v0.0.8#567964158e4f8bd0d52321b9801de44966422e1b" +source = "git+https://github.com/stackpop/edgezero?rev=35a72835322fe0127beb2ce998e4988e95974373#35a72835322fe0127beb2ce998e4988e95974373" dependencies = [ "anyhow", "async-compression", @@ -1599,7 +1599,7 @@ dependencies = [ [[package]] name = "edgezero-macros" version = "0.1.0" -source = "git+https://github.com/stackpop/edgezero?tag=v0.0.8#567964158e4f8bd0d52321b9801de44966422e1b" +source = "git+https://github.com/stackpop/edgezero?rev=35a72835322fe0127beb2ce998e4988e95974373#35a72835322fe0127beb2ce998e4988e95974373" dependencies = [ "log", "proc-macro2", diff --git a/Cargo.toml b/Cargo.toml index ac0cac621..e1fb260c1 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -55,12 +55,12 @@ cssparser = "0.36" derive_more = { version = "2.0", features = ["display", "error"] } directories = "5" ed25519-dalek = { version = "2.2", features = ["rand_core"] } -edgezero-adapter-axum = { git = "https://github.com/stackpop/edgezero", tag = "v0.0.8", default-features = false } -edgezero-adapter-cloudflare = { git = "https://github.com/stackpop/edgezero", tag = "v0.0.8", default-features = false } -edgezero-adapter-fastly = { git = "https://github.com/stackpop/edgezero", tag = "v0.0.8", default-features = false } -edgezero-adapter-spin = { git = "https://github.com/stackpop/edgezero", tag = "v0.0.8", default-features = false } -edgezero-cli = { git = "https://github.com/stackpop/edgezero", tag = "v0.0.8" } -edgezero-core = { git = "https://github.com/stackpop/edgezero", tag = "v0.0.8", default-features = false } +edgezero-adapter-axum = { git = "https://github.com/stackpop/edgezero", rev = "35a72835322fe0127beb2ce998e4988e95974373", default-features = false } +edgezero-adapter-cloudflare = { git = "https://github.com/stackpop/edgezero", rev = "35a72835322fe0127beb2ce998e4988e95974373", default-features = false } +edgezero-adapter-fastly = { git = "https://github.com/stackpop/edgezero", rev = "35a72835322fe0127beb2ce998e4988e95974373", default-features = false } +edgezero-adapter-spin = { git = "https://github.com/stackpop/edgezero", rev = "35a72835322fe0127beb2ce998e4988e95974373", default-features = false } +edgezero-cli = { git = "https://github.com/stackpop/edgezero", rev = "35a72835322fe0127beb2ce998e4988e95974373" } +edgezero-core = { git = "https://github.com/stackpop/edgezero", rev = "35a72835322fe0127beb2ce998e4988e95974373", default-features = false } env_logger = "0.11" error-stack = "0.6" esi = "0.7.2" diff --git a/crates/trusted-server-adapter-fastly/Cargo.toml b/crates/trusted-server-adapter-fastly/Cargo.toml index 65320faa6..710523766 100644 --- a/crates/trusted-server-adapter-fastly/Cargo.toml +++ b/crates/trusted-server-adapter-fastly/Cargo.toml @@ -10,6 +10,13 @@ version = { workspace = true } [lints] workspace = true +[features] +# Opt-in reusable sandbox lifecycle. Off by default, and `fastly.toml`'s build +# command does not pass it, so production builds keep the single-request entry +# point. Enabling the feature alone does not enable reuse: the bounds still +# come from the runtime environment (see `sandbox::resolve_mode`). +reusable-sandbox = [] + [dependencies] async-trait = { workspace = true } base64 = { workspace = true } diff --git a/crates/trusted-server-adapter-fastly/src/app.rs b/crates/trusted-server-adapter-fastly/src/app.rs index b19653c9e..fbcd5735d 100644 --- a/crates/trusted-server-adapter-fastly/src/app.rs +++ b/crates/trusted-server-adapter-fastly/src/app.rs @@ -173,8 +173,10 @@ impl RuntimeStoreConfig { /// Application state built once per Wasm instance and shared for its lifetime. /// -/// In Fastly Compute each request spawns a new Wasm instance, so this struct is -/// effectively per-request. It holds pre-parsed settings and all service handles. +/// By default Fastly serves one request per instance. With the opt-in reusable +/// sandbox feature and complete runtime bounds, a successful application build +/// retains this state across requests until the sandbox retires. Settings and +/// shared services live here; request-specific state must remain request-local. pub(crate) struct AppState { pub(crate) settings: Arc, pub(crate) orchestrator: Arc, diff --git a/crates/trusted-server-adapter-fastly/src/logging.rs b/crates/trusted-server-adapter-fastly/src/logging.rs index 5fbeb176d..8515c7a66 100644 --- a/crates/trusted-server-adapter-fastly/src/logging.rs +++ b/crates/trusted-server-adapter-fastly/src/logging.rs @@ -74,11 +74,21 @@ fn detect_max_level(explicit: Option<&str>, hostname: Option<&str>) -> log::Leve /// used as-is, so `error`/`off` lowers it just as `debug` raises it; see /// [`resolve_max_level`]. /// -/// # Panics +/// # Errors /// -/// Panics if the Fastly logger cannot be built or if the global logger has already -/// been set. -pub(crate) fn init_logger() { +/// Returns a message if the Fastly logger cannot be built or if a global +/// logger is already installed. Returning rather than panicking lets +/// `Sandbox::setup_once` mark setup complete only on success, which *allows* +/// a later callback to retry. It does not by itself make retrying safe: +/// `setup_once` does not roll back partial side effects, so repeating this +/// function must stay harmless. +/// +/// It is, for the two ways it can fail. A failed `Logger::builder().build()` +/// installs nothing. A failed `fern::Dispatch::apply()` means a global logger +/// is already installed, so the retry fails the same way and never installs a +/// second one; the only repeated side effect is constructing a `Logger` that +/// is then dropped. Neither path leaves the process partially configured. +pub(crate) fn init_logger() -> Result<(), String> { let explicit = std::env::var(LOG_LEVEL_ENV).ok(); let hostname = std::env::var(LOCAL_HOSTNAME_ENV).ok(); let max_level = detect_max_level(explicit.as_deref(), hostname.as_deref()); @@ -88,7 +98,7 @@ pub(crate) fn init_logger() { .echo_stdout(true) .max_level(max_level) .build() - .expect("should build Logger"); + .map_err(|e| format!("should build Logger: {e}"))?; fern::Dispatch::new() .level(max_level) @@ -103,7 +113,7 @@ pub(crate) fn init_logger() { }) .chain(Box::new(logger) as Box) .apply() - .expect("should initialize logger"); + .map_err(|e| format!("should initialize logger: {e}")) } #[cfg(test)] diff --git a/crates/trusted-server-adapter-fastly/src/main.rs b/crates/trusted-server-adapter-fastly/src/main.rs index ce7264e78..b9f3b86be 100644 --- a/crates/trusted-server-adapter-fastly/src/main.rs +++ b/crates/trusted-server-adapter-fastly/src/main.rs @@ -13,7 +13,7 @@ use error_stack::Report; use fastly::http::Method as FastlyMethod; use fastly::{Request as FastlyRequest, Response as FastlyResponse}; -use trusted_server_core::cache_policy::EdgeCacheHeader; +use trusted_server_core::cache_policy::{EdgeCacheHeader, cache_control_headers_have_directive}; use trusted_server_core::ec::device::DeviceSignals; use trusted_server_core::ec::finalize::ec_finalize_response; use trusted_server_core::ec::kv::KvIdentityGraph; @@ -38,16 +38,23 @@ mod management_api; mod middleware; mod platform; mod rate_limiter; +mod sandbox; mod template_cache; mod tinybird; use crate::app::{ - EcFinalizeState, RuntimeStoreConfig, TrustedServerApp, load_settings_from_config_store, + AppState, EcFinalizeState, RuntimeStoreConfig, TrustedServerApp, + load_settings_from_config_store, }; use crate::ec_kv::FastlyEcKvStore; use crate::middleware::{HEADER_X_TS_FINALIZED, apply_finalize_headers, resolve_geo_for_response}; use crate::platform::{FastlyPlatformGeo, client_info_from_request}; use crate::rate_limiter::{FastlyRateLimiter, RATE_COUNTER_NAME}; +use crate::sandbox::{RetainedApp, Sandbox, SandboxCounters, ServeMode, StartupDiagnostics}; +// Only the reuse path builds a serving loop, so the retirement snapshots have +// no consumer in the default build. +#[cfg(feature = "reusable-sandbox")] +use crate::sandbox::RetirementCounters; /// Opens the Fastly Config Store used by the `EdgeZero` dispatcher. /// @@ -73,9 +80,89 @@ fn health_response(req: &FastlyRequest) -> Option { /// /// Uses an undecorated `main()` with `FastlyRequest::from_client()` instead of /// `#[fastly::main]` so the `EdgeZero` streaming publisher path can call -/// [`fastly::Response::stream_to_client`] explicitly. +/// [`fastly::Response::stream_to_client`] explicitly. It owns the sandbox +/// lifecycle and delegates each request to [`handle_request`]. +/// +/// Without the `reusable-sandbox` feature this takes one request and returns, +/// which is the original behaviour. With the feature, the bounds still have to +/// come from the runtime environment before the SDK serving loop is entered; +/// an unconfigured or partially configured sandbox stays single-request. fn main() { - let req = FastlyRequest::from_client(); + let (mode, diagnostics) = serve_mode(); + + // Held outside the sandbox: `serve_custom` owns the `Sandbox` and exposes + // no slot for application state that is not the retained payload. + let mut startup = StartupDiagnostics::default(); + for message in diagnostics { + startup.push(message); + } + + match mode { + ServeMode::Single => { + // `run_custom` completes the callback's SDK result exactly once. + // The callback sends its own response and returns `()`, so there + // is no error for the SDK to turn into a second response. + let Ok(()) = edgezero_adapter_fastly::lifecycle::run_custom( + FastlyRequest::from_client(), + |request, sandbox: &mut Sandbox| handle_request(request, sandbox, &mut startup), + ); + } + ServeMode::Reuse(limits) => serve_loop(limits, startup), + } +} + +/// Serves requests from one sandbox under explicit bounds. +/// +/// Only compiled with the `reusable-sandbox` feature; [`serve_mode`] can never +/// return [`ServeMode::Reuse`] without it. +#[cfg(feature = "reusable-sandbox")] +fn serve_loop(limits: crate::sandbox::SandboxLimits, mut startup: StartupDiagnostics) { + // `serve_custom` owns the `Sandbox` and drops it when serving ends, so the + // retirement line reports snapshots taken while it was still borrowed. + // `RetirementCounters` reads the attempt count on the way out, so a build + // performed by the final callback is included. + let mut counters = RetirementCounters::default(); + + let summary = edgezero_adapter_fastly::lifecycle::serve_custom( + fastly::http::serve::Serve::new() + .with_max_requests(limits.max_requests) + .with_max_memory(limits.max_memory_mib) + .with_max_lifetime(limits.max_lifetime) + .with_timeout(limits.timeout), + |request, sandbox: &mut Sandbox| { + counters.observe(sandbox, |sandbox| { + handle_request(request, sandbox, &mut startup); + }); + }, + ); + + log::info!( + "sandbox retiring after {} attempted callback(s), {} observed, {} build attempt(s)", + summary.requests(), + counters.requests(), + counters.attempts() + ); +} + +#[cfg(not(feature = "reusable-sandbox"))] +fn serve_loop(_limits: crate::sandbox::SandboxLimits, _startup: StartupDiagnostics) { + unreachable!("serve_mode never selects reuse without the reusable-sandbox feature") +} + +/// Handles one request end to end, sending its own response. +/// +/// Returns `()` rather than a response so the streaming publisher path can +/// call [`fastly::Response::stream_to_client`] itself. The SDK's +/// `HandlerResult` impl for `()` treats that as already sent. +/// +/// Every non-panicking path through this function must send exactly once. In a +/// reused sandbox the SDK refuses to wait for the next request until the +/// current one is complete, so a missed send stalls the loop rather than +/// merely dropping one response. +fn handle_request(req: FastlyRequest, sandbox: &mut Sandbox, startup: &mut StartupDiagnostics) { + // The framework counts the callback before invoking it, including early + // returns, so this is already this request's 1-based ordinal. + let ordinal = sandbox.requests(); // Health probe bypasses logging, settings, and app construction as a cheap liveness signal. if let Some(response) = health_response(&req) { @@ -83,15 +170,179 @@ fn main() { return; } - logging::init_logger(); - edgezero_main(req); + // Marked complete only once installation succeeds, so a failed install is + // retried on a later callback. `setup_once` rolls nothing back, so that is + // only correct because `init_logger` is harmless to repeat — see its docs. + match sandbox.setup_once(crate::logging::init_logger) { + Ok(()) => startup.flush(), + Err(error) => { + // Logger installation failed, so its own error cannot rely on log. + // Keep startup diagnostics pending until installation succeeds. + #[allow(clippy::print_stderr, reason = "logger installation failed")] + { + eprintln!("logger installation failed, retrying next callback: {error}"); + } + } + } + + // Correlation is request-local and never retained. `FASTLY_TRACE_ID` names + // the sandbox, not the request, so it is not used here. The id rides on the + // response only when metrics are enabled; the client request is left + // untouched so nothing new reaches origin in the default configuration. + let request_id = req + .get_client_request_id() + .map(str::to_owned) + .unwrap_or_else(|| format!("{}-{ordinal}", instance_id())); + + edgezero_main(req, sandbox, ordinal, &request_id); +} + +/// Builds the counters snapshot response. +/// +/// A snapshot only. It reports the sandbox that served *this* probe, which is +/// not necessarily the sandbox that served any preceding workload request, so +/// reuse is established from the counters attached to workload responses +/// rather than from polling this. +/// +/// `ordinal` is this probe's own 1-based position in the sandbox, which is +/// also the number of callbacks the sandbox has served including this one. +/// The lifetime count is therefore not reported separately. +#[cfg(feature = "reusable-sandbox")] +fn sandbox_metrics_response(sandbox: &Sandbox, ordinal: u64) -> FastlyResponse { + let body = serde_json::json!({ + "instance": instance_id(), + "ordinal": ordinal, + "builds": sandbox.initialization_attempts(), + }); + + FastlyResponse::from_status(fastly::http::StatusCode::OK) + .with_header("cache-control", "private, no-store") + .with_body_json(&body) + .unwrap_or_else(|e| { + log::error!("failed to serialize sandbox metrics: {e}"); + FastlyResponse::from_status(fastly::http::StatusCode::INTERNAL_SERVER_ERROR) + }) +} + +/// Attaches sandbox counters only to private, no-store workload responses. +/// +/// Called before headers are committed, which on the streaming path means +/// before `stream_to_client`. The counters therefore describe the request up +/// to commitment and cannot report its eventual outcome; a failure after +/// commitment is recorded in logs instead and reconciled during analysis. +fn attach_sandbox_counters(response: &mut HttpResponse, counters: &SandboxCounters) { + // Per-request measurements must never be replayed from a cache. Requiring + // no-store also excludes browser-cacheable and field-qualified privacy. + if !cache_control_headers_have_directive(response.headers(), "private") + || !cache_control_headers_have_directive(response.headers(), "no-store") + { + return; + } + + let headers = response.headers_mut(); + for (name, value) in [ + (sandbox::HEADER_SANDBOX_INSTANCE, instance_id()), + ( + sandbox::HEADER_SANDBOX_ORDINAL, + counters.ordinal.to_string(), + ), + (sandbox::HEADER_SANDBOX_BUILDS, counters.builds.to_string()), + ( + sandbox::HEADER_SANDBOX_REQUEST_ID, + counters.request_id.clone(), + ), + (sandbox::HEADER_SANDBOX_VCPU_MS, vcpu_ms()), + (sandbox::HEADER_SANDBOX_HEAP_MIB, heap_mib()), + ] { + match edgezero_core::http::HeaderValue::from_str(&value) { + Ok(value) => { + headers.insert(name, value); + } + Err(e) => log::warn!("sandbox counter `{name}` is not a valid header value: {e}"), + } + } +} + +/// Resolves how this sandbox will serve requests. +/// +/// Without the `reusable-sandbox` feature this is unconditionally +/// [`ServeMode::Single`] and reads nothing, so the default build does no +/// startup work the original entry point did not do. +/// Returns the mode alongside diagnostics that must wait for the logger. +/// +/// This runs before any logger exists, so the reasons reuse was declined are +/// carried out rather than logged here, where they would be discarded. +#[cfg(feature = "reusable-sandbox")] +fn serve_mode() -> (ServeMode, Vec) { + let (raw, diagnostics) = crate::sandbox::read_raw_limits(); + (crate::sandbox::resolve_mode(raw), diagnostics) +} + +#[cfg(not(feature = "reusable-sandbox"))] +fn serve_mode() -> (ServeMode, Vec) { + (ServeMode::Single, Vec::new()) +} + +/// Guest-instance identifier used to attribute requests to a sandbox. +/// +/// `FASTLY_TRACE_ID` describes the sandbox, which is exactly what is wanted +/// here and exactly why it must not be used as a request id. An absent value +/// is reported rather than synthesized, so a measurement run cannot silently +/// claim reuse it never observed. +fn instance_id() -> String { + // The SDK documents this as the per-sandbox identifier; on wasm32-wasip1 + // it resolves to `FASTLY_TRACE_ID`, which is why that value must never be + // used as a request id. + let id = fastly::compute_runtime::sandbox_id(); + if id.is_empty() { + return sandbox::INSTANCE_ID_UNAVAILABLE.to_owned(); + } + id.to_owned() +} + +/// Cumulative guest vCPU milliseconds, or a marker when unsupported. +fn vcpu_ms() -> String { + fastly::compute_runtime::elapsed_vcpu_ms().map_or_else( + |_| sandbox::COUNTER_UNSUPPORTED.to_owned(), + |v| v.to_string(), + ) +} + +/// Guest heap snapshot in MiB, or a marker when unsupported. +fn heap_mib() -> String { + fastly::compute_runtime::heap_memory_snapshot_mib().map_or_else( + |_| sandbox::COUNTER_UNSUPPORTED.to_owned(), + |v| v.to_string(), + ) } /// Handles a request through the `EdgeZero` router path. -fn edgezero_main(mut req: FastlyRequest) { +fn edgezero_main(mut req: FastlyRequest, sandbox: &mut Sandbox, ordinal: u64, request_id: &str) { let runtime_env = runtime_env_config(TrustedServerApp::stores()); let runtime_stores = RuntimeStoreConfig::from_env(&runtime_env); + // Short-circuit the sandbox counters probe before app construction. It must + // not build the application: polling it would otherwise increment the very + // build counter it reports. + #[cfg(feature = "reusable-sandbox")] + if req.get_method() == FastlyMethod::GET && req.get_path() == sandbox::SANDBOX_METRICS_PATH { + match load_settings_from_config_store(&runtime_stores) { + Ok(settings) if sandbox::metrics_enabled(&settings) => { + sandbox_metrics_response(sandbox, ordinal).send_to_client(); + } + Ok(_) => { + FastlyResponse::from_status(fastly::http::StatusCode::NOT_FOUND).send_to_client(); + } + Err(e) => { + log::warn!("sandbox metrics endpoint: failed to load settings: {e:?}"); + FastlyResponse::from_status(fastly::http::StatusCode::INTERNAL_SERVER_ERROR) + .with_body_text_plain("Internal Server Error") + .send_to_client(); + } + } + return; + } + // Short-circuit the JA4 debug probe before app construction. Must run here // because TLS/JA4 accessors are only available on FastlyRequest before // conversion to edgezero types. @@ -125,8 +376,40 @@ fn edgezero_main(mut req: FastlyRequest) { } }; - let (app, app_state) = TrustedServerApp::build_app_with_state(&runtime_stores); + // Build lazily, once per sandbox. Reached only past the health, JA4, and + // counters short-circuits, so none of those pays for construction. + // + // `initialize` builds only when the sandbox is empty, retains only + // success, and returns the error unchanged. A failed build hands back its + // error router as the error payload: that serves this request and is then + // dropped, so a transient config-store failure cannot pin the sandbox into + // permanent error mode, and the next callback retries construction. + let failed_build = sandbox + .initialize(|| { + let (app, state) = TrustedServerApp::build_app_with_state(&runtime_stores); + match state { + Some(state) => Ok(RetainedApp { app, state }), + None => Err(app), + } + }) + .err(); + + let (app, app_state): (&edgezero_core::app::App, Option>) = + match (failed_build.as_ref(), sandbox.state()) { + (Some(app), _) => (app, None), + (None, Some(retained)) => (&retained.app, Some(Arc::clone(&retained.state))), + (None, None) => { + log::error!("no application available after initialization"); + FastlyResponse::from_status(fastly::http::StatusCode::INTERNAL_SERVER_ERROR) + .with_body_text_plain("Internal Server Error") + .send_to_client(); + return; + } + }; + let settings_snapshot = app_state.as_ref().map(|state| Arc::clone(&state.settings)); + let counters = + SandboxCounters::capture(sandbox, ordinal, request_id, settings_snapshot.as_deref()); let trusted_client_ip = settings_snapshot .as_deref() .and_then(|settings| settings.trusted_client_ip.as_ref()); @@ -219,7 +502,11 @@ fn edgezero_main(mut req: FastlyRequest) { if let Some(settings) = settings_snapshot.as_deref() { match apply_edgezero_ec_finalize(settings, &mut ec_state, &mut response) { Ok(partner_registry) => { - send_edgezero_response(response, request_filter_effects.as_ref()); + send_edgezero_response( + response, + request_filter_effects.as_ref(), + counters.as_ref(), + ); run_edgezero_pull_sync_after_send(settings, &partner_registry, &ec_state); return; } @@ -234,7 +521,11 @@ fn edgezero_main(mut req: FastlyRequest) { Ok(settings) => { match apply_edgezero_ec_finalize(&settings, &mut ec_state, &mut response) { Ok(partner_registry) => { - send_edgezero_response(response, request_filter_effects.as_ref()); + send_edgezero_response( + response, + request_filter_effects.as_ref(), + counters.as_ref(), + ); run_edgezero_pull_sync_after_send( &settings, &partner_registry, @@ -256,7 +547,7 @@ fn edgezero_main(mut req: FastlyRequest) { } } - send_edgezero_response(response, request_filter_effects.as_ref()); + send_edgezero_response(response, request_filter_effects.as_ref(), counters.as_ref()); } fn edge_error_response(error: EdgeError) -> HttpResponse { @@ -370,9 +661,26 @@ where fn send_edgezero_response( mut response: HttpResponse, request_filter_effects: Option<&RequestFilterEffects>, + counters: Option<&SandboxCounters>, ) { apply_terminal_response_effects(&mut response, request_filter_effects); + // Captured before the body is consumed so post-commitment failures can be + // matched back to the response that carried these counters. + let counter_context = counters.map_or_else(String::new, |counters| { + format!( + " [instance={} ordinal={} request={}]", + instance_id(), + counters.ordinal, + counters.request_id + ) + }); + + // Before headers commit, including before `stream_to_client` below. + if let Some(counters) = counters.as_ref() { + attach_sandbox_counters(&mut response, counters); + } + let (parts, body) = response.into_parts(); match body { @@ -385,11 +693,19 @@ fn send_edgezero_response( match futures::executor::block_on(stream_asset_body(body, &mut streaming_body)) { Ok(()) => { if let Err(e) = streaming_body.finish() { - log::error!("failed to finish EdgeZero streaming body: {e}"); + // Also post-commitment: same attribution as the + // streaming failure below. + log::error!( + "failed to finish EdgeZero streaming body{counter_context}: {e}" + ); } } Err(e) => { - log::error!("EdgeZero streaming failed: {e:?}"); + // After commitment: log and stop. Returning an error here + // would let the SDK attempt a second response. Counters + // already went out with the headers, so the failure is + // tagged with the same identity for reconciliation. + log::error!("EdgeZero streaming failed{counter_context}: {e:?}"); drop(streaming_body); } } @@ -616,6 +932,73 @@ mod tests { ); } + #[test] + fn sandbox_counters_never_appear_on_cacheable_responses() { + let counters = SandboxCounters { + ordinal: 2, + builds: 1, + request_id: "request-example".to_owned(), + }; + for policy in [ + None, + Some("public, s-maxage=3600"), + Some("private, max-age=60"), + Some("private=\"set-cookie\""), + Some("public, extension=\"private, no-store\""), + ] { + let mut response = HttpResponse::new(EdgeBody::empty()); + if let Some(policy) = policy { + response.headers_mut().insert( + "cache-control", + HeaderValue::from_str(policy).expect("should encode cache policy"), + ); + } + apply_terminal_response_effects(&mut response, None); + let original_headers = response.headers().clone(); + + attach_sandbox_counters(&mut response, &counters); + + assert_eq!( + response.headers(), + &original_headers, + "should preserve cache policy and omit all counters for {policy:?}" + ); + } + } + + #[test] + fn sandbox_counters_follow_final_private_no_store_policy() { + let counters = SandboxCounters { + ordinal: 2, + builds: 1, + request_id: "request-example".to_owned(), + }; + let mut response = response_builder() + .header("cache-control", "public, s-maxage=3600") + .body(EdgeBody::empty()) + .expect("should build response"); + response.extensions_mut().insert(TerminalPrivateResponse); + apply_terminal_response_effects(&mut response, None); + + attach_sandbox_counters(&mut response, &counters); + + assert_eq!( + response.headers()[sandbox::HEADER_SANDBOX_REQUEST_ID], + "request-example", + "should identify this uncached request" + ); + assert_eq!( + response.headers()[sandbox::HEADER_SANDBOX_ORDINAL], + "2", + "should report the current ordinal" + ); + assert_eq!( + response.headers()["cache-control"], + "no-store, private", + "should retain terminal privacy" + ); + } + #[test] fn late_filter_effects_cannot_make_an_assembled_response_public() { let mut response = response_builder() diff --git a/crates/trusted-server-adapter-fastly/src/sandbox.rs b/crates/trusted-server-adapter-fastly/src/sandbox.rs new file mode 100644 index 000000000..228106047 --- /dev/null +++ b/crates/trusted-server-adapter-fastly/src/sandbox.rs @@ -0,0 +1,962 @@ +//! Reusable-sandbox lifecycle state and limit resolution. +//! +//! Fastly Compute normally starts a fresh Wasm sandbox per request. SDK 0.12.1 +//! exposes [`fastly::http::serve::Serve`], which lets one sandbox serve +//! several. This module owns the opt-in decision, the bounds, and the +//! per-sandbox bookkeeping the entry point carries across requests. +//! +//! [`Sandbox`] owns the state the entry point carries across requests: the +//! logger guard, the measurement counters, and the retained application. + +use std::sync::Arc; +use std::time::Duration; + +use edgezero_core::app::App; +use trusted_server_core::settings::Settings; + +use crate::app::AppState; + +/// Header carrying the guest-instance identifier. +pub(crate) const HEADER_SANDBOX_INSTANCE: &str = "x-ts-sandbox-instance"; + +/// Header carrying the request ordinal within the current sandbox, 1-based. +pub(crate) const HEADER_SANDBOX_ORDINAL: &str = "x-ts-sandbox-ordinal"; + +/// Header carrying the number of application builds this sandbox has performed. +pub(crate) const HEADER_SANDBOX_BUILDS: &str = "x-ts-sandbox-builds"; + +/// Header carrying the per-request correlation id. +pub(crate) const HEADER_SANDBOX_REQUEST_ID: &str = "x-ts-sandbox-request-id"; + +/// Header carrying cumulative guest vCPU milliseconds, where supported. +pub(crate) const HEADER_SANDBOX_VCPU_MS: &str = "x-ts-sandbox-vcpu-ms"; + +/// Header carrying the guest heap snapshot in MiB, where supported. +pub(crate) const HEADER_SANDBOX_HEAP_MIB: &str = "x-ts-sandbox-heap-mib"; + +/// Value reported when a runtime counter is not supported by the host. +/// +/// Distinguished from a zero reading: unsupported is not the same as idle. +pub(crate) const COUNTER_UNSUPPORTED: &str = "unsupported"; + +/// Path of the counters snapshot endpoint. +/// +/// Gated on the feature alone, not `any(feature, test)`: its only caller is +/// the feature-gated short-circuit, so a `test` arm would make it dead code +/// in a feature-off test build. +#[cfg(feature = "reusable-sandbox")] +pub(crate) const SANDBOX_METRICS_PATH: &str = "/_ts/debug/sandbox"; + +/// Label used when the runtime reports no usable guest-instance identifier. +/// +/// Recorded rather than papered over: a measurement run that sees this value +/// has no instance identity and cannot claim observed reuse. +pub(crate) const INSTANCE_ID_UNAVAILABLE: &str = "unavailable"; + +/// Per-sandbox state, owned by `edgezero_adapter_fastly::lifecycle::Sandbox`. +/// +/// The framework owns lazy successful-only retention, the callback count, the +/// initialization-attempt count, and the one-time setup guard. This alias +/// names the application-specific payload it retains. +pub(crate) type Sandbox = edgezero_adapter_fastly::lifecycle::Sandbox; + +/// The application state a retained sandbox carries. +/// +/// Only ever reachable after a successful build: a failed build returns the +/// error router as the `initialize` error instead, so it is served for the +/// current request and dropped rather than retained. +pub(crate) struct RetainedApp { + pub(crate) app: App, + pub(crate) state: Arc, +} + +/// Startup diagnostics held until a logger exists. +/// +/// Limit resolution runs in `main`, before any logger is installed, so its +/// messages are collected here and flushed by the first callback that +/// installs logging. This is callback-local state rather than sandbox state: +/// `serve_custom` owns the `Sandbox` and exposes no slot for it. +#[derive(Debug, Default)] +pub(crate) struct StartupDiagnostics(Vec); + +impl StartupDiagnostics { + /// Records a message for emission once a logger exists. + pub(crate) fn push(&mut self, message: impl Into) { + self.0.push(message.into()); + } + + /// Emits and clears everything held so far. + pub(crate) fn flush(&mut self) { + for message in self.0.drain(..) { + log::info!("{message}"); + } + } + + /// Whether anything is still waiting to be emitted. + #[cfg(test)] + pub(crate) fn pending(&self) -> usize { + self.0.len() + } +} + +/// Snapshots of the framework's counters, taken around one callback. +/// +/// `serve_custom` owns the [`Sandbox`] and drops it when serving ends, so the +/// retirement line cannot read it afterwards. These are snapshots of +/// `EdgeZero`'s counters taken while the sandbox is still borrowed; nothing here +/// increments anything. +/// +/// The two are read at different points on purpose. The framework increments +/// its callback count *before* invoking the callback, so `requests` is correct +/// on entry. Initialization happens *during* the callback, so `attempts` must +/// be read on the way out. The count itself is never lost — it lives in the +/// sandbox until serving ends — but a snapshot taken on entry is stale by the +/// time the final callback finishes, so a build it performed goes unreported. +#[cfg(any(feature = "reusable-sandbox", test))] +#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)] +pub(crate) struct RetirementCounters { + requests: u64, + attempts: u64, +} + +#[cfg(any(feature = "reusable-sandbox", test))] +impl RetirementCounters { + /// Runs one callback against `sandbox`, capturing counters around it. + pub(crate) fn observe( + &mut self, + sandbox: &mut Sandbox, + callback: impl FnOnce(&mut Sandbox) -> R, + ) -> R { + self.requests = sandbox.requests(); + let outcome = callback(sandbox); + self.attempts = sandbox.initialization_attempts(); + outcome + } + + /// Callback count observed on entry to the last callback. + pub(crate) fn requests(&self) -> u64 { + self.requests + } + + /// Initialization attempts observed on exit from the last callback. + pub(crate) fn attempts(&self) -> u64 { + self.attempts + } +} + +/// Counter values captured for one response. +/// +/// An owned snapshot rather than a borrow of [`Sandbox`], so attaching counters +/// never competes with the mutable borrow the request path holds. +#[derive(Debug, Clone, PartialEq, Eq)] +pub(crate) struct SandboxCounters { + pub(crate) ordinal: u64, + pub(crate) builds: u64, + pub(crate) request_id: String, +} + +impl SandboxCounters { + /// Captures the current counters, or `None` when metrics are disabled. + /// + /// Returning `None` is what keeps the counters off every response in the + /// default configuration, including the correlation id. + pub(crate) fn capture( + sandbox: &Sandbox, + ordinal: u64, + request_id: &str, + settings: Option<&Settings>, + ) -> Option { + settings + .filter(|settings| metrics_enabled(settings)) + .map(|_| Self { + ordinal, + builds: sandbox.initialization_attempts(), + request_id: request_id.to_owned(), + }) + } +} + +/// Whether the counters endpoint and response counters are enabled. +pub(crate) fn metrics_enabled(settings: &Settings) -> bool { + settings.debug.sandbox_metrics_enabled +} + +/// Resolved sandbox limits, already normalized for the SDK. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(crate) struct SandboxLimits { + pub(crate) max_requests: usize, + pub(crate) max_memory_mib: u32, + pub(crate) max_lifetime: Duration, + pub(crate) timeout: Duration, +} + +/// How this sandbox will serve requests. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(crate) enum ServeMode { + /// Handle exactly one request, as a non-reusable sandbox does today. + Single, + /// Enter the SDK serving loop under the given bounds. + #[cfg_attr( + not(any(feature = "reusable-sandbox", test)), + expect(dead_code, reason = "constructed only on the reuse path") + )] + Reuse(SandboxLimits), +} + +/// Raw limit values as read from the runtime environment. +#[cfg(any(feature = "reusable-sandbox", test))] +#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)] +pub(crate) struct RawLimits { + pub(crate) max_requests: Option, + pub(crate) max_memory_mib: Option, + pub(crate) max_lifetime_ms: Option, + pub(crate) timeout_ms: Option, +} + +/// Resolves raw configuration into a serving mode. +/// +/// Reuse requires all four bounds. A bare request limit is refused because the +/// SDK's omitted lifetime and wait timeout both default to [`Duration::MAX`], +/// which would leave a sandbox waiting without a bound. An application-level +/// request limit of `0` normalizes to `1`, because the SDK reads +/// `with_max_requests(0)` as unlimited — the opposite of what an operator +/// writing `0` intends. +#[cfg(any(feature = "reusable-sandbox", test))] +pub(crate) fn resolve_mode(raw: RawLimits) -> ServeMode { + let (Some(max_requests), Some(max_lifetime_ms), Some(timeout_ms), Some(max_memory_mib)) = ( + raw.max_requests, + raw.max_lifetime_ms, + raw.timeout_ms, + raw.max_memory_mib, + ) else { + return ServeMode::Single; + }; + + let max_requests = max_requests.max(1); + if max_requests <= 1 || max_lifetime_ms == 0 || timeout_ms == 0 { + return ServeMode::Single; + } + + let Ok(max_requests) = usize::try_from(max_requests) else { + return ServeMode::Single; + }; + + let Ok(max_memory_mib) = u32::try_from(max_memory_mib) else { + return ServeMode::Single; + }; + // The SDK uses u32::MAX for unsupported heap snapshots. Keep the bound + // below that sentinel so an unavailable reading always retires the guest. + if max_memory_mib == 0 || max_memory_mib == u32::MAX { + return ServeMode::Single; + } + + ServeMode::Reuse(SandboxLimits { + max_requests, + max_memory_mib, + max_lifetime: Duration::from_millis(max_lifetime_ms), + timeout: Duration::from_millis(timeout_ms), + }) +} + +/// Suffixes of the service-scoped runtime-environment keys holding the bounds. +#[cfg(any(feature = "reusable-sandbox", test))] +const KEY_MAX_REQUESTS: &str = "TS__SANDBOX__MAX_REQUESTS"; +#[cfg(any(feature = "reusable-sandbox", test))] +const KEY_MAX_LIFETIME_MS: &str = "TS__SANDBOX__MAX_LIFETIME_MS"; +#[cfg(any(feature = "reusable-sandbox", test))] +const KEY_TIMEOUT_MS: &str = "TS__SANDBOX__TIMEOUT_MS"; +#[cfg(any(feature = "reusable-sandbox", test))] +const KEY_MAX_MEMORY_MIB: &str = "TS__SANDBOX__MAX_MEMORY_MIB"; + +/// Builds the service-scoped runtime-environment key for a bound. +/// +/// `EdgeZero`'s own `service_scoped_runtime_env_key` is private, so the shape is +/// reproduced here. It must stay identical to the one `edgezero provision` +/// writes. +/// +/// Follow-up: a public `EdgeZero` key-construction or lookup helper would let +/// this duplication go. Until such an API exists this implementation stays, so +/// the key shape has exactly one definition on our side. The `TS__SANDBOX__*` +/// suffixes and the limit-validation policy in [`resolve_mode`] are +/// application-owned either way and would not move. +#[cfg(any(feature = "reusable-sandbox", test))] +fn scoped_key(service_id: &str, suffix: &str) -> String { + format!("EDGEZERO__SERVICES__{service_id}__{suffix}") +} + +/// Collects raw limits from a fallible key lookup. +/// +/// Separated from the config-store call so the failure modes are testable +/// without a store: a lookup error, an absent key, and an unparseable value +/// must all degrade to "absent" rather than propagate. +/// +/// Diagnostics are returned rather than logged. This runs before the logger +/// exists, so anything logged here would be discarded. +#[cfg(any(feature = "reusable-sandbox", test))] +fn collect_raw_limits(service_id: &str, mut lookup: F) -> (RawLimits, Vec) +where + E: core::fmt::Display, + F: FnMut(&str) -> Result, E>, +{ + let mut diagnostics = Vec::new(); + + let mut read = |suffix: &str| -> Option { + let key = scoped_key(service_id, suffix); + match lookup(&key) { + Ok(Some(raw)) => match raw.trim().parse::() { + Ok(value) => Some(value), + Err(e) => { + diagnostics.push(format!( + "sandbox limit `{key}` is not a non-negative integer ({e}); ignoring" + )); + None + } + }, + Ok(None) => None, + Err(e) => { + diagnostics.push(format!( + "sandbox limit `{key}` lookup failed ({e}); ignoring" + )); + None + } + } + }; + + let limits = RawLimits { + max_requests: read(KEY_MAX_REQUESTS), + max_memory_mib: read(KEY_MAX_MEMORY_MIB), + max_lifetime_ms: read(KEY_MAX_LIFETIME_MS), + timeout_ms: read(KEY_TIMEOUT_MS), + }; + + (limits, diagnostics) +} + +/// Reads the sandbox bounds from the runtime-environment config store. +/// +/// These keys deliberately bypass `edgezero_adapter_fastly::runtime_env_config`. +/// That helper resolves a closed allowlist — adapter host and port, logging +/// settings, and per-store `__NAME`/`__KEY` selectors — and silently drops +/// everything else, so a sandbox key routed through it would always read as +/// absent and reuse would never engage. Still true at the pinned revision. +/// +/// Uses [`fastly::ConfigStore::try_get`], never `get`: `get` panics on a +/// lookup error, and this runs in `main` before the health probe, so a panic +/// here would take down liveness rather than merely disabling reuse. +/// +/// Any failure resolves to [`RawLimits::default`], and hence to +/// [`ServeMode::Single`]: an absent or unopenable store, an empty service id, +/// a failed lookup, or a value that does not parse as a `u64`. +#[cfg(feature = "reusable-sandbox")] +pub(crate) fn read_raw_limits() -> (RawLimits, Vec) { + use edgezero_adapter_fastly::RUNTIME_ENV_STORE_NAME; + + let Ok(store) = fastly::ConfigStore::try_open(RUNTIME_ENV_STORE_NAME) else { + return ( + RawLimits::default(), + vec![format!( + "sandbox reuse disabled: config store `{RUNTIME_ENV_STORE_NAME}` unavailable" + )], + ); + }; + + // Viceroy reports a service id of twenty-two zeros, which is a valid + // scope for key construction. Only an empty id is unusable. + let service_id = fastly::compute_runtime::service_id(); + if service_id.is_empty() { + return ( + RawLimits::default(), + vec!["sandbox reuse disabled: no service id available for key scoping".to_owned()], + ); + } + + collect_raw_limits(service_id, |key| store.try_get(key)) +} + +#[cfg(test)] +mod tests { + use std::cell::Cell; + + use super::*; + + #[test] + fn scoped_key_matches_the_edgezero_shape() { + // Viceroy reports twenty-two zeros locally; it is a valid scope. + let service_id = "0000000000000000000000"; + + assert_eq!( + scoped_key(service_id, KEY_MAX_REQUESTS), + "EDGEZERO__SERVICES__0000000000000000000000__TS__SANDBOX__MAX_REQUESTS", + "should reproduce the service-scoped key edgezero writes" + ); + assert_eq!( + scoped_key(service_id, KEY_MAX_LIFETIME_MS), + "EDGEZERO__SERVICES__0000000000000000000000__TS__SANDBOX__MAX_LIFETIME_MS", + "should scope the lifetime bound the same way" + ); + assert_eq!( + scoped_key(service_id, KEY_TIMEOUT_MS), + "EDGEZERO__SERVICES__0000000000000000000000__TS__SANDBOX__TIMEOUT_MS", + "should scope the wait timeout the same way" + ); + } + + fn full(max_requests: u64) -> RawLimits { + RawLimits { + max_requests: Some(max_requests), + max_memory_mib: Some(128), + max_lifetime_ms: Some(30_000), + timeout_ms: Some(500), + } + } + + #[test] + fn absent_configuration_stays_single_request() { + assert_eq!( + resolve_mode(RawLimits::default()), + ServeMode::Single, + "an unconfigured sandbox should behave as it does today" + ); + } + + #[test] + fn zero_request_limit_normalizes_to_single_not_unlimited() { + assert_eq!( + resolve_mode(full(0)), + ServeMode::Single, + "0 should mean one request, never the SDK's unlimited" + ); + } + + #[test] + fn one_request_limit_stays_single_request() { + assert_eq!( + resolve_mode(full(1)), + ServeMode::Single, + "a limit of 1 should not enter the serving loop" + ); + } + + #[test] + fn partial_configuration_refuses_to_reuse() { + let raw = RawLimits { + max_requests: Some(10), + max_memory_mib: Some(128), + max_lifetime_ms: None, + timeout_ms: Some(500), + }; + + assert_eq!( + resolve_mode(raw), + ServeMode::Single, + "a missing bound should not fall back to the SDK's unbounded default" + ); + } + + #[test] + fn zero_lifetime_or_timeout_refuses_to_reuse() { + assert_eq!( + resolve_mode(RawLimits { + max_lifetime_ms: Some(0), + ..full(10) + }), + ServeMode::Single, + "a zero lifetime should not mean unlimited" + ); + assert_eq!( + resolve_mode(RawLimits { + timeout_ms: Some(0), + ..full(10) + }), + ServeMode::Single, + "a zero wait timeout should not mean unlimited" + ); + } + + #[test] + fn complete_configuration_enables_reuse() { + assert_eq!( + resolve_mode(full(10)), + ServeMode::Reuse(SandboxLimits { + max_requests: 10, + max_memory_mib: 128, + max_lifetime: Duration::from_millis(30_000), + timeout: Duration::from_millis(500), + }), + "all four bounds present should enter the serving loop" + ); + } + + #[test] + fn setup_runs_once_on_success_and_retries_after_failure() { + let mut sandbox = Sandbox::default(); + let attempts = Cell::new(0_u32); + + // A failed install must leave setup eligible for retry. + let failed = sandbox.setup_once(|| { + attempts.set(attempts.get() + 1); + Err::<(), &str>("install failed") + }); + assert_eq!( + failed.err(), + Some("install failed"), + "the error should surface" + ); + + // The retry runs, because setup was never marked complete. + sandbox + .setup_once(|| { + attempts.set(attempts.get() + 1); + Ok::<(), &str>(()) + }) + .expect("the retry should succeed"); + assert_eq!(attempts.get(), 2, "a failed install should be retried"); + + // Once successful, it never runs again. Reinstalling the global logger + // panics inside `fern`, which is the failure this guards. + sandbox + .setup_once(|| -> Result<(), &str> { + attempts.set(attempts.get() + 1); + panic!("setup must not run again after success") + }) + .expect("a completed setup should be skipped"); + assert_eq!(attempts.get(), 2, "setup should stay complete"); + } + + /// Minimal settings sufficient to build real application state. + fn test_settings() -> Settings { + Settings::from_toml( + r#" + [[handlers]] + path = "^/_ts/admin" + username = "admin" + password = "admin-pass" + + [publisher] + domain = "test-publisher.example" + cookie_domain = ".test-publisher.example" + origin_url = "https://origin.test-publisher.example" + proxy_secret = "unit-test-proxy-secret" + + [ec] + passphrase = "test-secret-key-32-bytes-minimum" + "#, + ) + .expect("should parse sandbox test settings") + } + + /// A real `AppState`, so retention is exercised against the type the + /// entry point actually keeps rather than a stand-in. + fn test_state() -> Arc { + crate::app::build_state_from_settings(test_settings()) + .expect("should build sandbox test state") + } + + /// A real `RetainedApp`, so retention is exercised against the type the + /// entry point actually keeps rather than a stand-in. + fn test_retained() -> RetainedApp { + RetainedApp { + app: test_app(), + state: test_state(), + } + } + + /// An empty but real `App`. + fn test_app() -> App { + App::with_name( + edgezero_core::router::RouterService::builder().build(), + "sandbox-test", + ) + } + + /// Stand-in for `fastly::config_store::LookupError`, which cannot be + /// constructed outside the SDK. + #[derive(Debug, derive_more::Display)] + #[display("simulated lookup failure")] + struct LookupFailed; + + #[test] + fn a_failed_lookup_degrades_to_single_request_instead_of_panicking() { + // `ConfigStore::get` panics on a lookup error and runs before the + // health probe, so the fallible path must absorb the error. + let (limits, diagnostics) = + collect_raw_limits::("svc", |_key| Err(LookupFailed)); + + assert_eq!( + limits, + RawLimits::default(), + "a failed lookup should read as absent" + ); + assert_eq!( + resolve_mode(limits), + ServeMode::Single, + "a failed lookup should leave the sandbox single-request" + ); + assert_eq!( + diagnostics.len(), + 4, + "each failed key should report why reuse was declined" + ); + assert!( + diagnostics[0].contains("lookup failed"), + "the diagnostic should name the failure, got: {}", + diagnostics[0] + ); + } + + #[test] + fn an_unparseable_value_degrades_to_single_request() { + let (limits, diagnostics) = collect_raw_limits::("svc", |key| { + Ok(Some(if key.ends_with(KEY_MAX_REQUESTS) { + "ten".to_owned() + } else { + "500".to_owned() + })) + }); + + assert_eq!( + limits.max_requests, None, + "an unparseable limit should read as absent" + ); + assert_eq!( + resolve_mode(limits), + ServeMode::Single, + "an unparseable request limit should not enable reuse" + ); + assert!( + diagnostics + .iter() + .any(|d| d.contains("not a non-negative integer")), + "should explain the rejected value, got: {diagnostics:?}" + ); + } + + #[test] + fn positive_memory_limits_are_preserved_without_truncation() { + for memory in [1, 128, u32::MAX - 1] { + let mode = resolve_mode(RawLimits { + max_memory_mib: Some(u64::from(memory)), + ..full(10) + }); + assert!( + matches!(mode, ServeMode::Reuse(limits) if limits.max_memory_mib == memory), + "should pass a valid memory bound unchanged to the SDK" + ); + } + } + + #[test] + fn missing_or_invalid_memory_limit_refuses_reuse() { + for memory in [ + None, + Some("0"), + Some("-1"), + Some("many"), + Some("4294967295"), + Some("4294967296"), + ] { + let (limits, _) = collect_raw_limits::("example-service", |key| { + Ok(if key.ends_with("TS__SANDBOX__MAX_MEMORY_MIB") { + memory.map(str::to_owned) + } else { + Some("10".to_owned()) + }) + }); + + assert_eq!( + resolve_mode(limits), + ServeMode::Single, + "should refuse reuse with missing or invalid memory limit {memory:?}" + ); + } + } + + #[test] + fn absent_keys_report_nothing_and_stay_single_request() { + let (limits, diagnostics) = collect_raw_limits::("svc", |_key| Ok(None)); + + assert_eq!( + limits, + RawLimits::default(), + "absent keys should read absent" + ); + assert!( + diagnostics.is_empty(), + "an unconfigured sandbox is the normal case, not a diagnostic" + ); + } + + #[test] + fn a_complete_store_enables_reuse_end_to_end() { + let (limits, diagnostics) = collect_raw_limits::("svc", |key| { + Ok(Some( + if key.ends_with(KEY_MAX_REQUESTS) { + "10" + } else if key.ends_with(KEY_MAX_MEMORY_MIB) { + "128" + } else if key.ends_with(KEY_MAX_LIFETIME_MS) { + "30000" + } else { + "500" + } + .to_owned(), + )) + }); + + assert!(diagnostics.is_empty(), "a valid store should be quiet"); + assert_eq!( + resolve_mode(limits), + ServeMode::Reuse(SandboxLimits { + max_requests: 10, + max_memory_mib: 128, + max_lifetime: Duration::from_millis(30_000), + timeout: Duration::from_millis(500), + }), + "a fully configured store should enter the serving loop" + ); + } + + #[test] + fn deferred_diagnostics_are_held_until_flushed() { + let mut startup = StartupDiagnostics::default(); + startup.push("reuse declined"); + + assert_eq!( + startup.pending(), + 1, + "startup runs before the logger, so the message must be held" + ); + + startup.flush(); + + assert_eq!( + startup.pending(), + 0, + "flushing should emit and clear everything held" + ); + } + + #[test] + fn a_sandbox_starts_with_no_retained_application() { + let sandbox = Sandbox::default(); + + assert!( + sandbox.state().is_none(), + "construction must be lazy so the health probe never pays for it" + ); + assert_eq!( + sandbox.initialization_attempts(), + 0, + "a sandbox that has served nothing should report no build attempts" + ); + } + + #[test] + fn a_retained_application_is_reused_without_rebuilding() { + let mut sandbox = Sandbox::default(); + let builds = Cell::new(0_u32); + + sandbox + .initialize(|| { + builds.set(builds.get() + 1); + Ok::<_, &str>(test_retained()) + }) + .expect("the first build should succeed"); + assert_eq!(builds.get(), 1, "the first request should build"); + + for attempt in 2..=4_u32 { + sandbox + .initialize(|| -> Result { + builds.set(builds.get() + 1); + panic!("request {attempt} must not rebuild a retained application") + }) + .expect("a retained application should be reused"); + } + + assert_eq!( + builds.get(), + 1, + "four requests should cost exactly one build; that is the whole point" + ); + assert_eq!( + sandbox.initialization_attempts(), + 1, + "the framework counter should agree" + ); + } + + #[test] + fn a_failed_build_is_retried_and_a_later_success_is_retained() { + let mut sandbox = Sandbox::default(); + let builds = Cell::new(0_u32); + + // Request 1: the build fails. The error payload stands in for the + // error router `build_app_with_state` returns when state is `None`. + let failed = sandbox.initialize(|| { + builds.set(builds.get() + 1); + Err::("error router") + }); + assert_eq!( + failed.err(), + Some("error router"), + "a failed build must be handed back for this request only" + ); + assert!( + sandbox.state().is_none(), + "a failed build must not be retained" + ); + + // Request 2: the retry succeeds. Only reachable because the failure + // was not retained. + sandbox + .initialize(|| { + builds.set(builds.get() + 1); + Ok::<_, &str>(test_retained()) + }) + .expect("the retry should succeed"); + assert!( + sandbox.state().is_some(), + "the recovered application should be retained" + ); + + // Request 3: recovery is durable — no further build. + sandbox + .initialize(|| -> Result { + builds.set(builds.get() + 1); + panic!("a recovered application must not be rebuilt") + }) + .expect("request 3 should reuse the recovered app"); + + assert_eq!( + builds.get(), + 2, + "one failed build plus one successful build, then reuse" + ); + assert_eq!( + sandbox.initialization_attempts(), + 2, + "both attempts should be counted so a retry loop stays visible" + ); + } + + /// Drives the real reporting path: the same `observe` that `serve_loop` + /// wraps every callback in. A build performed by the FINAL callback must + /// appear in the retirement snapshot; reading the attempt count on entry + /// instead of on exit silently loses it. + #[test] + fn retirement_counters_include_a_build_from_the_final_callback() { + let mut sandbox = Sandbox::default(); + let mut counters = RetirementCounters::default(); + + // Callback 1 builds nothing, so nothing is attempted yet. + counters.observe(&mut sandbox, |_sandbox| {}); + assert_eq!( + counters.attempts(), + 0, + "a callback that never initializes should report no attempts" + ); + + // Callback 2 succeeds. The build happens DURING this callback, so it + // only shows up if the count is read on the way out. + counters.observe(&mut sandbox, |sandbox| { + sandbox + .initialize(|| Ok::<_, &str>(test_retained())) + .expect("the build should succeed"); + }); + assert_eq!( + counters.attempts(), + 1, + "a successful build on the final callback must be reported" + ); + } + + #[test] + fn retirement_counters_include_a_failed_build_from_the_final_callback() { + let mut sandbox = Sandbox::default(); + let mut counters = RetirementCounters::default(); + + // The application state is not retained, but the attempt itself stays + // in `initialization_attempts()` until the sandbox is dropped. What a + // stale snapshot loses is the report, not the count. + counters.observe(&mut sandbox, |sandbox| { + let failed = sandbox.initialize(|| Err::("error router")); + assert_eq!( + failed.err(), + Some("error router"), + "the failed build should be handed back" + ); + }); + + assert_eq!( + counters.attempts(), + 1, + "a failed build on the final callback must still be reported" + ); + assert!( + sandbox.state().is_none(), + "a failed build must not be retained" + ); + } + + #[test] + fn retirement_counters_report_the_ordinal_observed_on_entry() { + let mut sandbox = Sandbox::default(); + let mut counters = RetirementCounters::default(); + + // This covers snapshot copying only: that `observe` reports whatever + // the framework counter says rather than deriving a value of its own. + // That EdgeZero increments before invoking the callback is the + // framework's behaviour, covered by its tests and by runtime + // observation, not by this test — `requests()` advances only through + // `serve_custom` / `run_custom`, so a direct `observe` sees zero. + counters.observe(&mut sandbox, |sandbox| { + assert_eq!( + sandbox.requests(), + 0, + "no callback has been dispatched through the framework here" + ); + }); + + assert_eq!( + counters.requests(), + 0, + "the snapshot should mirror the framework counter, not a local one" + ); + } + + #[test] + fn counters_are_captured_only_when_metrics_are_enabled() { + let sandbox = Sandbox::default(); + let mut settings = Settings::default(); + + assert_eq!( + SandboxCounters::capture(&sandbox, 1, "req-1", None), + None, + "no settings should mean no counters" + ); + + settings.debug.sandbox_metrics_enabled = false; + assert_eq!( + SandboxCounters::capture(&sandbox, 1, "req-1", Some(&settings)), + None, + "counters must stay off in the default configuration" + ); + + settings.debug.sandbox_metrics_enabled = true; + assert_eq!( + SandboxCounters::capture(&sandbox, 7, "req-7", Some(&settings)), + Some(SandboxCounters { + ordinal: 7, + builds: 0, + request_id: "req-7".to_owned(), + }), + "enabling the flag should capture the current counters" + ); + } + + // Request ordinals and build attempts are counted by + // `edgezero_adapter_fastly::lifecycle::Sandbox` itself, incremented in its + // private per-callback hook. They are exercised through `serve_custom` / + // `run_custom` at runtime and covered by the framework's own tests, so + // there is nothing left here to unit test. +} diff --git a/crates/trusted-server-cli/Cargo.toml b/crates/trusted-server-cli/Cargo.toml index 19fb13366..eaa23f64a 100644 --- a/crates/trusted-server-cli/Cargo.toml +++ b/crates/trusted-server-cli/Cargo.toml @@ -57,6 +57,17 @@ trusted-server-core = { workspace = true } url = { workspace = true } which = { workspace = true } +# `ts dev sandbox-probe` needs a real HTTP/1.1 client with keep-alive. These +# are shared with the macOS-only proxy but scoped to every non-wasm host, so +# the probe is built and linted on Linux CI too. The exclusion still keeps +# `tokio`/`ring` off the repo-default `wasm32-wasip1` target. +[target.'cfg(not(target_family = "wasm"))'.dependencies] +bytes = { workspace = true } +http-body-util = { workspace = true } +hyper = { workspace = true, features = ["http1", "client"] } +hyper-util = { workspace = true, features = ["tokio"] } +tokio = { workspace = true, features = ["net", "rt-multi-thread"] } + # `ts dev proxy` is macOS-only — CA trust via the login keychain, Safari # automation via `networksetup`, and a native TLS / networking stack. Scoping # these dependencies to macOS keeps unsupported targets (notably the diff --git a/crates/trusted-server-cli/src/commands/dev/mod.rs b/crates/trusted-server-cli/src/commands/dev/mod.rs index 7a61d769b..702a5be32 100644 --- a/crates/trusted-server-cli/src/commands/dev/mod.rs +++ b/crates/trusted-server-cli/src/commands/dev/mod.rs @@ -3,6 +3,9 @@ // other host targets `ts dev` parses but exposes no subcommands. #[cfg(target_os = "macos")] pub mod proxy; +// Needs a real HTTP client; its dependencies are scoped to non-wasm hosts. +#[cfg(not(target_family = "wasm"))] +pub mod sandbox_probe; /// The `ts dev …` command group. #[derive(Debug, clap::Subcommand)] @@ -10,27 +13,20 @@ pub enum DevCommand { /// Run the local production-hostname dev proxy (macOS only). #[cfg(target_os = "macos")] Proxy(proxy::ProxyArgs), + /// Measure sandbox reuse from the counters on workload responses. + #[cfg(not(target_family = "wasm"))] + SandboxProbe(sandbox_probe::SandboxProbeArgs), } /// Dispatches a `dev` subcommand. /// /// # Errors -/// Returns the subcommand's failure rendered as a message. On non-macOS targets -/// `DevCommand` has no variants, so this never returns an error there. -// On non-macOS targets `DevCommand` is an empty enum: the by-value parameter is -// consumed by an empty `match`, which clippy reads as a needless by-value pass. -// Taking `&DevCommand` is not an option — a zero-arm `match` is not exhaustive -// over a reference type — so the owned parameter is required. -#[cfg_attr( - not(target_os = "macos"), - allow( - clippy::needless_pass_by_value, - reason = "empty enum requires owned value for exhaustive match" - ) -)] +/// Returns the subcommand's failure rendered as a message. pub fn run(command: DevCommand) -> Result<(), String> { match command { #[cfg(target_os = "macos")] DevCommand::Proxy(args) => proxy::run(&args).map_err(|report| format!("{report:?}")), + #[cfg(not(target_family = "wasm"))] + DevCommand::SandboxProbe(args) => sandbox_probe::run(&args), } } diff --git a/crates/trusted-server-cli/src/commands/dev/sandbox_probe.rs b/crates/trusted-server-cli/src/commands/dev/sandbox_probe.rs new file mode 100644 index 000000000..de5463838 --- /dev/null +++ b/crates/trusted-server-cli/src/commands/dev/sandbox_probe.rs @@ -0,0 +1,543 @@ +//! `ts dev sandbox-probe` — measures sandbox reuse from workload responses. +//! +//! Issues a keep-alive sequence against a locally served Trusted Server over a +//! single connection and reads the counters the Fastly adapter attaches to +//! each response. Reuse is reported only when one instance id serves several +//! *workload* requests with strictly increasing ordinals. +//! +//! The counters endpoint is deliberately not used to establish reuse: a +//! snapshot identifies the sandbox that served the probe, which need not be +//! the one that served the preceding request. + +use std::collections::BTreeMap; +use std::fmt::Write as _; +use std::time::{Duration, Instant}; + +use http_body_util::BodyExt as _; +use hyper_util::rt::TokioIo; + +/// Response headers the Fastly adapter attaches when sandbox metrics are on. +const HEADER_INSTANCE: &str = "x-ts-sandbox-instance"; +const HEADER_ORDINAL: &str = "x-ts-sandbox-ordinal"; +const HEADER_BUILDS: &str = "x-ts-sandbox-builds"; +const HEADER_REQUEST_ID: &str = "x-ts-sandbox-request-id"; + +/// Value the adapter reports when the runtime exposes no instance identity. +const INSTANCE_UNAVAILABLE: &str = "unavailable"; + +/// Paths that never carry counters because they short-circuit before dispatch. +const NON_WORKLOAD_PATHS: &[&str] = &["/health", "/_ts/debug/sandbox", "/_ts/debug/ja4"]; + +/// Arguments for `ts dev sandbox-probe`. +#[derive(Debug, clap::Args)] +pub struct SandboxProbeArgs { + /// Host and port of the locally served instance. + #[arg(long, default_value = "127.0.0.1:7676")] + pub authority: String, + + /// Workload path to exercise. + /// + /// Required, and it must reach the router. `/health` and the debug probes + /// short-circuit ahead of the point where counters are attached, so + /// probing one of those reports nothing. + #[arg(long)] + pub path: String, + + /// Number of requests to issue over one connection. + #[arg(long, default_value_t = 6)] + pub requests: u32, + + /// Per-request timeout in milliseconds. + #[arg(long, default_value_t = 10_000)] + pub timeout_ms: u64, +} + +/// One observation taken from a workload response. +#[derive(Debug, Clone)] +struct Observation { + instance: Option, + ordinal: Option, + builds: Option, + request_id: Option, + status: u16, + elapsed: Duration, +} + +/// Why the request sequence stopped, if it stopped early. +#[derive(Debug, Clone, PartialEq, Eq)] +enum Ending { + /// Every requested iteration completed. + Completed, + /// The connection closed. Expected when a bounded sandbox retires, but a + /// client cannot distinguish that from any other close. + ConnectionClosed(String), + /// A transport or body error, which invalidates the run. + TransportError(String), +} + +/// Runs the probe and prints a report. +/// +/// # Errors +/// Returns a message when the path is not a workload route, the async runtime +/// cannot start, or the connection cannot be established. +pub fn run(args: &SandboxProbeArgs) -> Result<(), String> { + if NON_WORKLOAD_PATHS.contains(&args.path.as_str()) { + return Err(format!( + "`{}` short-circuits before counters are attached; choose a workload route", + args.path + )); + } + + let runtime = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .map_err(|e| format!("failed to start async runtime: {e}"))?; + + let (observations, ending) = runtime.block_on(collect(args))?; + + // A transport error invalidates the run, so say so on stderr as well as + // in the report: a caller piping stdout to a file should still see it. + if let Ending::TransportError(reason) = &ending { + crate::output::warn(&format!( + "measurement run is unreliable, observations may be incomplete: {reason}" + )); + } + + crate::output::info(report(&observations, &ending).trim_end()); + Ok(()) +} + +/// Issues the keep-alive sequence over a single HTTP/1.1 connection. +async fn collect(args: &SandboxProbeArgs) -> Result<(Vec, Ending), String> { + let stream = tokio::net::TcpStream::connect(&args.authority) + .await + .map_err(|e| format!("failed to connect to {}: {e}", args.authority))?; + + let (mut sender, connection) = hyper::client::conn::http1::handshake(TokioIo::new(stream)) + .await + .map_err(|e| format!("HTTP/1.1 handshake with {} failed: {e}", args.authority))?; + + // The connection task drives the socket; it resolves when the peer closes. + let connection = tokio::spawn(connection); + + let timeout = Duration::from_millis(args.timeout_ms); + let mut observations = Vec::with_capacity(args.requests as usize); + let mut ending = Ending::Completed; + + for _ in 0..args.requests { + let request = hyper::Request::builder() + .method(hyper::Method::GET) + .uri(&args.path) + .header(hyper::header::HOST, &args.authority) + .body(String::new()) + .map_err(|e| format!("failed to build request: {e}"))?; + + let started = Instant::now(); + + let response = match tokio::time::timeout(timeout, sender.send_request(request)).await { + Ok(Ok(response)) => response, + Ok(Err(e)) if e.is_closed() || e.is_incomplete_message() => { + ending = Ending::ConnectionClosed(e.to_string()); + break; + } + Ok(Err(e)) => { + ending = Ending::TransportError(e.to_string()); + break; + } + Err(_) => { + ending = Ending::TransportError(format!("request timed out after {timeout:?}")); + break; + } + }; + + let (parts, body) = response.into_parts(); + + // Hyper handles chunk extensions, trailers, and close-delimited + // bodies. The body must be drained before the next request is sent. + match tokio::time::timeout(timeout, body.collect()).await { + Ok(Ok(collected)) => drop(collected.to_bytes()), + Ok(Err(e)) => { + ending = Ending::TransportError(format!("body read failed: {e}")); + break; + } + Err(_) => { + ending = Ending::TransportError(format!("body read timed out after {timeout:?}")); + break; + } + } + + observations.push(observation_from(&parts, started.elapsed())); + } + + drop(sender); + if let Ok(Err(e)) = connection.await + && ending == Ending::Completed + { + ending = Ending::ConnectionClosed(e.to_string()); + } + + Ok((observations, ending)) +} + +/// Extracts the counters from one response. +fn observation_from(parts: &hyper::http::response::Parts, elapsed: Duration) -> Observation { + let header = |name: &str| { + parts + .headers + .get(name) + .and_then(|value| value.to_str().ok()) + .map(str::to_owned) + }; + + Observation { + instance: header(HEADER_INSTANCE), + ordinal: header(HEADER_ORDINAL).and_then(|v| v.parse().ok()), + builds: header(HEADER_BUILDS).and_then(|v| v.parse().ok()), + request_id: header(HEADER_REQUEST_ID), + status: parts.status.as_u16(), + elapsed, + } +} + +/// Renders the report, stating explicitly what was and was not established. +fn report(observations: &[Observation], ending: &Ending) -> String { + let mut out = String::new(); + let _ = writeln!(out, "responses: {}", observations.len()); + + match ending { + Ending::Completed => {} + Ending::ConnectionClosed(reason) => { + let _ = writeln!( + out, + "ended: connection closed ({reason}); a client cannot tell sandbox \ + retirement from any other close" + ); + } + Ending::TransportError(reason) => { + let _ = writeln!(out, "ended: transport error ({reason})"); + } + } + + if observations.is_empty() { + let _ = writeln!(out, "reuse: unverified (no responses)"); + return out; + } + + for (index, observation) in observations.iter().enumerate() { + let _ = writeln!( + out, + " {:>3}. status={} instance={} ordinal={} builds={} request={} elapsed={:?}", + index + 1, + observation.status, + observation.instance.as_deref().unwrap_or("-"), + observation + .ordinal + .map_or_else(|| "-".to_owned(), |v| v.to_string()), + observation + .builds + .map_or_else(|| "-".to_owned(), |v| v.to_string()), + observation.request_id.as_deref().unwrap_or("-"), + observation.elapsed, + ); + } + + out.push_str(&verdict(observations, ending)); + out +} + +/// Decides whether the observations establish reuse. +/// +/// Deliberately conservative: only strictly increasing ordinals on one +/// instance count. A repeated or decreasing ordinal under one instance id +/// means the attribution is unreliable — the id is not identifying what it +/// claims to — so no conclusion is drawn either way. Counting bare repeats +/// would instead report those runs as reuse. +fn verdict(observations: &[Observation], ending: &Ending) -> String { + if let Ending::TransportError(reason) = ending { + return format!("reuse: unverified (transport error: {reason})\n"); + } + + if observations.iter().all(|o| o.instance.is_none()) { + return "reuse: unverified (no counters; enable debug.sandbox_metrics_enabled, \ + and probe a route whose response is private and no-store)\n" + .to_owned(); + } + + if observations + .iter() + .any(|o| o.instance.as_deref() == Some(INSTANCE_UNAVAILABLE)) + { + return "reuse: unverified (runtime reported no instance identity)\n".to_owned(); + } + + // Partial attribution means some responses cannot be placed in a sandbox, + // so neither reuse nor its absence can be concluded. + if observations + .iter() + .any(|o| o.instance.is_none() || o.ordinal.is_none()) + { + return "reuse: unverified (incomplete attribution: some responses carried no \ + instance id or ordinal)\n" + .to_owned(); + } + + let mut by_instance: BTreeMap<&str, Vec> = BTreeMap::new(); + for observation in observations { + if let (Some(instance), Some(ordinal)) = + (observation.instance.as_deref(), observation.ordinal) + { + by_instance.entry(instance).or_default().push(ordinal); + } + } + + let mut out = String::new(); + let _ = writeln!(out, "distinct instances: {}", by_instance.len()); + + // Ordinals from one sandbox must arrive strictly increasing. Anything else + // means the ids are not the identities they claim to be. + let anomalous: Vec<&str> = by_instance + .iter() + .filter(|(_, ordinals)| !ordinals.windows(2).all(|pair| pair[0] < pair[1])) + .map(|(instance, _)| *instance) + .collect(); + if !anomalous.is_empty() { + let _ = writeln!( + out, + "reuse: unverified (instance(s) {} reported repeated or decreasing ordinals)", + anomalous.join(", ") + ); + return out; + } + + let reused = by_instance + .values() + .filter(|ordinals| ordinals.len() > 1) + .count(); + let max_per_instance = by_instance.values().map(Vec::len).max().unwrap_or_default(); + let max_builds = observations + .iter() + .filter_map(|o| o.builds) + .max() + .unwrap_or_default(); + + let _ = writeln!(out, "max requests per instance: {max_per_instance}"); + let _ = writeln!(out, "max builds per instance: {max_builds}"); + + if reused == 0 { + let _ = writeln!( + out, + "reuse: not observed (every response came from a distinct instance)" + ); + } else { + let _ = writeln!( + out, + "reuse: observed ({reused} instance(s) served more than one workload request)" + ); + } + out +} + +#[cfg(test)] +mod tests { + use super::*; + + fn observation( + instance: Option<&str>, + ordinal: Option, + builds: Option, + ) -> Observation { + Observation { + instance: instance.map(str::to_owned), + ordinal, + builds, + request_id: Some("req".to_owned()), + status: 200, + elapsed: Duration::from_millis(1), + } + } + + #[test] + fn absent_counters_report_unverified_not_negative() { + let observations = vec![observation(None, None, None), observation(None, None, None)]; + + let verdict = verdict(&observations, &Ending::Completed); + + assert!( + verdict.contains("unverified"), + "missing counters are not evidence that reuse failed" + ); + assert!( + !verdict.contains("reusable-sandbox"), + "counters are feature-independent, so the hint must not name the feature: {verdict}" + ); + } + + #[test] + fn an_unavailable_instance_id_reports_unverified() { + let observations = vec![ + observation(Some(INSTANCE_UNAVAILABLE), Some(1), Some(1)), + observation(Some(INSTANCE_UNAVAILABLE), Some(2), Some(1)), + ]; + + assert!( + verdict(&observations, &Ending::Completed).contains("unverified"), + "a runtime with no identity cannot establish reuse" + ); + } + + #[test] + fn repeated_ordinals_on_one_instance_are_not_reuse() { + // Unreliable attribution: one id reporting ordinal 1 twice. A bare + // repeat count would call this reuse. + let observations = vec![ + observation(Some("a"), Some(1), Some(1)), + observation(Some("a"), Some(1), Some(1)), + ]; + + let verdict = verdict(&observations, &Ending::Completed); + + assert!( + verdict.contains("unverified"), + "duplicate ordinals must not read as reuse, got: {verdict}" + ); + assert!( + verdict.contains("repeated or decreasing"), + "should name the anomaly, got: {verdict}" + ); + } + + #[test] + fn decreasing_ordinals_are_not_reuse() { + let observations = vec![ + observation(Some("a"), Some(3), Some(1)), + observation(Some("a"), Some(2), Some(1)), + ]; + + assert!( + verdict(&observations, &Ending::Completed).contains("unverified"), + "ordinals going backwards mean the identity is untrustworthy" + ); + } + + #[test] + fn incomplete_attribution_reports_unverified() { + let observations = vec![ + observation(Some("a"), Some(1), Some(1)), + observation(None, None, None), + ]; + + let verdict = verdict(&observations, &Ending::Completed); + + assert!( + verdict.contains("incomplete attribution"), + "a response that cannot be placed in a sandbox blocks a verdict, got: {verdict}" + ); + } + + #[test] + fn a_transport_error_invalidates_the_run() { + let observations = vec![ + observation(Some("a"), Some(1), Some(1)), + observation(Some("a"), Some(2), Some(1)), + ]; + + let verdict = verdict(&observations, &Ending::TransportError("reset".to_owned())); + + assert!( + verdict.contains("unverified"), + "a transport error makes the run unreliable, got: {verdict}" + ); + } + + #[test] + fn distinct_instances_report_reuse_not_observed() { + let observations = vec![ + observation(Some("a"), Some(1), Some(1)), + observation(Some("b"), Some(1), Some(1)), + ]; + + let verdict = verdict(&observations, &Ending::Completed); + + assert!( + verdict.contains("not observed"), + "one request per instance is arm A, got: {verdict}" + ); + assert!( + verdict.contains("distinct instances: 2"), + "should report the instance count, got: {verdict}" + ); + } + + #[test] + fn increasing_ordinals_on_one_instance_report_reuse() { + let observations = vec![ + observation(Some("a"), Some(1), Some(1)), + observation(Some("a"), Some(2), Some(1)), + observation(Some("a"), Some(3), Some(1)), + ]; + + let verdict = verdict(&observations, &Ending::Completed); + + assert!( + verdict.contains("reuse: observed"), + "three increasing ordinals on one instance is reuse, got: {verdict}" + ); + assert!( + verdict.contains("max requests per instance: 3"), + "should report the depth reached, got: {verdict}" + ); + assert!( + verdict.contains("max builds per instance: 1"), + "one build across three requests is the arm-C signal, got: {verdict}" + ); + } + + #[test] + fn a_closed_connection_is_reported_without_being_read_as_retirement() { + let observations = vec![ + observation(Some("a"), Some(1), Some(1)), + observation(Some("a"), Some(2), Some(1)), + ]; + + let report = report( + &observations, + &Ending::ConnectionClosed("closed".to_owned()), + ); + + assert!( + report.contains("cannot tell sandbox"), + "a close is ambiguous and must be reported as such, got: {report}" + ); + assert!( + report.contains("reuse: observed"), + "a close does not invalidate observations already collected, got: {report}" + ); + } + + #[test] + fn non_workload_paths_are_rejected_before_connecting() { + for path in NON_WORKLOAD_PATHS { + let args = SandboxProbeArgs { + authority: "127.0.0.1:7676".to_owned(), + path: (*path).to_owned(), + requests: 2, + timeout_ms: 100, + }; + + let error = run(&args).expect_err("should refuse a short-circuiting path"); + + assert!( + error.contains("short-circuits"), + "should explain why `{path}` cannot measure reuse, got: {error}" + ); + } + } + + #[test] + fn an_empty_run_is_unverified() { + assert!( + report(&[], &Ending::Completed).contains("unverified"), + "no responses cannot establish anything" + ); + } +} diff --git a/crates/trusted-server-cli/src/lib.rs b/crates/trusted-server-cli/src/lib.rs index 74045ed7a..208d39776 100644 --- a/crates/trusted-server-cli/src/lib.rs +++ b/crates/trusted-server-cli/src/lib.rs @@ -25,5 +25,8 @@ pub use run::{RunOutcome, run_from_env}; // internals. #[cfg(not(target_arch = "wasm32"))] pub mod commands; -#[cfg(target_os = "macos")] +// Console output wrappers. Gated to non-wasm hosts rather than macOS: the +// macOS-only proxy was its first consumer, but `ts dev sandbox-probe` builds +// on every host target and needs it too. +#[cfg(not(target_arch = "wasm32"))] mod output; diff --git a/crates/trusted-server-core/src/integrations/datadome/protection_scope.rs b/crates/trusted-server-core/src/integrations/datadome/protection_scope.rs index a83f73217..73bd77070 100644 --- a/crates/trusted-server-core/src/integrations/datadome/protection_scope.rs +++ b/crates/trusted-server-core/src/integrations/datadome/protection_scope.rs @@ -864,12 +864,23 @@ mod tests { clear_ip_cidr_source_cache_for_tests(); let mut config = config_with_protection(); config.protection_excluded_ip_cidr_sources = vec![ProtectionIpCidrSourceConfig { - config_store: "datadome-ip-bypass".to_string(), - key: "googlebot_ips".to_string(), + config_store: "example-ip-bypass".to_string(), + key: "example_cidrs".to_string(), }]; + config.protection_exclusion_rules = vec![ProtectionExclusionRuleConfig { + id: "example-source-rule".to_owned(), + enabled: true, + methods: Vec::new(), + matcher: ProtectionMatcherConfig::IpCidrSource { + config_store: "example-rule-store".to_owned(), + key: "example_cidrs".to_owned(), + }, + }]; + // Force refreshes as well as lookups to exercise replacement of entries. + config.protection_ip_list_cache_ttl_seconds = 0; let scope = ProtectionScope::compile(&config).expect("should compile scope"); let mut data = HashMap::new(); - data.insert("googlebot_ips".to_string(), "203.0.113.0/24".to_string()); + data.insert("example_cidrs".to_string(), "203.0.113.0/24".to_string()); let services = build_services_with_config_and_secret(HashMapConfigStore::new(data), NoopSecretStore); @@ -892,6 +903,50 @@ mod tests { .. } )); + + let expected_keys = HashSet::from([ + ProtectionIpCidrSourceCacheKey { + config_store: "example-ip-bypass".to_owned(), + key: "example_cidrs".to_owned(), + }, + ProtectionIpCidrSourceCacheKey { + config_store: "example-rule-store".to_owned(), + key: "example_cidrs".to_owned(), + }, + ]); + for index in 1..=32 { + let path = format!("/example/{index}"); + let query = format!("example={index}"); + let decision = scope.evaluate( + &facts( + if index % 2 == 0 { "GET" } else { "POST" }, + &path, + Some(&query), + Some(IpAddr::V4(Ipv4Addr::new(192, 0, 2, index))), + Some(64512 + u32::from(index)), + ), + &services, + ); + assert!( + matches!(decision, ProtectionScopeDecision::Protect), + "should evaluate both configured sources for unmatched traffic" + ); + let keys: HashSet<_> = IP_CIDR_SOURCE_CACHE + .lock() + .expect("should lock CIDR cache") + .keys() + // Other native tests can populate unrelated sources concurrently. + .filter(|key| { + key.config_store == "example-ip-bypass" + || key.config_store == "example-rule-store" + }) + .cloned() + .collect(); + assert_eq!( + keys, expected_keys, + "should retain only config-derived keys across traffic and refreshes" + ); + } } #[test] diff --git a/crates/trusted-server-core/src/integrations/google_tag_manager.rs b/crates/trusted-server-core/src/integrations/google_tag_manager.rs index ae9ceb7cb..d292f5b38 100644 --- a/crates/trusted-server-core/src/integrations/google_tag_manager.rs +++ b/crates/trusted-server-core/src/integrations/google_tag_manager.rs @@ -12,7 +12,7 @@ //! | `GET/POST` | `.../collect` | Proxies GA analytics beacons | //! | `GET/POST` | `.../g/collect` | Proxies GA4 analytics beacons | -use std::sync::{Arc, LazyLock, Mutex}; +use std::sync::{Arc, LazyLock}; use async_trait::async_trait; use edgezero_core::body::Body as EdgeBody; @@ -24,6 +24,7 @@ use serde::{Deserialize, Serialize}; use validator::{Validate, ValidationError}; use crate::error::TrustedServerError; +use crate::integrations::ScriptTextAccumulator; use crate::integrations::{ AttributeRewriteAction, IntegrationAttributeContext, IntegrationAttributeRewriter, IntegrationEndpoint, IntegrationProxy, IntegrationRegistration, IntegrationScriptContext, @@ -356,14 +357,6 @@ pub struct GoogleTagManagerIntegration { /// allowlist as "any host", so without this a 3xx from the upstream would /// let an arbitrary origin's body be re-served as first-party JavaScript. proxy_allowed_domains: Vec, - /// Accumulates text fragments when `lol_html` splits a text node across - /// chunk boundaries. Drained on `is_last_in_text_node`. - /// - /// Uses `Mutex` to satisfy the `Sync` bound on `IntegrationScriptRewriter`. - /// The pipeline is single-threaded (`lol_html::HtmlRewriter` is `!Send`), - /// so the lock is uncontended. `lol_html` delivers text chunks sequentially - /// per element — the buffer is always empty when a new element's text begins. - accumulated_text: Mutex, } impl GoogleTagManagerIntegration { @@ -372,7 +365,6 @@ impl GoogleTagManagerIntegration { Arc::new(Self { config, proxy_allowed_domains, - accumulated_text: Mutex::new(String::new()), }) } @@ -1029,10 +1021,12 @@ impl IntegrationScriptRewriter for GoogleTagManagerIntegration { } fn rewrite(&self, content: &str, ctx: &IntegrationScriptContext<'_>) -> ScriptRewriteAction { - let mut buf = self - .accumulated_text - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner); + // Per document, never per registry: a registry-lifetime buffer would + // carry one document's partial script into the next. + let accumulator = ctx + .document_state + .get_or_insert_with(GTM_INTEGRATION_ID, ScriptTextAccumulator::default); + let mut buf = accumulator.buffer(); // Cheap gate: only engage the accumulation path for scripts whose // running text could plausibly contain a GTM/GA domain. Unrelated @@ -3465,6 +3459,76 @@ container_id = "GTM-DEFAULT" } } + #[test] + fn an_interrupted_document_leaves_no_residue_for_the_next_document() { + // The registry holds one `Arc` for the + // lifetime of the application, so both documents below go through the + // SAME rewriter instance. Only the document state differs. With the + // buffer owned by the rewriter this test fails: document two emits + // document one's secret. + let integration = GoogleTagManagerIntegration::new(tag_config("GTM-LEAK01", &[])); + + // Document one: a GTM snippet that is cut off before its final + // fragment ever arrives, as a client disconnect or truncated origin + // body would do. + let first_document = IntegrationDocumentState::default(); + let interrupted = IntegrationScriptContext { + selector: "script", + request_host: "first.example.com", + request_scheme: "https", + origin_host: "origin.example.com", + is_last_in_text_node: false, + max_buffered_script_bytes: 16 * 1024 * 1024, + document_state: &first_document, + }; + // Must end mid-domain: that is what makes the cheap prefix gate + // accumulate rather than pass the fragment through untouched. + let secret = + r#"(function(w,d,s,l,i){var token='SESSION-ONE-SECRET'; j.src='https://www.google"#; + + let action = IntegrationScriptRewriter::rewrite(&*integration, secret, &interrupted); + assert_eq!( + action, + ScriptRewriteAction::RemoveNode, + "the partial fragment should be withheld, which is what strands it" + ); + + // Document two: a different request, a fresh document state, the same + // registry and the same rewriter. + let second_document = IntegrationDocumentState::default(); + let fresh = IntegrationScriptContext { + selector: "script", + request_host: "second.example.com", + request_scheme: "https", + origin_host: "origin.example.com", + is_last_in_text_node: true, + max_buffered_script_bytes: 16 * 1024 * 1024, + document_state: &second_document, + }; + let benign = r#"(function(w,d,s,l,i){j.src='https://www.googletagmanager.com/gtm.js?id='+i;})(window,document,'script','dataLayer','GTM-LEAK01');"#; + + let action = IntegrationScriptRewriter::rewrite(&*integration, benign, &fresh); + + let emitted = match action { + ScriptRewriteAction::Replace(rewritten) => rewritten, + ScriptRewriteAction::Keep => benign.to_owned(), + other => panic!("expected the second document to be emitted, got {other:?}"), + }; + + assert!( + !emitted.contains("SESSION-ONE-SECRET"), + "the interrupted document's content must not reach the next document, got: {emitted}" + ); + assert!( + !emitted.contains("first.example.com"), + "no trace of the previous document should survive, got: {emitted}" + ); + assert!( + emitted.contains("/integrations/google_tag_manager/gtm.js"), + "the second document should still be rewritten correctly, got: {emitted}" + ); + } + #[test] fn non_gtm_fragmented_script_returns_keep_on_every_fragment() { // A script with no GTM marker in flight (and no plausible prefix) must diff --git a/crates/trusted-server-core/src/integrations/mod.rs b/crates/trusted-server-core/src/integrations/mod.rs index ee883b0d5..820e2c24b 100644 --- a/crates/trusted-server-core/src/integrations/mod.rs +++ b/crates/trusted-server-core/src/integrations/mod.rs @@ -37,6 +37,7 @@ pub use registry::{ IntegrationRequestFilter, IntegrationScriptContext, IntegrationScriptRewriter, ProxyDispatchInput, RequestFilterDecision, RequestFilterEffects, RequestFilterInput, RequestFilterRegistryInput, RequestFilterRegistryOutcome, ScriptRewriteAction, + ScriptTextAccumulator, }; /// Registers or retrieves a platform backend for the given URL. diff --git a/crates/trusted-server-core/src/integrations/nextjs/script_rewriter.rs b/crates/trusted-server-core/src/integrations/nextjs/script_rewriter.rs index b6de84449..4e76ff7b1 100644 --- a/crates/trusted-server-core/src/integrations/nextjs/script_rewriter.rs +++ b/crates/trusted-server-core/src/integrations/nextjs/script_rewriter.rs @@ -234,6 +234,53 @@ mod tests { } } + #[test] + fn an_interrupted_document_leaves_no_residue_for_the_next_document() { + // One rewriter serves every document the registry serves, so both + // documents below share it and only the document state differs. With + // the buffer owned by the rewriter this test fails: document two + // emits document one's payload. + let rewriter = NextJsNextDataRewriter::new(test_config()) + .expect("should build Next.js structured rewriter"); + + // Document one: cut off before its final fragment arrives. + let first_document = IntegrationDocumentState::default(); + let interrupted = IntegrationScriptContext { + is_last_in_text_node: false, + ..ctx("script#__NEXT_DATA__", &first_document) + }; + let secret = r#"{"props":{"pageProps":{"sessionToken":"SESSION-ONE-SECRET","href":"#; + + let action = rewriter.rewrite(secret, &interrupted); + assert_eq!( + action, + ScriptRewriteAction::RemoveNode, + "the partial fragment should be withheld, which is what strands it" + ); + + // Document two: fresh document state, same rewriter. + let second_document = IntegrationDocumentState::default(); + let fresh = ctx("script#__NEXT_DATA__", &second_document); + let benign = r#"{"props":{"pageProps":{"href":"https://origin.example.com/reviews"}}}"#; + + let action = rewriter.rewrite(benign, &fresh); + + let emitted = match action { + ScriptRewriteAction::Replace(value) => value, + ScriptRewriteAction::Keep => benign.to_owned(), + other => panic!("expected the second document to be emitted, got {other:?}"), + }; + + assert!( + !emitted.contains("SESSION-ONE-SECRET"), + "the interrupted document's content must not reach the next document, got: {emitted}" + ); + assert!( + emitted.contains("ts.example.com"), + "the second document should still be rewritten correctly, got: {emitted}" + ); + } + #[test] fn structured_rewriter_updates_next_data_payload() { let payload = r#"{"props":{"pageProps":{"primary":{"href":"https://origin.example.com/reviews"},"secondary":{"href":"http://origin.example.com/sign-in"},"fallbackHref":"http://origin.example.com/legacy","protoRelative":"//origin.example.com/assets/logo.png"}}}"#; diff --git a/crates/trusted-server-core/src/integrations/registry.rs b/crates/trusted-server-core/src/integrations/registry.rs index 47c521e61..d3c09cd7d 100644 --- a/crates/trusted-server-core/src/integrations/registry.rs +++ b/crates/trusted-server-core/src/integrations/registry.rs @@ -1,6 +1,6 @@ use std::any::{Any, TypeId}; use std::collections::BTreeMap; -use std::sync::{Arc, Mutex}; +use std::sync::{Arc, Mutex, MutexGuard, PoisonError}; use async_trait::async_trait; use edgezero_core::body::Body as EdgeBody; @@ -191,6 +191,38 @@ impl IntegrationDocumentState { } } +/// Per-document buffer for script text fragments split across chunks. +/// +/// `lol_html` can deliver one text node as several chunks, so a rewriter that +/// needs the whole script must accumulate until `is_last_in_text_node`. +/// +/// This lives in [`IntegrationDocumentState`] rather than on the rewriter. +/// Rewriters are registered once as `Arc` and +/// live as long as the [`IntegrationRegistry`], so a buffer owned by a +/// rewriter is shared by every document that registry serves. A document whose +/// stream ends before the final fragment — client disconnect, origin error, +/// truncated body — leaves its partial script in that buffer, and the next +/// document prepends the residue to its own accumulation. That corrupts the +/// response and can disclose the previous document's content. +/// +/// Keyed per integration id, so each integration gets its own buffer, and +/// dropped with the document state at end of document. +#[derive(Debug, Default)] +pub struct ScriptTextAccumulator { + buffer: Mutex, +} + +impl ScriptTextAccumulator { + /// Locks the buffer. + /// + /// Recovers from poisoning rather than panicking: a poisoned buffer holds + /// at worst a partial script, and the caller's `is_last_in_text_node` + /// handling already tolerates unexpected contents. + pub fn buffer(&self) -> MutexGuard<'_, String> { + self.buffer.lock().unwrap_or_else(PoisonError::into_inner) + } +} + /// Describes an HTTP endpoint exposed by an integration. #[derive(Clone, Debug)] pub struct IntegrationEndpoint { diff --git a/crates/trusted-server-core/src/settings.rs b/crates/trusted-server-core/src/settings.rs index 97a3500b6..bf4ff7999 100644 --- a/crates/trusted-server-core/src/settings.rs +++ b/crates/trusted-server-core/src/settings.rs @@ -2549,6 +2549,33 @@ pub struct DebugConfig { /// un-sanitized creative for diagnostics, so never enable in production. #[serde(default)] pub inject_adm_for_testing: bool, + + /// Expose the reusable-sandbox counters endpoint at `GET /_ts/debug/sandbox` + /// and attach the same counters to private, no-store workload responses. + /// + /// The counters are the guest-instance identifier, the request ordinal + /// within that instance, the application build count, and the request + /// correlation id. They carry no settings, secrets, or request content. + /// Cacheable responses omit counters without changing their cache policy; + /// probes of those routes cannot establish sandbox reuse. + /// + /// Independent of the adapter's `reusable-sandbox` Cargo feature by design: + /// the feature decides whether a `Serve` loop exists, this flag decides + /// whether counters are emitted. Keeping them separate is what lets the + /// feature-off baseline be measured on the same channel as the reuse arms. + /// + /// Skipped from serialization while false: [`DebugConfig`] denies unknown + /// fields, so a default blob must stay readable by a binary built before + /// this field existed. A blob with it enabled requires restoring a + /// compatible blob before rolling back, the same trade + /// [`DebugConfig::auction_html_comment_options`] makes. + #[serde(default, skip_serializing_if = "is_false")] + pub sandbox_metrics_enabled: bool, +} + +/// Serde predicate for omitting `false` flags from serialized config blobs. +fn is_false(value: &bool) -> bool { + !*value } /// Metadata keys safe to surface in the `ts-debug` auction comment. @@ -3697,6 +3724,57 @@ mod tests { use std::collections::BTreeSet; use std::sync::Arc; + /// `DebugConfig` denies unknown fields, so a binary built before + /// `sandbox_metrics_enabled` existed must still accept a default blob. + /// That only holds while the flag is skipped during serialization. + #[test] + fn default_debug_config_omits_sandbox_metrics_for_rollback() { + let serialized = + serde_json::to_value(DebugConfig::default()).expect("should serialize debug config"); + + assert!( + serialized.get("sandbox_metrics_enabled").is_none(), + "a default blob must not carry the field, or an older binary rejects it: {serialized}" + ); + } + + #[test] + fn enabled_sandbox_metrics_serializes_and_round_trips() { + let config = DebugConfig { + sandbox_metrics_enabled: true, + ..DebugConfig::default() + }; + + let serialized = serde_json::to_value(&config).expect("should serialize debug config"); + assert_eq!( + serialized.get("sandbox_metrics_enabled"), + Some(&json!(true)), + "an enabled flag must be written so the setting survives a round trip" + ); + + let restored: DebugConfig = + serde_json::from_value(serialized).expect("should deserialize debug config"); + assert!( + restored.sandbox_metrics_enabled, + "the flag should survive a round trip" + ); + } + + #[test] + fn debug_config_accepts_a_blob_without_the_sandbox_field() { + let restored: DebugConfig = serde_json::from_value(json!({"ja4_endpoint_enabled": true})) + .expect("should deserialize a blob written before the field existed"); + + assert!( + restored.ja4_endpoint_enabled, + "existing fields should still load" + ); + assert!( + !restored.sandbox_metrics_enabled, + "an absent flag should default to off" + ); + } + use crate::auction::build_orchestrator; use crate::integrations::{ IntegrationRegistry, gpt::GptConfig, nextjs::NextJsIntegrationConfig, diff --git a/docs/superpowers/plans/2026-09-21-fastly-review-resolution.md b/docs/superpowers/plans/2026-09-21-fastly-review-resolution.md new file mode 100644 index 000000000..cfe57b846 --- /dev/null +++ b/docs/superpowers/plans/2026-09-21-fastly-review-resolution.md @@ -0,0 +1,27 @@ +# Fastly review resolution implementation plan + +**Goal:** Resolve the reusable-sandbox review findings without changing public response caching. + +**Approved scope:** Gate workload counters on existing private/no-store policy; require a positive memory limit; report logger installation failures; cover config-bounded CIDR caching; correct retention and measurement documentation; prepare a current PR description. + +**Architecture:** Keep policy and lifecycle changes in the Fastly adapter. Exercise the existing DataDome cache without changing its implementation. Memory limits are checked between requests by the SDK, not enforced as an allocation ceiling during a request. Missing heap support retires the sandbox conservatively. + +**Tech stack:** Rust, Fastly Compute, EdgeZero, Viceroy. + +- [x] Add counter privacy regressions in `crates/trusted-server-adapter-fastly/src/main.rs`; run `cargo test-fastly-reuse --locked sandbox_counters` before and after the fix. Cover public/absent/quoted policies and terminal-private responses. +- [x] Add memory-limit validation regressions in `crates/trusted-server-adapter-fastly/src/sandbox.rs`; require `TS__SANDBOX__MAX_MEMORY_MIB` to fit a positive `u32`, and pass it to `Serve::with_max_memory`. Test missing, zero, malformed, overflowing and valid values. +- [x] Report failed logger installation through the approved, narrowly scoped stderr fallback; preserve successful-only setup and diagnostic flushing. Run the existing setup retry regression and inspect the direct stderr branch; avoid adding test-only indirection around a single output statement. +- [x] Extend the existing CIDR source test in `crates/trusted-server-core/src/integrations/datadome/protection_scope.rs` to assert configured cache keys remain unchanged across varied traffic and refreshes. +- [x] Correct `AppState` lifetime docs, configuration guidance and measurement limitations. Preserve provenance of historical measurements. +- [x] Prepare an updated PR body with the actual EdgeZero pin and missing commits; distinguish previous validation from checks run for these fixes. +- [x] Run target-matched regression tests, then the full CI gate list in `AGENTS.md`; report any environment blockers accurately. Review the final diff before handoff. + +## Verification outcome + +All `AGENTS.md` CI gates passed. Additional host CLI tests passed (532 passed, +18 ignored). After the final CIDR assertion adjustment, the targeted WASM test, +59 host DataDome tests, Fastly clippy and Rust formatting passed again. Independent +code review found no remaining issues. + +The counter-privacy and missing-memory regressions failed before their fixes. +Deployed memory behavior and historical performance measurements were not rerun. diff --git a/docs/superpowers/specs/2026-09-17-fastly-reusable-sandbox-design.md b/docs/superpowers/specs/2026-09-17-fastly-reusable-sandbox-design.md new file mode 100644 index 000000000..1dc0fe9c0 --- /dev/null +++ b/docs/superpowers/specs/2026-09-17-fastly-reusable-sandbox-design.md @@ -0,0 +1,599 @@ +# Fastly reusable sandbox adoption + +Issue: [#856](https://github.com/IABTechLab/trusted-server/issues/856) + +## Problem + +Every Fastly Compute request currently starts a fresh Wasm sandbox, so the +entry point repeats the whole initialization sequence: read the runtime env +config store, open the application config store, load and parse settings, +resolve secrets, compile the auction plan, build the orchestrator, build the +integration registry, and construct the telemetry sink. The `OnceLock` regex +caches in `settings.rs` are initialized and discarded within a single request. + +Fastly SDK 0.12.1 exposes `fastly::http::serve::Serve`, which lets one sandbox +handle several requests. This design adopts it so that initialization can be +amortized, while keeping reuse opt-in and leaving the default deployment +behaviour unchanged. + +## Scope + +In scope: the Fastly adapter entry point and the state it owns, plus one +`trusted-server-core` change — moving GTM's script accumulation buffer out of +its registry-lifetime object, without which retention is unsafe. Next.js now +uses the per-document state introduced by #1135. +See [Retained rewrite buffers](#retained-rewrite-buffers-blocking). + +Out of scope: Spin, Cloudflare, and Axum adapters; any other change to routing, +auction, EC, or integration behaviour. Memory-bounded retirement is included +following review; it is a between-request SDK check, not an allocation ceiling. + +## Current entry point + +The adapter does not use `#[fastly::main]`. It owns its own `main`, converts +the raw request itself, dispatches straight into the router, and sends the +response explicitly. That shape is load-bearing and must survive the change. + +| Concern | Location | +| ----------------------------------------------------- | ----------------------------------------------------------------- | +| Raw `FastlyRequest::from_client()` | `main.rs:78-89` | +| Health probe ahead of all initialization | `main.rs:65-72`, called at `main.rs:81` | +| Global logger install | `main.rs:87` → `logging.rs:82-106` | +| Runtime env read | `main.rs:93-94` | +| Native config-store handle → request extension | opened `main.rs:58-63`, `main.rs:117-127`; inserted `main.rs:178` | +| Application build | `main.rs:129` → `app.rs:1278-1297` | +| Direct router dispatch (keeps duplicate `Set-Cookie`) | `main.rs:170-180` | +| Progressive streaming via `stream_to_client` | `main.rs:338-368` | +| Synchronous post-send pull sync | `main.rs:450-466` | + +Two properties of this path block naive reuse. + +1. `logging.rs:105` ends in `.apply().expect("should initialize logger")`. + A second call panics, because a global logger is already installed. +2. `app.rs:1288-1296` returns `startup_error_router` with `state: None` when + settings fail to load. Retaining that result would pin a sandbox into + permanent error mode for every subsequent request it serves. + +## Dependencies + +The workspace pins EdgeZero at `35a72835322fe0127beb2ce998e4988e95974373` on +`feat/reusable-app-lifecycle`. (It briefly sat at `277544c4` on the same +branch; that revision carried the CLI and Cloudflare fixes but no lifecycle +module.) + +The current revision also exposes registry-aware Fastly request conversion +and uses typed `dispatch_app::` helpers inside the Cloudflare and Spin +`run_app::` entry points. Those entry points remain compatible with +Trusted Server. Its custom Fastly path continues to use raw request +conversion and application-managed stores; it does not opt into the +framework's additional required KV bindings. + +This includes `edgezero_adapter_fastly::lifecycle`, introduced at `76c59b44`, which this adapter +now uses instead of its own equivalents. The framework owns lazy +successful-only retention, the callback count, the initialization-attempt +count, the one-time setup guard, and the serving wrappers. The sections below +describe the design as originally implemented; where they name a local +mechanism such as `logger_installed`, `resolve_app` or a hand-rolled +`Serve::run_with_context`, that mechanism has since been replaced by +`Sandbox::setup_once`, `Sandbox::initialize`, and +`lifecycle::serve_custom` / `run_custom`. What stays application-owned is +unchanged. + +The pin is **not** for the `Serve` re-export: that type comes from the +already-pinned `fastly 0.12.1` SDK, and EdgeZero's own contract says not to +repin merely to swap the import. It is for the CLI fixes at that revision +(push/diff validation scoped to the selected adapter, and secret-reference +redaction in Spin diagnostics) and the Cloudflare duplicate-header fix. Every +`edgezero-core` and `edgezero-core::router` change across the range is test +only, so no runtime behaviour we depend on moved. + +`Serve`, `ServeSummary`, and `HandlerResult` all come from `fastly 0.12.1`, +which is already resolved in `Cargo.lock`. EdgeZero PR #379 re-exports only +`Serve` and `ServeSummary` (not `HandlerResult`) and adds `serve_app` / +`serve_app_with_request_extensions`. + +Those helpers are unusable here, but not because of header handling: EdgeZero's +`to_fastly_response` uses `append_header` +(`crates/edgezero-adapter-fastly/src/response.rs`), so duplicate `Set-Cookie` +values survive that conversion. The disqualifying part is that the same +function drains `Body::Stream` into a buffered `fastly::Body` before sending, +which ends progressive streaming, and that the helper path leaves no place for +the response-extension finalization and post-send work this adapter performs. +That holds regardless of whether PR #379 merges, so this work takes no +dependency on it. + +`HandlerResult` is implemented for `()`, for handlers that have already called +`send_to_client` or `stream_to_client`. That is exactly the current handler +shape, so the existing send logic is reused unchanged. + +## Design + +### Three execution modes + +Compatibility is layered, and the layers are not equivalent. Only the first +reproduces today's execution. + +```text +feature off → original entry path +feature on, effective max_requests <= 1 → single-request handler path +feature on, validated reuse limits → Serve loop +``` + +The feature-on single-request path is **not** byte-identical to the feature-off +path: it performs the startup mode lookup, which opens a config store and reads +four keys before the first request. It does not construct a `Serve` or enter +its loop — it calls the handler once directly, as the code below shows. It is +described as _single-request operation_, not as identical execution. + +The startup mode lookup must not be able to break the health probe. If the +lookup fails for any reason, the adapter falls back to single-request operation +and the health probe continues to answer. + +### Commit 1 — lifecycle, no retention + +Add a `reusable-sandbox` Cargo feature to `trusted-server-adapter-fastly`, off +by default. `fastly.toml`'s build command does not pass it, so production +builds are unaffected. + +Split `main` into a loop owner and a handler: + +```rust +fn main() { + let mut sandbox = Sandbox::default(); + match serve_mode() { + ServeMode::Single => { handle_request(FastlyRequest::from_client(), &mut sandbox); } + ServeMode::Reuse(serve) => { let _ = serve.run_with_context(handle_request, &mut sandbox); } + } +} + +fn handle_request(req: FastlyRequest, sandbox: &mut Sandbox) { + // today's `main` body: health probe, logger, edgezero_main +} +``` + +`handle_request` returns `()`. It keeps sending its own response, so streaming, +duplicate `Set-Cookie`, response extensions, and post-send pull sync are +untouched. + +In this commit `Sandbox` holds no application state. It holds: + +- `logger_installed: bool`, which fixes the `logging.rs:105` double-`apply()` + panic. +- The measurement bookkeeping from + [Ownership of the measurement surface](#ownership-of-the-measurement-surface): + the validated guest-instance identifier, the request ordinal within that + instance, and the application build counter. + +The build counter lives here from commit 1 even though nothing increments it +past one until commit 3 — that is what makes arm B's "one build per request" +baseline measurable rather than assumed. + +Counters are attached to the response **before headers are committed**. On the +streaming path that is before `stream_to_client`, which means they describe the +request up to commitment and cannot report its eventual outcome. A failure +after commitment is therefore recorded separately, in logs, and is reconciled +with the counters during analysis rather than being expected to appear in +them. + +Every non-panicking path through `handle_request` must send exactly once. This +is already true of the current code (each early return sends before returning), +but it becomes a correctness requirement of the loop rather than an incidental +property of a process that is about to exit, so it is asserted by test. + +The loop owner introduces one new failure mode. `run_with_context` panics +(`serve.rs:320`) if `RequestPromise::new` fails mid-loop, which cannot happen in +today's one-shot `main`. This is accepted rather than worked around: it fires +only between requests, after the current response has been sent, and the SDK +treats it as a sandbox failure — the same outcome as any other guest panic. It +is recorded here so it is not mistaken for an application fault during +validation. + +### Commit 2 — per-document rewrite buffers + +A `trusted-server-core` prerequisite for retention, specified under +[Retained rewrite buffers](#retained-rewrite-buffers-blocking). No adapter +change and no behavioural change under today's one-request-per-sandbox model. + +### Commit 3 — retention + +`Sandbox` gains `app: Option`, where `RetainedApp` owns the built +`App` and its `Arc`. It is constructed lazily on the first request +that is not one of the three pre-build short-circuits — the health probe, the +JA4 debug probe, and the counters endpoint from +[Ownership of the measurement surface](#ownership-of-the-measurement-surface) — +so none of those paths pays for construction. + +A failed build is never stored. `router_with_state` already returns +`(startup_error_router, None)` on failure; that result serves the current +request and is dropped, so a transient config-store failure cannot poison the +sandbox. + +### Ownership + +| Retained in `Sandbox` | Rebuilt every request | +| --------------------------------------------------------------------------------------------------------------------- | ------------------------------------------------------------------- | +| `logger_installed` | `ConfigStoreHandle` (native handle) | +| `App` + `Arc`: settings, auction plan, orchestrator, integration registry, telemetry sink, default KV store | `EnvConfig` / `RuntimeStoreConfig` | +| | `ClientInfo`, `DeviceSignals`, TLS/JA4 metadata, resolved client IP | +| | `RuntimeServices` (`app.rs:285`) | +| | Correlation id from `get_client_request_id()` | +| | `EcFinalizeState`, `RequestFilterEffects`, auth results, bodies | + +Native store handles are never retained alongside the app. + +Correlation is request-local. `get_client_request_id()` supplies it; +`FASTLY_TRACE_ID` identifies the sandbox, not the request, and is never used as +a request id. Correlation values are never written into app state or into +global logger configuration. + +### Nested-state audit + +Retaining `AppState` is only safe if everything reachable from it is +config-derived and free of request state and native handles. + +- `AuctionOrchestrator` (`orchestrator.rs:307-317`): `enabled`, `plan_backed`, + `plan`, `planned_providers`, `mediator`. All config-derived. The per-auction + `PlannedLaunchState` is a local, not a field. +- `IntegrationRegistry` (`registry.rs:780-783`): `Arc` + plus an optional plan. **The original audit found request content — see + [Retained rewrite buffers](#retained-rewrite-buffers-blocking).** +- `FastlyTinybirdAuctionTelemetrySink` (`tinybird.rs:36-48`): owned strings and + a backend spec. No handles. +- `default_kv_store` is `UnavailableKvStore`, inert. + +Two process-wide statics outlive a request for the first time under reuse. +Neither is reachable from `AppState`, but both change behaviour: + +- `IP_CIDR_SOURCE_CACHE` (`protection_scope.rs:200`), covered under + [Behavioural tests](#behavioural-tests). +- `MISSING_GEO_WARNING_LOGGED` (`consent/mod.rs:64`), an `AtomicBool` that + becomes log-once-per-sandbox instead of log-once-per-request. Benign and + arguably the intent, but recorded here so a reduced warning count during + validation is not mistaken for a dropped warning. + +A sweep of the core and adapter crates for `static` interior mutability finds +only these two; everything else is immutable `LazyLock` regexes and sets. + +### Retained rewrite buffers (blocking) + +At the original design revision (`c56ff745e`), the registry was **not** safe to +retain. Two script rewriters held request content in interior-mutable state on +the rewriter object itself: + +- `GoogleTagManagerIntegration.accumulated_text: Mutex` + (`google_tag_manager.rs:366`, used at `:1033` at that revision) +- `NextJsNextDataRewriter.accumulated_text: Mutex` + (`nextjs/script_rewriter.rs:24`, used at `:78` at that revision; subsequently + replaced by per-document `rsc_stream::FragmentState` in #1135) + +Both accumulated fragments of an inline `