diff --git a/.claude/agents/pr-reviewer.md b/.claude/agents/pr-reviewer.md index 4f65a426f..66bc88c10 100644 --- a/.claude/agents/pr-reviewer.md +++ b/.claude/agents/pr-reviewer.md @@ -27,13 +27,13 @@ The agent resolves the input into exactly one of three **modes** before starting any fetch. Modes determine the per-invocation variables defined in step 1; later steps reference those variables and don't restate mode logic. -| Input | Mode (after resolution) | -|---|---| -| A PR number (e.g. `#165`) | **PR** | -| A branch name, and PR lookup returns exactly one matching PR | **PR** (via lookup) | -| A branch name, no PR exists | **BRANCH-REMOTE** — always. No probe on the current checkout. | -| User explicitly says "review my local working tree" | **BRANCH-LOCAL** — current checkout only. If the user also names a branch, the agent verifies it matches `git branch --show-current`; otherwise it stops and asks the user to either check out that branch first or drop the name. The agent does **not** silently review whatever HEAD happens to be. | -| No input | Run the PR lookup probe with `$(git branch --show-current)`. If it returns a PR → PR mode. Otherwise apply the no-input rule below. | +| Input | Mode (after resolution) | +| ------------------------------------------------------------ | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------ | +| A PR number (e.g. `#165`) | **PR** | +| A branch name, and PR lookup returns exactly one matching PR | **PR** (via lookup) | +| A branch name, no PR exists | **BRANCH-REMOTE** — always. No probe on the current checkout. | +| User explicitly says "review my local working tree" | **BRANCH-LOCAL** — current checkout only. If the user also names a branch, the agent verifies it matches `git branch --show-current`; otherwise it stops and asks the user to either check out that branch first or drop the name. The agent does **not** silently review whatever HEAD happens to be. | +| No input | Run the PR lookup probe with `$(git branch --show-current)`. If it returns a PR → PR mode. Otherwise apply the no-input rule below. | **Branch-to-PR lookup rule.** `gh pr list --head ` does not support `:` syntax, so fork PRs with the same branch name can collide. @@ -52,7 +52,7 @@ matches=$(gh pr list --head "$REQUESTED_HEAD" \ - `>1` matches → stop and ask the user which PR number to review; do not pick the first row. -**No-input / no-PR rule** (the only place an inferred BRANCH-* mode happens +**No-input / no-PR rule** (the only place an inferred BRANCH-\* mode happens — a named branch never triggers this probe; the user said the name, the agent honours it): @@ -60,7 +60,7 @@ agent honours it): configured **and** `git rev-list --left-right --count "@{upstream}...HEAD"` returns `0 0` → resolve to **BRANCH-REMOTE**, with `` bound to the branch name the upstream points at, not to the local branch name. The - probe approves the *upstream* state; the agent must fetch that exact + probe approves the _upstream_ state; the agent must fetch that exact branch, not `origin/$(git branch --show-current)` which could be a different ref (the local branch might track `origin/main-fork`). Because BRANCH-REMOTE fetches from `origin`, this inference only applies when the @@ -87,6 +87,7 @@ agent honours it): When `$REQUESTED_HEAD` is bound, BRANCH-REMOTE proceeds with that name. When it's not bound (upstream on a non-`origin` remote), the agent asks the user instead. + - Anything else → **ask the user** which mode they want. The probe uses only `git status` / `git rev-parse` on existing local refs — it does **not** fetch, so the choice is made before any network or worktree side effect. @@ -283,7 +284,7 @@ fi This block adds the worktree **at most once** per invocation, and reuses the one from a prior invocation rather than adding a second. The decision keys off -git's worktree *registry* — the source of truth — not just the directory's +git's worktree _registry_ — the source of truth — not just the directory's existence, because the two can disagree (a registered worktree whose directory was manually `rm -rf`'d, or a leftover directory git never registered). Keying off the directory alone would send a dangling-but-registered path into @@ -345,9 +346,9 @@ After this step every later step uses `${WT:-.}` for cwd (so BRANCH-LOCAL implicitly runs from the project root) and `$DIFF_RANGE` for diffs. There are no more per-mode forks until step 7e (which checks `[ -n "$WT" ]` for scratch verification) and step 8, where only **step 8b (the GitHub review -submission)** is skipped for BRANCH-* modes — no PR to submit a review to. +submission)** is skipped for BRANCH-\* modes — no PR to submit a review to. Step 8a still runs in every mode to compose the review artifact and, in -BRANCH-* modes, render the would-have-been verdict and findings into chat. +BRANCH-\* modes, render the would-have-been verdict and findings into chat. Stash the `HEAD_OID_EXPECTED` value — step 8 re-checks it immediately before submission and pins it into the review payload as `commit_id`. @@ -449,7 +450,7 @@ fi If the PR has passing CI checks, report them as PASS in the review. Only run CI locally if checks haven't run yet or if you need to verify a specific failure. Note any CI failures in the review but continue with the code review -regardless. (This governs CI *status reporting* only — the suggestion +regardless. (This governs CI _status reporting_ only — the suggestion scratch-verification in step 7e is independent and runs regardless of what GitHub's checks say.) @@ -647,7 +648,7 @@ numbers it covers; pick `line` (and `start_line` when multi-line) from inside the same hunk; don't span hunk boundaries. **Deleted lines and base-side context.** Findings about something the PR -*removed* don't live on the RIGHT side at all — there are no new-file lines +_removed_ don't live on the RIGHT side at all — there are no new-file lines to anchor to. Two options: - Anchor the inline comment on the LEFT side: `"side": "LEFT"` (and @@ -662,7 +663,7 @@ to anchor to. Two options: see step 1's "base-side reads" note.) Same rule for renamed files: the file's new path can carry a `suggestion` -block normally; comments about content the rename *also dropped* anchor on +block normally; comments about content the rename _also dropped_ anchor on the old path with `side: "LEFT"`. Use a `suggestion` block when: @@ -693,12 +694,12 @@ fenced code block when: cleanly revised. In all of those cases, give the proposed code in a plain fenced block (e.g. -```` ```rust ````) and end with a short "Apply manually — can't be auto-applied +` ```rust `) and end with a short "Apply manually — can't be auto-applied as a suggestion because …" sentence. ##### Fence length when the replacement itself contains backticks -The default `` ```suggestion `` fence is three backticks. If the replacement +The default ` ```suggestion ` fence is three backticks. If the replacement bytes contain a line that is itself a run of three-or-more backticks — common for Markdown/docs suggestions that include a nested code fence — that inner run closes the outer `suggestion` block early, and the rendered one-click @@ -706,14 +707,14 @@ suggestion is truncated or malformed rather than matching the bytes the user approved. Before displaying or submitting any suggestion: 1. Scan the replacement for the longest run of consecutive backticks, `N`. -2. If `N < 3`, use the normal three-backtick `` ```suggestion `` fence. +2. If `N < 3`, use the normal three-backtick ` ```suggestion ` fence. 3. If `N >= 3`, open and close the block with a fence of `N + 1` backticks - (e.g. ` ````suggestion ` for an inner ```` ``` ````), so the outer fence is + (e.g. ` ````suggestion ` for an inner ` ``` `), so the outer fence is strictly longer than any inner run — GitHub follows the CommonMark rule that a fence closes only on a run of **at least** as many backticks. This rule is deterministic; whether GitHub renders the widened fence as a one-click suggestion is a server-side property that cannot be checked locally (7e's - scratch pass verifies replacement *bytes*, not rendering). When in doubt — + scratch pass verifies replacement _bytes_, not rendering). When in doubt — e.g. an unusually exotic replacement — **demote the finding to prose-only** (a plain fenced block plus the "Apply manually …" sentence) rather than risk posting a malformed suggestion. @@ -753,7 +754,7 @@ For a multi-line suggestion, add `start_line` and `start_side`: **Indentation matters**: the block replaces the original lines verbatim, so leading whitespace must match exactly what the file expects after the fix. -**Fence length matters too**: the `` ```suggestion `` fences above use three +**Fence length matters too**: the ` ```suggestion ` fences above use three backticks, which only holds when the replacement contains no three-or-more backtick run of its own. When it does (e.g. a docs suggestion with a nested code fence), widen the outer fence per step 7a's fence-length rule or demote @@ -776,7 +777,7 @@ the comment body and tell the author it has to be applied manually: #### 7c-bis. Inline comment on a removed (LEFT-side) line -A finding about a line the PR *removed* has no RIGHT-side anchor — pin it +A finding about a line the PR _removed_ has no RIGHT-side anchor — pin it on the LEFT (base) side instead. `suggestion` blocks aren't applicable (GitHub only commits suggestions from the RIGHT side), so the body uses a plain code block: @@ -805,8 +806,8 @@ can't straddle sides. - The total number of inline comments has a soft cap of ~30. If you would exceed that, consolidate the lowest-severity findings into the review body with file/line references but no inline comment. -- A given inline comment may contain at most one ```` ```suggestion ```` block. - Prose context blocks (e.g. ```` ```rust ````) are fine alongside it. +- A given inline comment may contain at most one ` ```suggestion ` block. + Prose context blocks (e.g. ` ```rust `) are fine alongside it. - If the user changed an emoji tag during triage, the comment uses the new tag. - Don't post suggestions on lines outside the RIGHT side of the diff — they'll fail GitHub's "position could not be resolved" check. Comments about @@ -838,7 +839,7 @@ suggestion A only compiles because suggestion B was also applied, the agent has labelled A as verified but A-alone can break the build. So the inner loop tests each suggestion against a clean worktree first; a final batch pass (all approved suggestions applied together) is a nice-to-have to -catch *interactions*, but the per-suggestion runs are the real gate: +catch _interactions_, but the per-suggestion runs are the real gate: ```bash # Confirm clean starting state at $WT (HEAD = $HEAD_REF, status empty). @@ -931,8 +932,9 @@ gate when **any** of these is true: - The suggestion touches a `#[cfg(test)]` module, a test, or a feature gate. - The finding is 🔧 wrench (blocking) — release-blocking fixes must clear the release gate. -- The touched code is shared (`crates/trusted-server-core/src/{auction,ec, - http_util,publisher,html_processor,settings,constants}` and similar). +- The touched code is shared + (`crates/trusted-server-core/src/{auction,ec,http_util,publisher,html_processor,settings,constants}` + and similar). - The suggestion changes program behaviour and the agent prefers to ship it **without** the compile-verified-only disclaimer. @@ -967,11 +969,11 @@ at build time. Run the build whenever the suggestion touches files under **Post-verify drift check (snapshot the approved patch, hard-fail on any deviation).** Filename-level comparison isn't enough — a formatter or -codegen step can change a different range *inside* an approved file, while +codegen step can change a different range _inside_ an approved file, while the posted GitHub suggestion still contains only the originally-approved range. The correct check is byte-exact: snapshot the full patch immediately after applying the approved suggestions but **before** running any -verification command, then compare with the patch *after* verification. Any +verification command, then compare with the patch _after_ verification. Any delta — different range in the same file, an extra tracked file, a whitespace change in `Cargo.toml` from a build script — means verification mutated the tree beyond what the agent approved, and the suggestion as @@ -1060,7 +1062,7 @@ After cleanup, the worktree's HEAD must be at `$HEAD_REF` Step 8 splits into two halves. **8a** is mode-agnostic: determine the verdict and compose the review body + inline comments. **8b** is PR-only: post the -review to GitHub. BRANCH-* modes still produce the artifact in 8a (so step 10 +review to GitHub. BRANCH-\* modes still produce the artifact in 8a (so step 10 has a "would-have-been verdict" to report) but skip 8b — the artifact is rendered into the chat instead. @@ -1164,14 +1166,14 @@ inline comment", not "this finding is out of scope". Omit any section that has no findings — don't include empty headings. -In BRANCH-* modes (`[ -z "$NUMBER" ]`) the artifact is now complete: render it +In BRANCH-\* modes (`[ -z "$NUMBER" ]`) the artifact is now complete: render it into the chat exactly as the GitHub UI would have shown it (body markdown followed by each inline comment, labelled with the file:line it would have anchored to). Stop after rendering — there is no review to submit. #### 8b. Submit the GitHub review (PR mode only — `[ -n "$NUMBER" ]`) -Skip this entire sub-step in BRANCH-* modes. +Skip this entire sub-step in BRANCH-\* modes. ##### Re-check the PR head before submission @@ -1197,7 +1199,7 @@ fi If `SKIP_SUBMISSION` is set, the agent **must skip every command in the "Submit the review" sub-section below**, render the artifact in chat (the -same way BRANCH-* mode would in 8a), and report `submission skipped: +same way BRANCH-\* mode would in 8a), and report `submission skipped: $SKIP_REASON` as the stop reason in step 10. The agent does not restart inside the same invocation — single-invocation = single pass; the user re-invokes if they want another pass against the new head. @@ -1218,7 +1220,7 @@ fi ``` When `STOP_SUBMISSION=1`, render the composed artifact into chat exactly the -way BRANCH-* mode does in 8a, report `submission skipped: $SKIP_REASON` per +way BRANCH-\* mode does in 8a, report `submission skipped: $SKIP_REASON` per step 10, and do not run any of the `gh api` submit/delete commands below. Use the GitHub API to submit. Handle these known issues: @@ -1307,7 +1309,7 @@ Use the GitHub API to submit. Handle these known issues: fallthrough, no `exit 1`. **Re-gate after the pending-review check.** The `SKIP_SUBMISSION` guard - at the top of "Submit the review" runs *before* the pending-review check, + at the top of "Submit the review" runs _before_ the pending-review check, so a `keep` decision must be caught by a second gate immediately after the case block — otherwise the JSON-assembly / `gh api … -X POST` below would still run: @@ -1349,7 +1351,7 @@ head-recheck above. **Invariant — never heredoc-interpolate user content into the payload.** Inline-comment bodies routinely contain `$`, backticks, backslashes, -quoted code, and the literal `` ```suggestion `` fence. A `cat <` (binary name `ts`), or install it with - `cargo install --path crates/trusted-server-cli`. +- The `ts` CLI runs from source without installing: + `cargo run -p trusted-server-cli -- config ` (binary name `ts`), or + install it with `cargo install --path crates/trusted-server-cli`. - The app-config store's logical id and blob key are `trusted_server_config` (`settings_data.rs` `DEFAULT_CONFIG_STORE_ID`, `config_payload.rs` `CONFIG_BLOB_KEY`); older builds used `app_config`. The override env-var key diff --git a/.github/workflows/format.yml b/.github/workflows/format.yml index 7253b5ea2..93bc99229 100644 --- a/.github/workflows/format.yml +++ b/.github/workflows/format.yml @@ -149,5 +149,11 @@ jobs: - name: Run Prettier (check) run: npm run format + - name: Run Prettier (check) — Markdown outside docs/ + working-directory: . + run: >- + docs/node_modules/.bin/prettier --config docs/.prettierrc --check + "*.md" ".claude/**/*.md" ".github/**/*.md" "crates/**/*.md" "scripts/**/*.md" "tinybird/**/*.md" + - name: Build with VitePress (fails on dead links) run: npm run build diff --git a/.gitignore b/.gitignore index 24b9e06aa..238e6b133 100644 --- a/.gitignore +++ b/.gitignore @@ -63,3 +63,7 @@ src/*.html # leftover local build artifacts (node_modules, target, dist) that remain on disk. /crates/js/ /crates/integration-tests/ + +# Wrangler config generated by the Cloudflare integration-test harness from +# wrangler.toml at run time; regenerated on every run. +wrangler.integration.generated.toml diff --git a/AGENTS.md b/AGENTS.md index b6e61ecf9..738fcab34 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -329,14 +329,14 @@ IntegrationRegistration::builder(ID) ## Configuration Files -| File | Purpose | -| --------------------- | ---------------------------------------------------------- | -| `edgezero.toml` | EdgeZero app/platform manifest and logical stores | -| `fastly.toml` | Fastly service configuration and build settings | -| `trusted-server.example.toml` | Source-controlled Trusted Server app-config template | -| `trusted-server.toml` | Operator-owned app config; gitignored; `ts config push` publishes it as an EdgeZero blob envelope | -| `rust-toolchain.toml` | Pins Rust version to 1.95.0 | -| `.env.dev` | Local development environment variables | +| File | Purpose | +| ----------------------------- | ------------------------------------------------------------------------------------------------- | +| `edgezero.toml` | EdgeZero app/platform manifest and logical stores | +| `fastly.toml` | Fastly service configuration and build settings | +| `trusted-server.example.toml` | Source-controlled Trusted Server app-config template | +| `trusted-server.toml` | Operator-owned app config; gitignored; `ts config push` publishes it as an EdgeZero blob envelope | +| `rust-toolchain.toml` | Pins Rust version to 1.95.0 | +| `.env.dev` | Local development environment variables | --- @@ -351,6 +351,7 @@ Every PR must pass: 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`) 7. Docs format (`cd docs && npm run format`) +8. Markdown format outside `docs/` (requires `cd docs && npm ci` first): `docs/node_modules/.bin/prettier --config docs/.prettierrc --check "*.md" ".claude/**/*.md" ".github/**/*.md" "crates/**/*.md" "scripts/**/*.md" "tinybird/**/*.md"`; fix with `--write` in place of `--check` --- @@ -439,18 +440,18 @@ both runtime behavior and build/tooling changes. ## Key Files -| File | Purpose | -| -------------------------------------------- | ------------------------------------------------- | -| `crates/trusted-server-core/src/integrations/registry.rs` | IntegrationRegistry, `js_module_ids()` | -| `crates/trusted-server-core/src/tsjs.rs` | Script tag generation with module IDs | -| `crates/trusted-server-core/src/html_processor.rs` | Injects `"), + "other", + "an unbounded client-controlled token must not reach the row verbatim" + ); + } + + #[test] + fn row_normalizes_method_even_when_snapshot_carries_a_raw_token() { + // The normalizer runs inside `access_event_row` so every row-building + // path is covered, regardless of what the snapshot's `method` field + // holds — a caller-controlled extension method must never leak into + // the row unnormalized. + let mut snapshot = unknown_snapshot(RouteClass::Other, "/other/*"); + snapshot.method = "PROPFIND".to_owned(); + let row = access_event_row(&snapshot, &TimingSnapshot::default(), 0); + let parsed: serde_json::Value = + serde_json::from_str(&row).expect("should serialize valid JSON"); + + assert_eq!(parsed["method"], "other"); + } +} diff --git a/crates/trusted-server-core/src/auction/endpoints.rs b/crates/trusted-server-core/src/auction/endpoints.rs index ab3585e3d..355d7d400 100644 --- a/crates/trusted-server-core/src/auction/endpoints.rs +++ b/crates/trusted-server-core/src/auction/endpoints.rs @@ -22,6 +22,7 @@ use crate::ec::registry::PartnerRegistry; use crate::error::TrustedServerError; use crate::openrtb::{Eid, Uid}; use crate::platform::RuntimeServices; +use crate::request_timing::RequestTimings; use crate::settings::Settings; use super::AuctionOrchestrator; @@ -145,6 +146,16 @@ pub async fn handle_auction( } let (parts, body) = req.into_parts(); + // T0-anchored timeline (spec section 18). This route is the auction, so + // dispatch and resolve bracket `run_auction` rather than the origin + // fetch, and the commit mark lands once the OpenRTB response carrying the + // targeting has been built. A defaulted handle records into nothing that + // is ever read, so direct-handler tests are unaffected. + let timings = parts + .extensions + .get::() + .cloned() + .unwrap_or_default(); let body_bytes = body.into_bytes().unwrap_or_default(); if body_bytes.len() > MAX_AUCTION_BODY_SIZE { return Response::builder() @@ -362,10 +373,17 @@ pub async fn handle_auction( ec_context, ); + timings.set_auction_id(observation.auction_id); + // Run the auction + timings.mark_auction_dispatched(); let result = match orchestrator.run_auction(&auction_request, &context).await { - Ok(result) => result, + Ok(result) => { + timings.mark_auction_resolved(); + result + } Err(err) => { + timings.mark_auction_resolved(); let elapsed_ms = observation.elapsed_ms(); emit_auction_events_best_effort_lazy(services, || { build_auction_events( @@ -410,6 +428,10 @@ pub async fn handle_auction( } }; + // Targeting is available to the caller: for this route the response body + // is the commit, since there is no page state to write into. + timings.mark_auction_committed(); + emit_auction_events_best_effort_lazy(services, || { build_auction_events( observation, diff --git a/crates/trusted-server-core/src/auction/openrtb.rs b/crates/trusted-server-core/src/auction/openrtb.rs index b361386de..74553c9bc 100644 --- a/crates/trusted-server-core/src/auction/openrtb.rs +++ b/crates/trusted-server-core/src/auction/openrtb.rs @@ -16,7 +16,7 @@ use super::profile::{ use super::routing::{ PrebidTransportHeaders, ProviderAuctionInput, ProviderSlotInput, RoutedAuction, }; -use super::types::{AdFormat, AuctionResponse, Bid}; +use super::types::{AdFormat, AdSlot, AuctionResponse, Bid}; use crate::consent::ConsentSource; use crate::error::TrustedServerError; use crate::openrtb::{ @@ -154,11 +154,20 @@ pub(crate) type BidDimensionIndex = BTreeMap; /// Build the requested-dimension index once for one provider response. pub(crate) fn build_bid_dimension_index(input: &ProviderAuctionInput) -> BidDimensionIndex { + build_bid_dimension_index_from_slots(input.slots().iter().map(ProviderSlotInput::slot)) +} + +/// Build the requested-dimension index from plain [`AdSlot`]s. +/// +/// The first slot wins when several share an ID. +pub(crate) fn build_bid_dimension_index_from_slots<'a>( + slots: impl IntoIterator, +) -> BidDimensionIndex { let mut index = BidDimensionIndex::new(); - for slot in input.slots() { + for slot in slots { index - .entry(slot.slot().id.clone()) - .or_insert_with(|| SlotBidDimensions::from_formats(slot.slot().formats.as_slice())); + .entry(slot.id.clone()) + .or_insert_with(|| SlotBidDimensions::from_formats(slot.formats.as_slice())); } index } diff --git a/crates/trusted-server-core/src/auction/orchestrator.rs b/crates/trusted-server-core/src/auction/orchestrator.rs index 204201e60..8d5e5952f 100644 --- a/crates/trusted-server-core/src/auction/orchestrator.rs +++ b/crates/trusted-server-core/src/auction/orchestrator.rs @@ -4171,7 +4171,8 @@ mod tests { .expect("current adapters should accept completed late mediator responses"); assert!( current_mediator.response_time_ms >= 50, - "mediator timing should preserve actual elapsed duration" + "mediator timing should preserve actual elapsed duration, got {} ms", + current_mediator.response_time_ms ); assert_eq!(current.winning_bids["slot-1"].bidder, "mediated"); diff --git a/crates/trusted-server-core/src/config.rs b/crates/trusted-server-core/src/config.rs index 79f100e2f..24d407121 100644 --- a/crates/trusted-server-core/src/config.rs +++ b/crates/trusted-server-core/src/config.rs @@ -182,6 +182,10 @@ impl edgezero_core::app_config::AppConfigMeta for TrustedServerAppConfig { vec![optional_object("tinybird"), object("auction_token_secret")], true, ), + field( + vec![optional_object("tinybird"), object("access_token_secret")], + true, + ), field( vec![ optional_object("integrations"), @@ -421,7 +425,7 @@ fn validate_secret_key_references(settings: &Settings) -> Result<(), Report Result<(), Report("datadome")? { if datadome.enable_protection { @@ -830,6 +842,7 @@ formats = [{ width = 300, height = 250 }] ("handlers[*].password".to_owned(), false), ("trusted_client_ip.shared_secret".to_owned(), false), ("tinybird.auction_token_secret".to_owned(), true), + ("tinybird.access_token_secret".to_owned(), true), ( "integrations.datadome.server_side_key_secret_name".to_owned(), true, @@ -1252,10 +1265,14 @@ password = "production-admin-password-32-bytes" ); } - /// Integrations that default to disabled do not validate inactive fields. + /// `enabled` defaults to `false`, so a section that omits the flag resolves + /// to disabled and must not have its fields validated. #[test] fn deploy_validation_skips_field_validation_for_integrations_with_omitted_enabled() { let mut settings = valid_settings(); + // `endpoint` parses as a plain string but would fail the `url` + // validator, so this section only survives if validation is skipped for + // integrations that resolve to disabled. settings .integrations .insert_config( diff --git a/crates/trusted-server-core/src/config_payload.rs b/crates/trusted-server-core/src/config_payload.rs index c9690b01f..f530968a4 100644 --- a/crates/trusted-server-core/src/config_payload.rs +++ b/crates/trusted-server-core/src/config_payload.rs @@ -72,16 +72,30 @@ pub fn settings_from_config_blob( } fn remove_inactive_secret_references(data: &mut serde_json::Value) { - if data - .pointer("/tinybird/enabled") - .and_then(serde_json::Value::as_bool) - != Some(true) - && let Some(tinybird) = data - .get_mut("tinybird") - .and_then(serde_json::Value::as_object_mut) + if let Some(tinybird) = data + .get_mut("tinybird") + .and_then(serde_json::Value::as_object_mut) { - tinybird.remove("auction_token_secret"); - tinybird.remove("access_token_secret"); + let enabled = tinybird.get("enabled").and_then(serde_json::Value::as_bool) == Some(true); + if !enabled { + tinybird.remove("auction_token_secret"); + tinybird.remove("access_token_secret"); + } else { + if tinybird + .get("auction_enabled") + .and_then(serde_json::Value::as_bool) + == Some(false) + { + tinybird.remove("auction_token_secret"); + } + if tinybird + .get("access_enabled") + .and_then(serde_json::Value::as_bool) + != Some(true) + { + tinybird.remove("access_token_secret"); + } + } } if let Some(partners) = data @@ -138,7 +152,10 @@ fn json_bool_or_string_is_true(value: Option<&serde_json::Value>) -> bool { #[cfg(test)] mod tests { + use std::sync::Arc; + use super::*; + use crate::integrations::IntegrationRegistry; use crate::integrations::didomi::DidomiIntegrationConfig; use crate::platform::{PlatformError, StoreId}; use crate::redacted::Redacted; @@ -845,9 +862,15 @@ mod tests { #[test] fn runtime_blob_accepts_disabled_browser_bidder_ownership_overlap() { let original = settings_with_browser_bidder_overlap(false); + let reconstructed = load_settings(&envelope_json(&original)) + .expect("should decode dormant conflicting runtime blob"); + let plan = Arc::new( + crate::auction::compile_auction_plan(&reconstructed) + .expect("should compile decoded disabled auction plan"), + ); - load_settings(&envelope_json(&original)) - .expect("runtime should accept disabled browser bidder ownership overlap"); + IntegrationRegistry::with_plan(&reconstructed, plan) + .expect("runtime registry should accept disabled ownership overlap"); } #[test] diff --git a/crates/trusted-server-core/src/constants.rs b/crates/trusted-server-core/src/constants.rs index 1b55679e2..8e39f2d19 100644 --- a/crates/trusted-server-core/src/constants.rs +++ b/crates/trusted-server-core/src/constants.rs @@ -38,6 +38,8 @@ pub const HEADER_X_TS_ENV: HeaderName = HeaderName::from_static("x-ts-env"); // Fastly environment variables pub const ENV_FASTLY_SERVICE_VERSION: &str = "FASTLY_SERVICE_VERSION"; pub const ENV_FASTLY_IS_STAGING: &str = "FASTLY_IS_STAGING"; +pub const ENV_FASTLY_SERVICE_ID: &str = "FASTLY_SERVICE_ID"; +pub const ENV_FASTLY_POP: &str = "FASTLY_POP"; // Common standard header names used across modules pub const HEADER_USER_AGENT: HeaderName = HeaderName::from_static("user-agent"); diff --git a/crates/trusted-server-core/src/ec/kv.rs b/crates/trusted-server-core/src/ec/kv.rs index b0eb83d3f..8653f2342 100644 --- a/crates/trusted-server-core/src/ec/kv.rs +++ b/crates/trusted-server-core/src/ec/kv.rs @@ -1717,6 +1717,28 @@ mod tests { assert!(ts > 0, "should return a nonzero timestamp"); } + #[test] + fn kv_span_accumulates_across_graph_operations() { + let timings = crate::request_timing::RequestTimings::new(); + let graph = KvIdentityGraph::new(crate::platform::TimedKvStore::new( + crate::ec::kv_backend::test_support::InMemoryEcKv::new("test-store"), + timings.clone(), + )); + + graph + .create("ec-1", &live_entry()) + .expect("should create entry through the timed store"); + graph + .get("ec-1") + .expect("should read the entry back through the timed store"); + + timings.mark_headers_ready(); + assert!( + timings.snapshot().kv_ms.is_some(), + "should accumulate Phase::EcKv across both graph operations, not just the last write" + ); + } + #[test] fn serialize_entry_produces_valid_json() { let entry = KvEntry::tombstone(1000); diff --git a/crates/trusted-server-core/src/geo.rs b/crates/trusted-server-core/src/geo.rs index 63f7907f5..fe5785d26 100644 --- a/crates/trusted-server-core/src/geo.rs +++ b/crates/trusted-server-core/src/geo.rs @@ -48,6 +48,24 @@ impl GeoInfo { } } +/// Carries the outcome of a request-phase geo lookup across to +/// response-phase finalization, so a finalize consumer can reuse it instead +/// of performing a second lookup for the same request. +/// +/// Attached as a response extension on every exit path that attempted a +/// lookup, including the asset-route fallback (which does not carry an EC +/// finalize state). +#[derive(Debug, Clone)] +pub enum GeoLookupState { + /// No lookup has been attempted for this request. + NotAttempted, + /// A lookup ran and failed (or returned no result). This must not be + /// retried: finalize treats it the same as no geo info being available. + Attempted, + /// A lookup ran and resolved geo info. + Resolved(GeoInfo), +} + fn insert_geo_header(headers: &mut http::HeaderMap, name: http::header::HeaderName, value: &str) { match HeaderValue::from_str(value) { Ok(header_value) => { diff --git a/crates/trusted-server-core/src/integrations/adserver_mock.rs b/crates/trusted-server-core/src/integrations/adserver_mock.rs index d4fbd3231..640574e03 100644 --- a/crates/trusted-server-core/src/integrations/adserver_mock.rs +++ b/crates/trusted-server-core/src/integrations/adserver_mock.rs @@ -16,6 +16,10 @@ use std::time::Duration; use validator::Validate; use crate::auction::context::{ContextQueryParams, build_url_with_context_params}; +use crate::auction::openrtb::{ + BidDimensionIndex, BidRejectionReason, build_bid_dimension_index_from_slots, + parse_optional_bid_dimension, resolve_bid_dimensions, +}; use crate::auction::provider::{AuctionProvider, ProviderRequestOutcome}; use crate::auction::types::{ AuctionContext, AuctionRequest, AuctionResponse, Bid, BidStatus, MediaType, @@ -264,6 +268,7 @@ impl AdServerMockProvider { json: &Json, response_time_ms: u64, bid_index: &BidIndex, + dimensions_by_slot: Option<&BidDimensionIndex>, ) -> AuctionResponse { let empty_array = vec![]; let seatbid = json["seatbid"].as_array().unwrap_or(&empty_array); @@ -292,14 +297,32 @@ impl AdServerMockProvider { let restored_bidder = original.map_or_else(|| seat_name.to_string(), |b| b.bidder.clone()); - let width = bid["w"].as_u64().unwrap_or(0) as u32; - let height = bid["h"].as_u64().unwrap_or(0) as u32; - if width == 0 || height == 0 { - log::debug!( - "adserver_mock: bid for slot '{slot_id}' has zero dimension ({width}×{height}), skipping" - ); - continue; - } + // Reuse the shared `OpenRTB` dimension parser and slot-format + // validation so this provider admits the same dimensions as the + // other auction providers. Keep `{:?}` for the raw upstream + // values below: Debug escapes newlines and quotes, which + // prevents log injection. + let dimensions = match ( + parse_optional_bid_dimension(bid, "w"), + parse_optional_bid_dimension(bid, "h"), + ) { + (Ok(width), Ok(height)) => match dimensions_by_slot { + Some(index) => resolve_bid_dimensions(index, &slot_id, width, height), + None => width.zip(height).ok_or(BidRejectionReason::InvalidBid), + }, + _ => Err(BidRejectionReason::InvalidBid), + }; + let (width, height) = match dimensions { + Ok(dimensions) => dimensions, + Err(reason) => { + log::debug!( + "adserver_mock: bid for slot '{slot_id}' has unusable dimensions {:?}×{:?} ({reason:?}), skipping", + bid["w"], + bid["h"] + ); + continue; + } + }; all_bids.push(Bid { slot_id, @@ -360,6 +383,7 @@ impl AdServerMockProvider { response: PlatformResponse, response_time_ms: u64, bid_index: &BidIndex, + dimensions_by_slot: Option<&BidDimensionIndex>, ) -> Result> { let response = response.response; @@ -385,8 +409,12 @@ impl AdServerMockProvider { log::trace!("AdServer Mock response: {:?}", response_json); - let auction_response = - self.parse_mediation_response(&response_json, response_time_ms, bid_index); + let auction_response = self.parse_mediation_response( + &response_json, + response_time_ms, + bid_index, + dimensions_by_slot, + ); log::info!( "AdServer Mock returned {} bids in {}ms", @@ -505,7 +533,7 @@ impl AuctionProvider for AdServerMockProvider { // [`parse_response_with_context`], so this path only serves callers // outside the orchestration flow. log::debug!("adserver_mock: parsing without context — SSP bid metadata unavailable"); - self.parse_response_inner(response, response_time_ms, &BidIndex::new()) + self.parse_response_inner(response, response_time_ms, &BidIndex::new(), None) .await } @@ -513,15 +541,21 @@ impl AuctionProvider for AdServerMockProvider { &self, response: PlatformResponse, response_time_ms: u64, - _request: &AuctionRequest, + request: &AuctionRequest, context: &AuctionContext<'_>, ) -> Result> { // Rebuild the SSP-bid lookup from the orchestrator-provided bidder // responses so nurl/burl/ad_id survive mediation. Request-scoped data // travels on the context instead of provider-instance state. let bid_index = build_bid_index(context.provider_responses.unwrap_or(&[])); - self.parse_response_inner(response, response_time_ms, &bid_index) - .await + let dimensions_by_slot = build_bid_dimension_index_from_slots(&request.slots); + self.parse_response_inner( + response, + response_time_ms, + &bid_index, + Some(&dimensions_by_slot), + ) + .await } fn supports_media_type(&self, media_type: &MediaType) -> bool { @@ -789,7 +823,7 @@ mod tests { }); let auction_response = - provider.parse_mediation_response(&mediation_response, 200, &BidIndex::new()); + provider.parse_mediation_response(&mediation_response, 200, &BidIndex::new(), None); assert_eq!(auction_response.provider, "adserver_mock"); assert_eq!(auction_response.status, BidStatus::Success); @@ -836,7 +870,8 @@ mod tests { ] }); - let response = provider.parse_mediation_response(&mediation_response, 10, &BidIndex::new()); + let response = + provider.parse_mediation_response(&mediation_response, 10, &BidIndex::new(), None); assert_eq!(response.bids.len(), 2); assert_eq!(response.bids[0].bidder, "provider-instance"); @@ -912,7 +947,7 @@ mod tests { ); let auction_response = - provider.parse_mediation_response(&mediation_response, 42, &bid_index); + provider.parse_mediation_response(&mediation_response, 42, &bid_index, None); assert_eq!(auction_response.status, BidStatus::Success); assert_eq!(auction_response.bids.len(), 1); @@ -1016,7 +1051,7 @@ mod tests { ); let auction_response = - provider.parse_mediation_response(&mediation_response, 42, &bid_index); + provider.parse_mediation_response(&mediation_response, 42, &bid_index, None); assert_eq!( auction_response.bids[0].bid_id.as_deref(), @@ -1080,6 +1115,7 @@ mod tests { }), 2, &reduced_index, + None, ); let winner = mediated .bids @@ -1111,7 +1147,7 @@ mod tests { }); let auction_response = - provider.parse_mediation_response(&mediation_response, 100, &BidIndex::new()); + provider.parse_mediation_response(&mediation_response, 100, &BidIndex::new(), None); assert_eq!(auction_response.status, BidStatus::NoBid); assert_eq!(auction_response.bids.len(), 0); @@ -1304,7 +1340,7 @@ mod tests { }); let auction_response = - provider.parse_mediation_response(&mediation_response, 200, &BidIndex::new()); + provider.parse_mediation_response(&mediation_response, 200, &BidIndex::new(), None); assert_eq!(auction_response.status, BidStatus::Success); assert_eq!(auction_response.bids.len(), 2); @@ -1323,6 +1359,359 @@ mod tests { ); } + #[test] + fn test_parse_mediation_response_skips_oversized_dimensions() { + // A dimension above u32::MAX must be rejected rather than silently + // wrapped into a small, plausible-looking value. The offset of 101 is + // load-bearing: u32::MAX + 1 truncates to 0, which is already rejected + // as a zero dimension, so it would not pin this fix. + let config = AdServerMockConfig::default(); + let provider = AdServerMockProvider::new(config); + + let oversized_dimension = u64::from(u32::MAX) + 101; + let mediation_response = json!({ + "id": "test-auction-123", + "seatbid": [ + { + "seat": "test-bidder", + "bid": [ + { + "id": "bid-oversized-width", + "impid": "header-banner", + "price": 3.50, + "adm": "
Oversized width
", + "w": oversized_dimension, + "h": 90, + }, + { + "id": "bid-oversized-height", + "impid": "sidebar", + "price": 1.25, + "adm": "
Oversized height
", + "w": 300, + "h": oversized_dimension, + }, + { + "id": "bid-valid", + "impid": "footer", + "price": 2.00, + "adm": "
Valid Ad
", + "w": 728, + "h": 90, + } + ] + } + ], + "cur": "USD" + }); + + let auction_response = + provider.parse_mediation_response(&mediation_response, 200, &BidIndex::new(), None); + + assert_eq!( + auction_response.bids.len(), + 1, + "Bids with oversized w/h should be skipped, only the valid bid should remain" + ); + assert_eq!(auction_response.bids[0].slot_id, "footer"); + assert_eq!(auction_response.bids[0].width, 728); + assert_eq!(auction_response.bids[0].height, 90); + } + + #[test] + fn test_parse_mediation_response_skips_zero_dimensions() { + // Zero or missing dimensions must keep their existing skip behavior. + let config = AdServerMockConfig::default(); + let provider = AdServerMockProvider::new(config); + + let mediation_response = json!({ + "id": "test-auction-123", + "seatbid": [ + { + "seat": "test-bidder", + "bid": [ + { + "id": "bid-zero-width", + "impid": "header-banner", + "price": 3.50, + "adm": "
Zero width
", + "w": 0, + "h": 90, + }, + { + "id": "bid-zero-height", + "impid": "sidebar", + "price": 1.25, + "adm": "
Zero height
", + "w": 300, + "h": 0, + }, + { + "id": "bid-missing-dimensions", + "impid": "skyscraper", + "price": 1.10, + "adm": "
Missing dimensions
", + }, + { + "id": "bid-valid", + "impid": "footer", + "price": 2.00, + "adm": "
Valid Ad
", + "w": 728, + "h": 90, + } + ] + } + ], + "cur": "USD" + }); + + let auction_response = + provider.parse_mediation_response(&mediation_response, 200, &BidIndex::new(), None); + + let slots: Vec<&str> = auction_response + .bids + .iter() + .map(|bid| bid.slot_id.as_str()) + .collect(); + assert_eq!( + slots, + ["footer"], + "should drop the zero-width, zero-height, and missing-dimension bids" + ); + } + + #[test] + fn test_parse_mediation_response_accepts_u32_max_dimensions() { + // u32::MAX is the largest representable dimension and must be accepted, + // pinning the upper boundary of the oversized-dimension check. + let config = AdServerMockConfig::default(); + let provider = AdServerMockProvider::new(config); + + let mediation_response = json!({ + "id": "test-auction-123", + "seatbid": [ + { + "seat": "test-bidder", + "bid": [ + { + "id": "bid-max-dimensions", + "impid": "header-banner", + "price": 3.50, + "adm": "
Max dimensions
", + "w": u32::MAX, + "h": u32::MAX, + } + ] + } + ], + "cur": "USD" + }); + + let auction_response = + provider.parse_mediation_response(&mediation_response, 200, &BidIndex::new(), None); + + assert_eq!( + auction_response.bids.len(), + 1, + "should accept a bid whose w/h equal u32::MAX" + ); + assert_eq!( + auction_response.bids[0].width, + u32::MAX, + "should keep width at u32::MAX" + ); + assert_eq!( + auction_response.bids[0].height, + u32::MAX, + "should keep height at u32::MAX" + ); + } + + #[test] + fn test_parse_mediation_response_skips_negative_dimensions() { + // Negative dimensions are not valid u64 values and must be skipped. + let config = AdServerMockConfig::default(); + let provider = AdServerMockProvider::new(config); + + let mediation_response = json!({ + "id": "test-auction-123", + "seatbid": [ + { + "seat": "test-bidder", + "bid": [ + { + "id": "bid-negative-width", + "impid": "header-banner", + "price": 3.50, + "adm": "
Negative width
", + "w": -1, + "h": 90, + }, + { + "id": "bid-negative-height", + "impid": "sidebar", + "price": 1.25, + "adm": "
Negative height
", + "w": 300, + "h": -1, + }, + { + "id": "bid-valid", + "impid": "footer", + "price": 2.00, + "adm": "
Valid Ad
", + "w": 728, + "h": 90, + } + ] + } + ], + "cur": "USD" + }); + + let auction_response = + provider.parse_mediation_response(&mediation_response, 200, &BidIndex::new(), None); + + let slots: Vec<&str> = auction_response + .bids + .iter() + .map(|bid| bid.slot_id.as_str()) + .collect(); + assert_eq!( + slots, + ["footer"], + "should drop the negative-width and negative-height bids" + ); + } + + #[test] + fn test_parse_mediation_response_accepts_integral_float_dimensions() { + // Integral floats (`300.0`) are a legitimate dimension encoding, matching + // the shared `OpenRTB` parser; fractional floats are still skipped. + let config = AdServerMockConfig::default(); + let provider = AdServerMockProvider::new(config); + + let mediation_response = json!({ + "id": "test-auction-123", + "seatbid": [ + { + "seat": "test-bidder", + "bid": [ + { + "id": "bid-float-dimensions", + "impid": "header-banner", + "price": 3.50, + "adm": "
Float dimensions
", + "w": 300.0, + "h": 250.0, + }, + { + "id": "bid-fractional-width", + "impid": "sidebar", + "price": 1.25, + "adm": "
Fractional width
", + "w": 300.5, + "h": 250, + } + ] + } + ], + "cur": "USD" + }); + + let auction_response = + provider.parse_mediation_response(&mediation_response, 200, &BidIndex::new(), None); + + assert_eq!( + auction_response.bids.len(), + 1, + "should accept the integral-float bid and skip the fractional one" + ); + assert_eq!( + auction_response.bids[0].slot_id, "header-banner", + "should keep the integral-float bid" + ); + assert_eq!( + ( + auction_response.bids[0].width, + auction_response.bids[0].height + ), + (300, 250), + "should convert integral floats to u32 dimensions" + ); + } + + #[test] + fn test_parse_mediation_response_validates_dimensions_against_slots() { + // With the request slots available, mediated bids must match a + // requested format, and missing dimensions are inferred from the + // slot's single format, matching the other auction providers. + let config = AdServerMockConfig::default(); + let provider = AdServerMockProvider::new(config); + let request = create_test_auction_request(); + let dimensions_by_slot = build_bid_dimension_index_from_slots(&request.slots); + + let mediation_response = json!({ + "id": "test-auction-123", + "seatbid": [ + { + "seat": "test-bidder", + "bid": [ + { + "id": "bid-exact", + "impid": "header-banner", + "price": 3.50, + "w": 728, + "h": 90, + }, + { + "id": "bid-mismatch", + "impid": "header-banner", + "price": 3.00, + "w": 300, + "h": 250, + }, + { + "id": "bid-inferred", + "impid": "header-banner", + "price": 2.50, + }, + { + "id": "bid-unrequested", + "impid": "sidebar", + "price": 1.25, + "w": 728, + "h": 90, + } + ] + } + ], + "cur": "USD" + }); + + let auction_response = provider.parse_mediation_response( + &mediation_response, + 200, + &BidIndex::new(), + Some(&dimensions_by_slot), + ); + + let admitted: Vec<(Option<&str>, u32, u32)> = auction_response + .bids + .iter() + .map(|bid| (bid.bid_id.as_deref(), bid.width, bid.height)) + .collect(); + assert_eq!( + admitted, + [ + (Some("bid-exact"), 728, 90), + (Some("bid-inferred"), 728, 90) + ], + "should keep the exact and inferred bids and drop the mismatched and unrequested ones" + ); + } + #[test] fn test_build_endpoint_url_with_context_query_params() { let config = AdServerMockConfig { diff --git a/crates/trusted-server-core/src/integrations/aps.rs b/crates/trusted-server-core/src/integrations/aps.rs index 4f45fbeb9..53b30ed6f 100644 --- a/crates/trusted-server-core/src/integrations/aps.rs +++ b/crates/trusted-server-core/src/integrations/aps.rs @@ -78,42 +78,54 @@ const APS_RENDERER_DOCUMENT: &str = r#" var match=/^#tsaps=([A-Za-z0-9_-]{22,128})$/.exec(location.hash); var expected=match&&match[1]; try{history.replaceState(null,'',location.pathname+location.search);}catch(_error){} -if(!expected)return; +var reported=false; +function report(reason,nonce){ + if(reported)return; + reported=true; + try{parent.postMessage({message:'trusted-server/aps/renderer-failed',nonce:nonce,reason:reason},'*');}catch(_error){} +} +if(!expected){report('bad_hash');return;} function keys(value,expectedKeys){ if(!value||typeof value!=='object'||Array.isArray(value))return false; var actual=Object.keys(value).sort(); return actual.length===expectedKeys.length&&actual.every(function(key,index){return key===expectedKeys[index];}); } -function validRenderer(renderer){ +function rendererProblem(renderer){ if(!keys(renderer,['aaxResponse','accountId','bidId','creativeId','creativeUrl','height','tagType','type','version','width'])&& - !keys(renderer,['aaxResponse','accountId','bidId','creativeUrl','height','tagType','type','version','width']))return false; - if(renderer.type!=='aps'||renderer.version!==1||typeof renderer.accountId!=='string'||!renderer.accountId||new TextEncoder().encode(renderer.accountId).length>1024)return false; - if(typeof renderer.bidId!=='string'||!renderer.bidId||!Number.isInteger(renderer.width)||renderer.width<=0||!Number.isInteger(renderer.height)||renderer.height<=0)return false; - if(Object.prototype.hasOwnProperty.call(renderer,'creativeId')&&(typeof renderer.creativeId!=='string'||!renderer.creativeId||new TextEncoder().encode(renderer.creativeId).length>1024))return false; - if(renderer.tagType!=='iframe'&&renderer.tagType!=='script')return false; - if(typeof renderer.creativeUrl!=='string'||new TextEncoder().encode(renderer.creativeUrl).length>4096)return false; - if(typeof renderer.aaxResponse!=='string'||!renderer.aaxResponse||renderer.aaxResponse.length>349528)return false; + !keys(renderer,['aaxResponse','accountId','bidId','creativeUrl','height','tagType','type','version','width']))return 'descriptor_keys'; + if(renderer.type!=='aps'||renderer.version!==1||typeof renderer.accountId!=='string'||!renderer.accountId||new TextEncoder().encode(renderer.accountId).length>1024)return 'descriptor_fields'; + if(typeof renderer.bidId!=='string'||!renderer.bidId||!Number.isInteger(renderer.width)||renderer.width<=0||!Number.isInteger(renderer.height)||renderer.height<=0)return 'descriptor_fields'; + if(Object.prototype.hasOwnProperty.call(renderer,'creativeId')&&(typeof renderer.creativeId!=='string'||!renderer.creativeId||new TextEncoder().encode(renderer.creativeId).length>1024))return 'descriptor_fields'; + if(renderer.tagType!=='iframe'&&renderer.tagType!=='script')return 'descriptor_fields'; + if(typeof renderer.creativeUrl!=='string'||new TextEncoder().encode(renderer.creativeUrl).length>4096)return 'descriptor_fields'; + if(typeof renderer.aaxResponse!=='string'||!renderer.aaxResponse||renderer.aaxResponse.length>349528)return 'descriptor_fields'; try{ var url=new URL(renderer.creativeUrl); - if(url.protocol!=='https:'||url.username||url.password)return false; + if(url.protocol!=='https:'||url.username||url.password)return 'descriptor_envelope'; var binary=atob(renderer.aaxResponse); - if(binary.length>262144||btoa(binary)!==renderer.aaxResponse)return false; + if(binary.length>262144||btoa(binary)!==renderer.aaxResponse)return 'descriptor_envelope'; var bytes=Uint8Array.from(binary,function(character){return character.charCodeAt(0);}); var decoded=JSON.parse(new TextDecoder('utf-8',{fatal:true}).decode(bytes)); - if(!keys(decoded,['seatbid'])||!Array.isArray(decoded.seatbid)||decoded.seatbid.length!==1)return false; + if(!keys(decoded,['seatbid'])||!Array.isArray(decoded.seatbid)||decoded.seatbid.length!==1)return 'descriptor_envelope'; var seat=decoded.seatbid[0]; - if(!keys(seat,['bid'])||!Array.isArray(seat.bid)||seat.bid.length!==1)return false; + if(!keys(seat,['bid'])||!Array.isArray(seat.bid)||seat.bid.length!==1)return 'descriptor_envelope'; var bid=seat.bid[0]; - if(!keys(bid,['ext','h','id','price','w'])||!keys(bid.ext,['creativeurl','tagtype']))return false; - return bid.id===renderer.bidId&&bid.w===renderer.width&&bid.h===renderer.height&& + if(!keys(bid,['ext','h','id','price','w'])||!keys(bid.ext,['creativeurl','tagtype']))return 'descriptor_envelope'; + if(bid.id===renderer.bidId&&bid.w===renderer.width&&bid.h===renderer.height&& bid.ext.creativeurl===renderer.creativeUrl&&bid.ext.tagtype===renderer.tagType&& - typeof bid.price==='number'&&Number.isFinite(bid.price)&&bid.price>=0; - }catch(_error){return false;} + typeof bid.price==='number'&&Number.isFinite(bid.price)&&bid.price>=0)return undefined; + return 'descriptor_envelope'; + }catch(_error){return 'descriptor_envelope';} } function receive(event){ - if(event.source!==parent)return; var message=event.data; - if(!keys(message,['nonce','renderer'])||message.nonce!==expected||!validRenderer(message.renderer))return; + // Stay silent for traffic that is not shaped like the render handshake, so an + // unrelated sender cannot consume this frame's single report. + if(!keys(message,['nonce','renderer']))return; + if(event.source!==parent){report('source_mismatch');return;} + if(message.nonce!==expected){report('nonce_mismatch');return;} + var problem=rendererProblem(message.renderer); + if(problem){report(problem,message.nonce);return;} removeEventListener('message',receive); var acceptedNonce=expected; expected=''; @@ -128,7 +140,7 @@ function receive(event){ var script=document.createElement('script'); script.src='https://client.aps.amazon-adsystem.com/prebid-creative.js'; script.onload=function(){parent.postMessage({message:'trusted-server/aps/renderer-ready',nonce:acceptedNonce},'*');}; - script.onerror=function(){parent.postMessage({message:'trusted-server/aps/renderer-failed',nonce:acceptedNonce},'*');}; + script.onerror=function(){report('amazon_script_error',acceptedNonce);}; document.head.appendChild(script); } addEventListener('message',receive); @@ -3308,4 +3320,41 @@ mod tests { assert!(APS_RENDERER_CSP.contains("sandbox allow-forms")); assert!(!APS_RENDERER_CSP.contains("allow-same-origin")); } + + #[test] + fn renderer_document_reports_a_reason_for_every_silent_guard() { + for reason in [ + "bad_hash", + "source_mismatch", + "nonce_mismatch", + "descriptor_keys", + "descriptor_fields", + "descriptor_envelope", + "amazon_script_error", + ] { + assert!( + APS_RENDERER_DOCUMENT.contains(reason), + "renderer document should report a `{reason}` reason instead of returning silently" + ); + } + + // Reasons travel on the existing failure message rather than a new channel. + assert!( + APS_RENDERER_DOCUMENT.contains("reason:reason"), + "should attach the reason to the failure message" + ); + + // A reason is a fixed category, never a copy of the rejected descriptor. + assert!(!APS_RENDERER_DOCUMENT.contains("JSON.stringify(renderer)")); + assert!(!APS_RENDERER_DOCUMENT.contains("reason:message")); + + // Reporting is one-shot so a hostile sender cannot flood the parent. + assert!( + APS_RENDERER_DOCUMENT.contains("if(reported)return"), + "should report at most one reason per frame" + ); + + // A foreign sender is answered through the parent, never the sender. + assert!(!APS_RENDERER_DOCUMENT.contains("event.source.postMessage")); + } } diff --git a/crates/trusted-server-core/src/integrations/gpt_bootstrap.js b/crates/trusted-server-core/src/integrations/gpt_bootstrap.js index afae77535..0c2f00133 100644 --- a/crates/trusted-server-core/src/integrations/gpt_bootstrap.js +++ b/crates/trusted-server-core/src/integrations/gpt_bootstrap.js @@ -326,7 +326,7 @@ // and deliberately identical to the bundle scheduler — the impression is // spent on a viewed tab, and the post-hydration guarantee holds whenever // the request is actually issued. - ts.scheduleInitialAdInit = function (initialBids, initialSlots) { + ts.scheduleInitialAdInit = function (initialBids, initialSlots, initialAuctionDiagnostics) { // The bundle may replace this scheduler after the fallback claims the initial // pass. Keep the latch on the shared document API so replacement cannot reset it. if ((ts.navGeneration || 0) !== 0 || ts.initialAdInitScheduled) return; @@ -336,6 +336,9 @@ // would overwrite a committed SPA navigation's slots. if (initialSlots !== undefined) ts.adSlots = initialSlots; if (initialBids !== undefined) ts.bids = initialBids; + if (initialAuctionDiagnostics !== undefined) { + ts.auctionDiagnostics = initialAuctionDiagnostics; + } var fire = function () { if ((ts.navGeneration || 0) !== 0) return; if (typeof ts.adInit === "function") ts.adInit(); diff --git a/crates/trusted-server-core/src/integrations/gpt_diagnostics.rs b/crates/trusted-server-core/src/integrations/gpt_diagnostics.rs index 1447a8358..de1aedd36 100644 --- a/crates/trusted-server-core/src/integrations/gpt_diagnostics.rs +++ b/crates/trusted-server-core/src/integrations/gpt_diagnostics.rs @@ -62,6 +62,7 @@ pub enum GptDiagnosticsCookieAction { #[derive(Clone, Debug, Default, PartialEq, Eq)] pub struct GptDiagnosticsRequestDecision { active: bool, + browser_session_active: bool, clean_browser_path_and_query: Option, cookie_action: GptDiagnosticsCookieAction, } @@ -73,6 +74,16 @@ impl GptDiagnosticsRequestDecision { self.active } + /// Whether this request came from an activated diagnostics browser session. + /// + /// Unlike [`Self::active`], this remains true for non-document requests such + /// as the SPA page-bids fetch. It is captured before the private activation + /// cookie is stripped from the request. + #[must_use] + pub(crate) fn browser_session_active(&self) -> bool { + self.browser_session_active + } + /// Whether the response must be private and non-storeable. #[must_use] pub fn requires_private_no_store(&self) -> bool { @@ -121,6 +132,7 @@ impl GptDiagnosticsRequestDecision { pub(crate) fn active_for_tests() -> Self { Self { active: true, + browser_session_active: true, clean_browser_path_and_query: None, cookie_action: GptDiagnosticsCookieAction::None, } @@ -143,6 +155,7 @@ mod head_seam_invariant_tests { ] { out.push(GptDiagnosticsRequestDecision { active, + browser_session_active: active, clean_browser_path_and_query: clean.clone(), cookie_action, }); @@ -279,12 +292,19 @@ pub fn prepare_request( replace_path_and_query(request, &clean_path)?; } - let mut decision = GptDiagnosticsRequestDecision::default(); + let mut decision = GptDiagnosticsRequestDecision { + browser_session_active: integration_enabled + && directive == QueryDirective::Absent + && cookie_state.occurrences == 1 + && cookie_state.canonical, + ..GptDiagnosticsRequestDecision::default() + }; if integration_enabled && eligible_navigation && had_reserved_query { decision.clean_browser_path_and_query = Some(clean_path); match directive { QueryDirective::Enable => { decision.active = true; + decision.browser_session_active = true; decision.cookie_action = GptDiagnosticsCookieAction::SetSession; } QueryDirective::Disable => { @@ -547,6 +567,22 @@ mod tests { assert_eq!(duplicate.headers()[header::COOKIE], "other=value"); } + #[test] + fn active_cookie_marks_non_document_requests_without_activating_document_behavior() { + let mut request = Request::builder() + .method(Method::GET) + .uri("https://publisher.example/_ts/page-bids?path=/article") + .header(header::COOKIE, "__Host-ts-console=1; other=value") + .body(EdgeBody::empty()) + .expect("should build page-bids request"); + + let decision = prepare_request(&settings(true), &mut request).expect("should prepare"); + + assert!(!decision.active()); + assert!(decision.browser_session_active()); + assert_eq!(request.headers()[header::COOKIE], "other=value"); + } + #[test] fn invalid_duplicate_and_disable_directives_fail_closed() { for query in [ diff --git a/crates/trusted-server-core/src/integrations/prebid.rs b/crates/trusted-server-core/src/integrations/prebid.rs index a031c2554..4dc7d4024 100644 --- a/crates/trusted-server-core/src/integrations/prebid.rs +++ b/crates/trusted-server-core/src/integrations/prebid.rs @@ -1,5 +1,4 @@ -use std::collections::HashMap; -use std::collections::HashSet; +use std::collections::{HashMap, HashSet}; use std::sync::{Arc, LazyLock}; #[cfg(test)] use std::time::Duration; @@ -526,6 +525,14 @@ pub struct PrebidIntegrationConfig { pub enabled: bool, #[serde(default)] pub account_id: Option, + /// Prebid User ID modules that Trusted Server installs and keeps installed. + /// + /// Each entry is forwarded to Prebid.js verbatim; publisher-configured + /// entries with other names are preserved. Names must be unique, and no two + /// names may resolve to the same Prebid User ID submodule. + #[serde(default)] + #[validate(nested, custom(function = "validate_unique_managed_user_id_names"))] + pub managed_user_ids: Vec, #[serde(default = "default_timeout_ms")] pub timeout_ms: u32, #[serde(default)] @@ -552,14 +559,6 @@ pub struct PrebidIntegrationConfig { #[serde(default, deserialize_with = "crate::settings::vec_from_seq_or_map")] #[validate(custom(function = "validate_excluded_gam_ad_unit_path_suffixes"))] pub excluded_gam_ad_unit_path_suffixes: Vec, - /// Prebid User ID modules that Trusted Server installs and keeps installed. - /// - /// Each entry is forwarded to Prebid.js verbatim; publisher-configured - /// entries with other names are preserved. Names must be unique, and no two - /// names may resolve to the same Prebid User ID submodule. - #[serde(default)] - #[validate(nested, custom(function = "validate_unique_managed_user_id_names"))] - pub managed_user_ids: Vec, /// CLI-only external bundle build inputs; runtime registration ignores these fields. #[serde(default)] pub bundle: PrebidBundleBuildConfig, @@ -570,6 +569,7 @@ impl Default for PrebidIntegrationConfig { Self { enabled: default_enabled(), account_id: None, + managed_user_ids: Vec::new(), timeout_ms: default_timeout_ms(), debug: false, script_patterns: default_script_patterns(), @@ -578,7 +578,6 @@ impl Default for PrebidIntegrationConfig { external_bundle_sri: None, client_side_bidders: Vec::new(), excluded_gam_ad_unit_path_suffixes: Vec::new(), - managed_user_ids: Vec::new(), bundle: PrebidBundleBuildConfig::default(), } } @@ -596,6 +595,7 @@ impl From<&LegacyPrebidServerConfig> for PrebidIntegrationConfig { Self { enabled: config.enabled, account_id: config.account_id.clone(), + managed_user_ids: config.managed_user_ids.clone(), timeout_ms: config.timeout_ms, debug: config.debug, script_patterns: config.script_patterns.clone(), @@ -604,7 +604,6 @@ impl From<&LegacyPrebidServerConfig> for PrebidIntegrationConfig { external_bundle_sri: config.external_bundle_sri.clone(), client_side_bidders: config.client_side_bidders.clone(), excluded_gam_ad_unit_path_suffixes: config.excluded_gam_ad_unit_path_suffixes.clone(), - managed_user_ids: config.managed_user_ids.clone(), bundle: PrebidBundleBuildConfig::default(), } } @@ -4387,7 +4386,6 @@ server_url = "https://prebid.example/openrtb2/auction" "should manage no User ID modules by default" ); } - #[test] fn excluded_gam_ad_unit_path_suffixes_reject_invalid_values() { for (suffix, expected_message) in [ @@ -5322,6 +5320,97 @@ external_bundle_sri = "sha384-AAAA" ); } + #[test] + #[allow(clippy::field_reassign_with_default)] + fn prepared_browser_injection_uses_only_plan_routes_and_browser_timeout_debug() { + let integration = PrebidIntegration::new(base_config()); + let mut browser_config = PrebidIntegrationConfig::default(); + browser_config.account_id = Some("browser-account".to_string()); + browser_config.timeout_ms = 1750; + browser_config.debug = false; + let plan = AuctionPlan::compile(AuctionPlanConfig { + timeout_ms: 2500, + providers: BTreeMap::from([ + ( + ProviderId::from_str("pbs-primary").expect("should parse provider ID"), + ProviderConfig { + protocol: "openrtb-2.6".to_string(), + profile: "prebid-server".to_string(), + endpoint: "https://primary.example.test/openrtb".to_string(), + timeout_ms: Some(3000), + routing: RoutingMode::Explicit, + notifications: NotificationConfig::default(), + profile_config: json!({"debug": true}), + }, + ), + ( + ProviderId::from_str("pbs-secondary").expect("should parse provider ID"), + ProviderConfig { + protocol: "openrtb-2.6".to_string(), + profile: "prebid-server".to_string(), + endpoint: "https://secondary.example.test/openrtb".to_string(), + timeout_ms: Some(4000), + routing: RoutingMode::Explicit, + notifications: NotificationConfig::default(), + profile_config: json!({"debug": true}), + }, + ), + ]), + bidders: BTreeMap::from([ + ( + BidderId::from_str("secondaryRoute").expect("should parse bidder ID"), + BidderRouteConfig { + provider: ProviderId::from_str("pbs-secondary") + .expect("should parse provider ID"), + }, + ), + ( + BidderId::from_str("primaryRoute").expect("should parse bidder ID"), + BidderRouteConfig { + provider: ProviderId::from_str("pbs-primary") + .expect("should parse provider ID"), + }, + ), + ]), + mediator: None, + request_signing: None, + }) + .expect("should compile plan while browser integration is not part of compilation"); + + let inserts = integration.head_inserts_for_plan(&browser_config, &plan); + let script = &inserts[0]; + + assert!(script.contains(r#""timeout":1750,"debug":false"#)); + assert!( + script.contains(r#""serverSideBidders":["primaryRoute","secondaryRoute"]"#), + "should inject deterministic browser route codes: {script}" + ); + assert!(!script.contains("pbs-primary")); + assert!(!script.contains("pbs-secondary")); + assert!(!script.contains("3000")); + assert!(!script.contains("4000")); + + let disabled_plan = plan.clone().with_enabled(false); + let disabled_inserts = integration.head_inserts_for_plan(&browser_config, &disabled_plan); + assert!( + disabled_inserts[0].contains(r#""serverSideBidders":[]"#), + "auction kill switch should suppress browser server-side bidders: {}", + disabled_inserts[0] + ); + } + + #[test] + fn browser_only_config_defaults_are_independent_and_can_be_disabled() { + let config = PrebidIntegrationConfig { + enabled: false, + ..PrebidIntegrationConfig::default() + }; + + assert!(!config.enabled); + assert_eq!(config.timeout_ms, 1000); + assert!(!config.debug); + } + #[test] fn planned_registration_injects_managed_user_ids() { let mut settings = make_settings(); @@ -5456,85 +5545,6 @@ external_bundle_sri = "sha384-AAAA" ); } - #[test] - #[allow(clippy::field_reassign_with_default)] - fn prepared_browser_injection_uses_only_plan_routes_and_browser_timeout_debug() { - let integration = PrebidIntegration::new(base_config()); - let mut browser_config = PrebidIntegrationConfig::default(); - browser_config.account_id = Some("browser-account".to_string()); - browser_config.timeout_ms = 1750; - browser_config.debug = false; - let plan = AuctionPlan::compile(AuctionPlanConfig { - timeout_ms: 2500, - providers: BTreeMap::from([ - ( - ProviderId::from_str("pbs-primary").expect("should parse provider ID"), - ProviderConfig { - protocol: "openrtb-2.6".to_string(), - profile: "prebid-server".to_string(), - endpoint: "https://primary.example.test/openrtb".to_string(), - timeout_ms: Some(3000), - routing: RoutingMode::Explicit, - notifications: NotificationConfig::default(), - profile_config: json!({"debug": true}), - }, - ), - ( - ProviderId::from_str("pbs-secondary").expect("should parse provider ID"), - ProviderConfig { - protocol: "openrtb-2.6".to_string(), - profile: "prebid-server".to_string(), - endpoint: "https://secondary.example.test/openrtb".to_string(), - timeout_ms: Some(4000), - routing: RoutingMode::Explicit, - notifications: NotificationConfig::default(), - profile_config: json!({"debug": true}), - }, - ), - ]), - bidders: BTreeMap::from([ - ( - BidderId::from_str("secondaryRoute").expect("should parse bidder ID"), - BidderRouteConfig { - provider: ProviderId::from_str("pbs-secondary") - .expect("should parse provider ID"), - }, - ), - ( - BidderId::from_str("primaryRoute").expect("should parse bidder ID"), - BidderRouteConfig { - provider: ProviderId::from_str("pbs-primary") - .expect("should parse provider ID"), - }, - ), - ]), - mediator: None, - request_signing: None, - }) - .expect("should compile plan while browser integration is not part of compilation"); - - let inserts = integration.head_inserts_for_plan(&browser_config, &plan); - let script = &inserts[0]; - - assert!(script.contains(r#""timeout":1750,"debug":false"#)); - assert!( - script.contains(r#""serverSideBidders":["primaryRoute","secondaryRoute"]"#), - "should inject deterministic browser route codes: {script}" - ); - assert!(!script.contains("pbs-primary")); - assert!(!script.contains("pbs-secondary")); - assert!(!script.contains("3000")); - assert!(!script.contains("4000")); - - let disabled_plan = plan.clone().with_enabled(false); - let disabled_inserts = integration.head_inserts_for_plan(&browser_config, &disabled_plan); - assert!( - disabled_inserts[0].contains(r#""serverSideBidders":[]"#), - "auction kill switch should suppress browser server-side bidders: {}", - disabled_inserts[0] - ); - } - #[test] fn head_injector_escapes_script_breakout_in_managed_user_ids() { let mut config = base_config(); @@ -5568,18 +5578,6 @@ external_bundle_sri = "sha384-AAAA" ); } - #[test] - fn browser_only_config_defaults_are_independent_and_can_be_disabled() { - let config = PrebidIntegrationConfig { - enabled: false, - ..PrebidIntegrationConfig::default() - }; - - assert!(!config.enabled); - assert_eq!(config.timeout_ms, 1000); - assert!(!config.debug); - } - #[test] fn head_injector_includes_excluded_gam_ad_unit_path_suffixes() { let mut config = base_config(); diff --git a/crates/trusted-server-core/src/integrations/registry.rs b/crates/trusted-server-core/src/integrations/registry.rs index 47c521e61..c137fd855 100644 --- a/crates/trusted-server-core/src/integrations/registry.rs +++ b/crates/trusted-server-core/src/integrations/registry.rs @@ -691,7 +691,17 @@ impl IntegrationRegistrationBuilder { } } -type RouteValue = (Arc, &'static str); +/// Proxy handler, integration id, and the registered route pattern (kept so +/// telemetry can label responses with the integration-defined literal, e.g. +/// `/integrations/prebid/*`, instead of deriving anything from the request +/// path). +type RouteValue = (Arc, &'static str, String); + +/// A test-constructor route entry: method, path pattern, and the proxy with +/// its integration id ([`IntegrationRegistry::from_routes`] fills the +/// pattern into [`RouteValue`] itself). +#[cfg(test)] +type RouteEntry<'a> = (Method, &'a str, (Arc, &'static str)); struct IntegrationRegistryInner { // Method-specific routers for O(log n) lookups @@ -837,7 +847,11 @@ impl IntegrationRegistry { for proxy in registration.proxies { for route in proxy.routes() { - let value = (proxy.clone(), registration.integration_id); + let value = ( + proxy.clone(), + registration.integration_id, + route.path.clone(), + ); // Convert /* wildcard to matchit's {*rest} syntax let matchit_path = if route.path.ends_with("/*") { @@ -934,6 +948,27 @@ impl IntegrationRegistry { self.find_route(method, path).is_some() } + /// The registered route pattern matched by `method` and `path`, if any. + /// + /// Patterns are integration-defined literals (for example + /// `/integrations/prebid/*`), so they are bounded and content-free and + /// safe to store as a telemetry dimension, unlike the request path. + #[must_use] + pub fn matched_route_pattern(&self, method: &Method, path: &str) -> Option<&str> { + self.find_route(method, path).map(|value| value.2.as_str()) + } + + /// Return true when at least one integration request filter is + /// registered. + /// + /// Adapters use this to decide whether to record a request-filter phase + /// timing span, so unconfigured deployments (no request filters) omit + /// that entry from observability output entirely. + #[must_use] + pub fn has_request_filters(&self) -> bool { + !self.inner.request_filters.is_empty() + } + /// Run pre-routing request filters. /// /// Request header mutations are applied immediately so later filters and @@ -1009,7 +1044,7 @@ impl IntegrationRegistry { services, mut req, } = input; - if let Some((proxy, _)) = self.find_route(method, path) { + if let Some((proxy, _, _)) = self.find_route(method, path) { // Organic proxy handler: generate if needed (best effort). // Only generate for document navigations — subresource requests // may lack consent signals such as the Sec-GPC header. @@ -1322,7 +1357,7 @@ impl IntegrationRegistry { /// # Panics /// /// Panics if route registration fails due to duplicate or invalid paths. - pub fn from_routes(routes: Vec<(Method, &str, RouteValue)>) -> Self { + pub fn from_routes(routes: Vec>) -> Self { let mut get_router = Router::new(); let mut post_router = Router::new(); let mut put_router = Router::new(); @@ -1331,7 +1366,8 @@ impl IntegrationRegistry { let mut head_router = Router::new(); let mut options_router = Router::new(); - for (method, path, value) in routes { + for (method, path, (proxy, integration_id)) in routes { + let value: RouteValue = (proxy, integration_id, path.to_owned()); // Convert /* wildcard to matchit's {*rest} syntax let matchit_path = if path.ends_with("/*") { format!( diff --git a/crates/trusted-server-core/src/integrations/sourcepoint.rs b/crates/trusted-server-core/src/integrations/sourcepoint.rs index 44979c296..d22cc4ffe 100644 --- a/crates/trusted-server-core/src/integrations/sourcepoint.rs +++ b/crates/trusted-server-core/src/integrations/sourcepoint.rs @@ -45,10 +45,11 @@ use crate::settings::{IntegrationConfig, Settings}; const SOURCEPOINT_INTEGRATION_ID: &str = "sourcepoint"; const SOURCEPOINT_CDN_HOST: &str = "cdn.privacy-mgmt.com"; const SOURCEPOINT_CDN_PREFIX: &str = "/integrations/sourcepoint/cdn"; +const SOURCEPOINT_SITE_DATA_PATH: &str = "/mms/v2/get_site_data"; -/// Maximum response body size (5 MB) that will be read into memory for -/// JavaScript rewriting. Responses larger than this are passed through -/// unmodified to avoid unbounded memory consumption. +/// Maximum input body size (5 MiB) collected for JavaScript or HTML rewriting. +/// Declared larger bodies pass through without collection. Bodies that exceed +/// this limit during collection return an integration error (502). const MAX_REWRITE_BODY_SIZE: u64 = 5 * 1024 * 1024; /// Sourcepoint cookie names that are safe to round-trip to the upstream CDN. @@ -609,7 +610,10 @@ impl SourcepointIntegration { /// is a conservative preflight — false negatives just mean we skip the /// `Accept-Encoding: identity` optimisation for that request. fn is_likely_javascript_path(path: &str) -> bool { - path.ends_with(".js") || path.ends_with(".mjs") || path.starts_with("/unified/") + path.ends_with(".js") + || path.ends_with(".mjs") + || path.starts_with("/unified/") + || path == SOURCEPOINT_SITE_DATA_PATH } /// Returns `true` when the response `Content-Type` looks like JavaScript. @@ -650,9 +654,22 @@ impl SourcepointIntegration { } } - fn rewrite_javascript_response(&self, response: &mut Response, rewritten: String) { + fn rewrite_javascript_response( + &self, + response: &mut Response, + rewritten: String, + target_path: &str, + forwarded_cookies: bool, + ) { self.finalize_rewritten_body(response, rewritten, "application/javascript; charset=utf-8"); + // Site data is a dynamic API response despite its JavaScript content + // type. Preserve its upstream cache policy and cookie-aware defaults. + if target_path == SOURCEPOINT_SITE_DATA_PATH { + self.apply_cache_headers(response, forwarded_cookies); + return; + } + // Rewritten JavaScript bundles are static, versioned files (hashed chunk // names, `/unified/4.40.1/…` paths), so we apply a fixed public cache // policy regardless of what upstream sent. This intentionally diverges @@ -870,9 +887,15 @@ impl IntegrationProxy for SourcepointIntegration { None, )?; + // Keep the body streaming where supported so the rewrite collector + // enforces its limit before the adapter buffers the entire response. + let mut platform_request = PlatformHttpRequest::new(proxy_req, backend_name); + if services.http_client().supports_streaming_responses() { + platform_request = platform_request.with_stream_response(); + } let mut response = services .http_client() - .send(PlatformHttpRequest::new(proxy_req, backend_name)) + .send(platform_request) .await .change_context(Self::error("Sourcepoint upstream request failed"))? .response; @@ -930,27 +953,21 @@ impl IntegrationProxy for SourcepointIntegration { .and_then(|v| v.to_str().ok()) .and_then(|s| s.parse::().ok()); - match content_length { - Some(len) if len > MAX_REWRITE_BODY_SIZE => { - log::warn!( - "Sourcepoint: response body for {path} exceeds {} bytes \ - (Content-Length: {len}), skipping rewrite (reason: known_length_too_large)", - MAX_REWRITE_BODY_SIZE - ); - self.apply_cache_headers(&mut response, forwarded_cookies); - return Ok(response); - } - None => { - log::warn!( - "Sourcepoint: no Content-Length for {path}, \ - skipping rewrite to avoid unbounded memory read (reason: missing_content_length)" - ); - self.apply_cache_headers(&mut response, forwarded_cookies); - return Ok(response); - } - Some(_) => {} + if let Some(len) = content_length + && len > MAX_REWRITE_BODY_SIZE + { + log::warn!( + "Sourcepoint: response body for {path} exceeds {} bytes \ + (Content-Length: {len}), skipping rewrite (reason: known_length_too_large)", + MAX_REWRITE_BODY_SIZE + ); + self.apply_cache_headers(&mut response, forwarded_cookies); + return Ok(response); } + // Content-Length is optional and advisory. Stop at the actual + // byte limit even when the header is absent or understates the + // size. Overflow discards the partial body and returns a 502. let (resp_parts, resp_body) = response.into_parts(); let body_bytes = collect_response_bounded( resp_body, @@ -975,7 +992,12 @@ impl IntegrationProxy for SourcepointIntegration { }; if response_is_javascript { let rewritten = Self::rewrite_script_content(&body); - self.rewrite_javascript_response(&mut response, rewritten); + self.rewrite_javascript_response( + &mut response, + rewritten, + target_path, + forwarded_cookies, + ); } else { let rewritten = Self::rewrite_html_content(&body); self.rewrite_html_response(&mut response, rewritten, forwarded_cookies); @@ -1085,9 +1107,539 @@ impl IntegrationHeadInjector for SourcepointIntegration { #[cfg(test)] mod tests { use super::*; + use crate::error::IntoHttpResponse as _; use crate::integrations::{IntegrationDocumentState, IntegrationRegistry}; + use crate::platform::test_support::{StubHttpClient, build_services_with_http_client}; + use crate::platform::{ + PlatformError, PlatformHttpClient, PlatformPendingRequest, PlatformResponse, + PlatformSelectResult, + }; use crate::test_support::tests::create_test_settings; use serde_json::json; + use std::sync::atomic::{AtomicUsize, Ordering}; + + const TEST_CHUNK_SIZE: usize = 8192; + + struct StreamingHttpClient { + stub: StubHttpClient, + reads: Arc, + } + + impl StreamingHttpClient { + fn new() -> Self { + Self { + stub: StubHttpClient::new(), + reads: Arc::new(AtomicUsize::new(0)), + } + } + } + + #[async_trait(?Send)] + impl PlatformHttpClient for StreamingHttpClient { + fn supports_streaming_responses(&self) -> bool { + true + } + + async fn send( + &self, + request: PlatformHttpRequest, + ) -> Result> { + let mut response = self.stub.send(request).await?; + let body = std::mem::replace(response.response.body_mut(), EdgeBody::empty()); + let EdgeBody::Once(bytes) = body else { + panic!("should receive a buffered stub body"); + }; + let reads = Arc::clone(&self.reads); + let chunks = futures::stream::unfold((bytes, 0), move |(bytes, offset)| { + reads.fetch_add(1, Ordering::Relaxed); + let end = (offset + TEST_CHUNK_SIZE).min(bytes.len()); + futures::future::ready( + (offset < bytes.len()).then(|| (bytes.slice(offset..end), (bytes, end))), + ) + }); + *response.response.body_mut() = EdgeBody::stream(chunks); + Ok(response) + } + + async fn send_async( + &self, + request: PlatformHttpRequest, + ) -> Result> { + self.stub.send_async(request).await + } + + async fn select( + &self, + pending_requests: Vec, + ) -> Result> { + self.stub.select(pending_requests).await + } + } + + #[test] + fn handle_rewrites_streamed_javascript_without_content_length() { + futures::executor::block_on(async { + let settings = create_test_settings(); + let integration = SourcepointIntegration::new(Arc::new(config(true))); + let client = Arc::new(StreamingHttpClient::new()); + let input = format!(r#"var api="https://{SOURCEPOINT_CDN_HOST}/consent/tcfv2";"#); + client.stub.push_response_with_headers( + 200, + input.into_bytes(), + vec![("content-type", "application/javascript")], + ); + let services = build_services_with_http_client(client.clone()); + + let response = integration + .handle( + &settings, + &services, + make_req( + Method::GET, + "https://publisher.example.com/integrations/sourcepoint/cdn/wrapper.js", + ), + ) + .await + .expect("should proxy JavaScript without Content-Length"); + + assert_eq!( + response + .into_body() + .into_bytes_bounded(1024) + .await + .expect("should collect JavaScript response") + .as_ref(), + br#"var api="/integrations/sourcepoint/cdn/consent/tcfv2";"#, + "should rewrite a streamed CDN URL without Content-Length" + ); + assert_eq!( + client.stub.recorded_stream_response_flags(), + vec![true], + "should request streaming before collecting the upstream body" + ); + }); + } + + #[test] + fn handle_rewrites_streamed_html_without_content_length() { + futures::executor::block_on(async { + let settings = create_test_settings(); + let integration = SourcepointIntegration::new(Arc::new(config(true))); + let client = Arc::new(StreamingHttpClient::new()); + client.stub.push_response_with_headers( + 200, + br#""#.to_vec(), + vec![("content-type", "text/html"), ("cache-control", "no-store")], + ); + let services = build_services_with_http_client(client.clone()); + + let response = integration + .handle( + &settings, + &services, + make_req(Method::GET, "https://publisher.example.com/integrations/sourcepoint/cdn/us_pm/index.html"), + ) + .await + .expect("should proxy HTML without Content-Length"); + + assert_eq!( + get_header_str(&response, header::CACHE_CONTROL), + Some("no-store"), + "should preserve upstream HTML cache policy" + ); + assert_eq!( + response.into_body().into_bytes_bounded(1024).await.expect("should collect HTML response").as_ref(), + br#""#, + "should rewrite streamed privacy-manager assets without Content-Length" + ); + }); + } + + #[test] + fn handle_accepts_exact_rewrite_limit_and_stops_reading_on_overflow() { + futures::executor::block_on(async { + let settings = create_test_settings(); + let integration = SourcepointIntegration::new(Arc::new(config(true))); + let limit = MAX_REWRITE_BODY_SIZE as usize; + for content_type in ["application/javascript", "text/html"] { + for declared_length in [None, Some("1")] { + for extra_bytes in [0, 1, TEST_CHUNK_SIZE * 2] { + let client = Arc::new(StreamingHttpClient::new()); + let mut headers = vec![("content-type", content_type)]; + if let Some(length) = declared_length { + headers.push(("content-length", length)); + } + client.stub.push_response_with_headers( + 200, + vec![b' '; limit + extra_bytes], + headers, + ); + let services = build_services_with_http_client(client.clone()); + + let result = integration.handle( + &settings, + &services, + make_req(Method::GET, "https://publisher.example.com/integrations/sourcepoint/cdn/asset"), + ).await; + + if extra_bytes == 0 { + let response = result.expect("should accept exactly 5 MiB"); + assert!( + response.headers().get(header::CONTENT_LENGTH).is_none(), + "should remove the advisory length after rewriting" + ); + assert_eq!( + take_body_bytes(response).len(), + limit, + "should retain the entire body at the limit" + ); + } else { + let error = + result.expect_err("should reject an oversized streamed body"); + assert_eq!( + error.current_context().status_code(), + StatusCode::BAD_GATEWAY, + "should report upstream overflow as 502" + ); + assert!( + matches!(error.current_context(), TrustedServerError::Integration { integration, message } if integration == SOURCEPOINT_INTEGRATION_ID && message.contains("exceeds")), + "should identify Sourcepoint response overflow" + ); + } + assert_eq!( + client.reads.load(Ordering::Relaxed), + limit / TEST_CHUNK_SIZE + 1, + "should stop at EOF or the first overflowing chunk without draining the stream" + ); + } + } + } + }); + } + + #[test] + fn handle_passes_through_declared_oversize_without_reading() { + futures::executor::block_on(async { + let settings = create_test_settings(); + let integration = SourcepointIntegration::new(Arc::new(config(true))); + for content_type in ["application/javascript", "text/html"] { + let client = Arc::new(StreamingHttpClient::new()); + let length = (MAX_REWRITE_BODY_SIZE + 1).to_string(); + client.stub.push_response_with_headers( + 200, + vec![b' '; MAX_REWRITE_BODY_SIZE as usize + 1], + vec![ + ("content-type", content_type), + ("content-length", &length), + ("cache-control", "no-store"), + ], + ); + let services = build_services_with_http_client(client.clone()); + + let response = integration + .handle( + &settings, + &services, + make_req( + Method::GET, + "https://publisher.example.com/integrations/sourcepoint/cdn/asset", + ), + ) + .await + .expect("should pass through a declared oversized response"); + + assert!( + matches!(response.body(), EdgeBody::Stream(_)), + "should retain the original stream" + ); + assert_eq!( + client.reads.load(Ordering::Relaxed), + 0, + "should not poll a declared oversized body" + ); + assert_eq!( + get_header_str(&response, header::CONTENT_LENGTH), + Some(length.as_str()), + "should preserve the pass-through length" + ); + assert_eq!( + get_header_str(&response, header::CACHE_CONTROL), + Some("no-store"), + "should preserve upstream cache policy" + ); + } + }); + } + + #[test] + fn handle_keeps_ineligible_responses_streaming() { + futures::executor::block_on(async { + let settings = create_test_settings(); + for (content_type, method, status, rewrite_sdk) in [ + ("application/json", Method::GET, 200, true), + ("text/css", Method::GET, 200, true), + ("image/png", Method::GET, 200, true), + ("application/javascript", Method::GET, 200, false), + ("text/html", Method::GET, 200, false), + ("application/javascript", Method::HEAD, 200, true), + ("text/html", Method::POST, 200, true), + ("application/javascript", Method::GET, 206, true), + ("text/html", Method::GET, 404, true), + ] { + let mut cfg = config(true); + cfg.rewrite_sdk = rewrite_sdk; + let integration = SourcepointIntegration::new(Arc::new(cfg)); + let client = Arc::new(StreamingHttpClient::new()); + client.stub.push_response_with_headers( + status, + b"unchanged".to_vec(), + vec![ + ("content-type", content_type), + ("content-encoding", "gzip"), + ("cache-control", "no-store"), + ], + ); + let services = build_services_with_http_client(client.clone()); + let expected_body: &[u8] = if method == Method::HEAD { + b"" + } else { + b"unchanged" + }; + let mut request = make_req( + method, + "https://publisher.example.com/integrations/sourcepoint/cdn/asset", + ); + set_req_header(&mut request, header::ACCEPT_ENCODING, "gzip, br"); + + let response = integration + .handle(&settings, &services, request) + .await + .expect("should pass through an ineligible response"); + + assert!( + matches!(response.body(), EdgeBody::Stream(_)), + "should leave ineligible bodies streaming" + ); + assert_eq!( + client.reads.load(Ordering::Relaxed), + 0, + "should not poll an ineligible body" + ); + assert_eq!( + get_header_str(&response, header::CONTENT_ENCODING), + Some("gzip"), + "should preserve pass-through encoding" + ); + assert_eq!( + get_header_str(&response, header::CACHE_CONTROL), + Some("no-store"), + "should preserve pass-through cache policy" + ); + assert_eq!( + response + .into_body() + .into_bytes_bounded(1024) + .await + .expect("should collect pass-through body") + .as_ref(), + expected_body, + "should preserve pass-through bytes or omit them for HEAD" + ); + assert!( + client.stub.recorded_request_headers()[0] + .iter() + .any(|(name, value)| name == "accept-encoding" && value == "gzip, br"), + "should forward the client's encoding for non-script paths" + ); + } + }); + } + + #[test] + fn handle_rewrites_on_buffered_adapters_with_or_without_content_length() { + futures::executor::block_on(async { + let settings = create_test_settings(); + let integration = SourcepointIntegration::new(Arc::new(config(true))); + for has_length in [false, true] { + let client = Arc::new(StubHttpClient::new()); + let input = format!(r#"var api="https://{SOURCEPOINT_CDN_HOST}/consent/tcfv2";"#); + let length = input.len().to_string(); + let mut headers = vec![ + ("content-type", "application/javascript"), + ("content-encoding", "identity"), + ("vary", "Accept-Encoding, Origin"), + ]; + if has_length { + headers.push(("content-length", &length)); + } + client.push_response_with_headers(200, input.into_bytes(), headers); + let services = build_services_with_http_client(client.clone()); + + let response = integration + .handle( + &settings, + &services, + make_req( + Method::GET, + "https://publisher.example.com/integrations/sourcepoint/cdn/wrapper.js", + ), + ) + .await + .expect("should rewrite a buffered response"); + + assert_eq!( + client.recorded_stream_response_flags(), + vec![false], + "should respect adapters without streaming support" + ); + assert!( + response.headers().get(header::CONTENT_LENGTH).is_none(), + "should not forward a stale upstream length" + ); + assert!( + response.headers().get(header::CONTENT_ENCODING).is_none(), + "should remove upstream encoding after rewriting" + ); + assert_eq!( + get_header_str(&response, header::VARY), + Some("Origin"), + "should remove only Accept-Encoding from Vary" + ); + assert_eq!( + get_header_str(&response, header::CACHE_CONTROL), + Some("public, max-age=3600"), + "should retain the static JavaScript cache policy" + ); + assert_eq!( + get_header_str(&response, header::CONTENT_TYPE), + Some("application/javascript; charset=utf-8"), + "should identify rewritten JavaScript" + ); + assert_eq!( + take_body_bytes(response), + br#"var api="/integrations/sourcepoint/cdn/consent/tcfv2";"#, + "should rewrite with either header state" + ); + } + }); + } + + #[test] + fn handle_site_data_requests_identity_and_preserves_dynamic_cache_policy() { + futures::executor::block_on(async { + let settings = create_test_settings(); + let integration = SourcepointIntegration::new(Arc::new(config(true))); + for (upstream_cache, forwarded_cookies, sets_cookie, expected_cache) in [ + (Some("no-store"), false, false, "no-store"), + ( + Some("private, max-age=60"), + true, + false, + "private, max-age=60", + ), + (None, true, false, "private, max-age=0"), + (None, false, false, "public, max-age=3600"), + ( + Some("public, max-age=3600"), + false, + true, + "private, no-store", + ), + ] { + let client = Arc::new(StreamingHttpClient::new()); + let input = format!(r#"var api="https://{SOURCEPOINT_CDN_HOST}/consent/tcfv2";"#); + let mut headers = vec![("content-type", "application/javascript")]; + if let Some(cache) = upstream_cache { + headers.push(("cache-control", cache)); + } + if sets_cookie { + headers.push(("set-cookie", "consentUUID=example; Path=/")); + } + client + .stub + .push_response_with_headers(200, input.into_bytes(), headers); + let services = build_services_with_http_client(client.clone()); + let mut request = make_req( + Method::GET, + "https://publisher.example.com/integrations/sourcepoint/cdn/mms/v2/get_site_data?account_id=123", + ); + set_req_header(&mut request, header::ACCEPT_ENCODING, "gzip, br"); + if forwarded_cookies { + set_req_header(&mut request, header::COOKIE, "consentUUID=example"); + } + + let response = integration + .handle(&settings, &services, request) + .await + .expect("should rewrite dynamic site data"); + + assert!( + client.stub.recorded_request_headers()[0] + .iter() + .any(|(name, value)| name == "accept-encoding" && value == "identity"), + "should request uncompressed site data despite its extensionless path" + ); + assert_eq!( + get_header_str(&response, header::CACHE_CONTROL), + Some(expected_cache), + "should use the dynamic endpoint's cache policy" + ); + assert_eq!( + take_body_bytes(response), + br#"var api="/integrations/sourcepoint/cdn/consent/tcfv2";"#, + "should rewrite unknown-length site data" + ); + } + }); + } + + #[test] + fn handle_preserves_invalid_utf8_bytes_and_headers() { + futures::executor::block_on(async { + let settings = create_test_settings(); + let integration = SourcepointIntegration::new(Arc::new(config(true))); + let client = Arc::new(StreamingHttpClient::new()); + let bytes = vec![0x1f, 0x8b, 0xff]; + client.stub.push_response_with_headers( + 200, + bytes.clone(), + vec![ + ("content-type", "application/javascript"), + ("content-encoding", "gzip"), + ("cache-control", "no-store"), + ], + ); + let services = build_services_with_http_client(client); + + let response = integration + .handle( + &settings, + &services, + make_req( + Method::GET, + "https://publisher.example.com/integrations/sourcepoint/cdn/wrapper.js", + ), + ) + .await + .expect("should retain non-UTF-8 content unchanged"); + + assert_eq!( + get_header_str(&response, header::CONTENT_ENCODING), + Some("gzip"), + "should preserve encoding when no rewrite occurs" + ); + assert_eq!( + get_header_str(&response, header::CACHE_CONTROL), + Some("no-store"), + "should preserve cache policy when no rewrite occurs" + ); + assert_eq!( + take_body_bytes(response), + bytes, + "should preserve invalid UTF-8 bytes" + ); + }); + } fn config(enabled: bool) -> SourcepointConfig { SourcepointConfig { @@ -1444,9 +1996,12 @@ mod tests { assert!(SourcepointIntegration::is_likely_javascript_path( "/module/sourcepoint.mjs" )); - assert!(!SourcepointIntegration::is_likely_javascript_path( + assert!(SourcepointIntegration::is_likely_javascript_path( "/mms/v2/get_site_data" )); + assert!(!SourcepointIntegration::is_likely_javascript_path( + "/mms/v2/get_site_data/other" + )); assert!(!SourcepointIntegration::is_likely_javascript_path( "/consent/tcfv2" )); @@ -1950,7 +2505,12 @@ mod tests { set_header(&mut response, header::CACHE_CONTROL, "no-store"); *response.body_mut() = EdgeBody::from(b"payload".to_vec()); - integration.rewrite_javascript_response(&mut response, "rewritten".to_string()); + integration.rewrite_javascript_response( + &mut response, + "rewritten".to_string(), + "/wrapper.js", + false, + ); assert_eq!(response.status(), StatusCode::OK); assert_eq!( @@ -1989,7 +2549,12 @@ mod tests { set_header(&mut response, header::CACHE_CONTROL, "public, max-age=3600"); *response.body_mut() = EdgeBody::from(b"payload".to_vec()); - integration.rewrite_javascript_response(&mut response, "rewritten".to_string()); + integration.rewrite_javascript_response( + &mut response, + "rewritten".to_string(), + "/wrapper.js", + false, + ); assert_eq!( get_header_str(&response, header::CACHE_CONTROL), @@ -2009,7 +2574,12 @@ mod tests { set_header(&mut response, header::VARY, "Accept-Encoding"); *response.body_mut() = EdgeBody::from(b"payload".to_vec()); - integration.rewrite_javascript_response(&mut response, "rewritten".to_string()); + integration.rewrite_javascript_response( + &mut response, + "rewritten".to_string(), + "/wrapper.js", + false, + ); assert!( response.headers().get(header::VARY).is_none(), diff --git a/crates/trusted-server-core/src/lib.rs b/crates/trusted-server-core/src/lib.rs index 249f741ab..3083a95eb 100644 --- a/crates/trusted-server-core/src/lib.rs +++ b/crates/trusted-server-core/src/lib.rs @@ -31,6 +31,7 @@ ) )] +pub mod access_telemetry; pub(crate) mod asset_image_optimizer; pub mod auction; pub mod auction_config_types; @@ -62,6 +63,7 @@ pub mod proxy; pub mod publisher; pub mod redacted; pub mod request_signing; +pub mod request_timing; pub mod response_privacy; pub mod rsc_flight; pub(crate) mod s3_sigv4; diff --git a/crates/trusted-server-core/src/platform/mod.rs b/crates/trusted-server-core/src/platform/mod.rs index 63f5c2a94..2802bd3e5 100644 --- a/crates/trusted-server-core/src/platform/mod.rs +++ b/crates/trusted-server-core/src/platform/mod.rs @@ -43,6 +43,7 @@ mod template_assembly; mod template_cache; #[cfg(test)] pub(crate) mod test_support; +mod timed_kv; mod traits; mod types; @@ -72,6 +73,7 @@ pub use template_cache::{ TemplateCookieValue, TemplateEntry, TemplateMetadata, TemplateMetadataEncodeError, UnavailableTemplateCache, VaryHeaderValues, VarySpec, reader_url_surrogate_key, }; +pub use timed_kv::TimedKvStore; pub use traits::{PlatformBackend, PlatformConfigStore, PlatformGeo, PlatformSecretStore}; pub use types::{ ClientInfo, GeoInfo, PlatformBackendSpec, RuntimeServices, RuntimeServicesBuilder, StoreId, diff --git a/crates/trusted-server-core/src/platform/test_support.rs b/crates/trusted-server-core/src/platform/test_support.rs index 17813c77d..aeab49ff3 100644 --- a/crates/trusted-server-core/src/platform/test_support.rs +++ b/crates/trusted-server-core/src/platform/test_support.rs @@ -594,6 +594,7 @@ impl PlatformHttpClient for StubHttpClient { .pop_front() .ok_or_else(|| Report::new(PlatformError::HttpClient))?; + let stream_response = stream_response || response.stream_body; let edge_response = build_stub_pending_response( StubPendingResponse { backend_name: request.backend_name, @@ -601,7 +602,7 @@ impl PlatformHttpClient for StubHttpClient { body: response.body, headers: response.headers, }, - stream_response || response.stream_body, + stream_response, request_is_head, )?; diff --git a/crates/trusted-server-core/src/platform/timed_kv.rs b/crates/trusted-server-core/src/platform/timed_kv.rs new file mode 100644 index 000000000..5697f5e68 --- /dev/null +++ b/crates/trusted-server-core/src/platform/timed_kv.rs @@ -0,0 +1,253 @@ +//! Latency-only timing decorator for KV store handles. +//! +//! [`TimedKvStore`] wraps an inner store plus a [`RequestTimings`] handle and +//! records [`Phase::EcKv`] around every call. It implements both +//! [`PlatformKvStore`] (for consent-store access obtained through +//! [`RuntimeServices`](super::RuntimeServices)) and [`EcKvStore`] (for +//! [`KvIdentityGraph`](crate::ec::kv::KvIdentityGraph) construction sites), +//! because no single existing abstraction covers the whole `ts-kv` taxonomy: +//! EC graph operations go through [`EcKvStore`] while consent persistence +//! uses [`PlatformKvStore`] directly. +//! +//! The decorator measures store-call latency only: it never reads, parses, +//! or logs any value passing through it. + +use std::sync::Arc; +use std::time::Duration; + +use async_trait::async_trait; +use bytes::Bytes; +use edgezero_core::key_value_store::{KvError, KvPage, KvStore as PlatformKvStore}; +use error_stack::Report; + +use crate::ec::kv_backend::{EcKvLookup, EcKvStore, EcKvWrite, EcKvWriteOutcome}; +use crate::error::TrustedServerError; +use crate::request_timing::{Phase, RequestTimings}; + +/// Wraps `inner` plus a [`RequestTimings`] handle, recording [`Phase::EcKv`] +/// around every store call made through it. +pub struct TimedKvStore { + /// The wrapped store handle. + inner: S, + /// The request's phase-timing collector. + timings: RequestTimings, +} + +impl TimedKvStore { + /// Creates a decorator around `inner` that records into `timings`. + #[must_use] + pub fn new(inner: S, timings: RequestTimings) -> Self { + Self { inner, timings } + } +} + +#[async_trait(?Send)] +impl PlatformKvStore for TimedKvStore> { + async fn get_bytes(&self, key: &str) -> Result, KvError> { + let _span = self.timings.span(Phase::EcKv); + self.inner.get_bytes(key).await + } + + async fn put_bytes(&self, key: &str, value: Bytes) -> Result<(), KvError> { + let _span = self.timings.span(Phase::EcKv); + self.inner.put_bytes(key, value).await + } + + // Forwarded explicitly: the trait's default body falls back to + // `get_bytes`, which would silently downgrade a backend's cheap + // metadata-only existence probe (the Spin adapter has one) into a full + // value transfer just because the store was decorated. + async fn exists(&self, key: &str) -> Result { + let _span = self.timings.span(Phase::EcKv); + self.inner.exists(key).await + } + + async fn put_bytes_with_ttl( + &self, + key: &str, + value: Bytes, + ttl: Duration, + ) -> Result<(), KvError> { + let _span = self.timings.span(Phase::EcKv); + self.inner.put_bytes_with_ttl(key, value, ttl).await + } + + async fn delete(&self, key: &str) -> Result<(), KvError> { + let _span = self.timings.span(Phase::EcKv); + self.inner.delete(key).await + } + + async fn list_keys_page( + &self, + prefix: &str, + cursor: Option<&str>, + limit: usize, + ) -> Result { + let _span = self.timings.span(Phase::EcKv); + self.inner.list_keys_page(prefix, cursor, limit).await + } +} + +impl EcKvStore for TimedKvStore { + fn store_name(&self) -> &str { + self.inner.store_name() + } + + fn lookup(&self, key: &str) -> Result, Report> { + let _span = self.timings.span(Phase::EcKv); + self.inner.lookup(key) + } + + fn key_exists(&self, key: &str) -> Result> { + let _span = self.timings.span(Phase::EcKv); + self.inner.key_exists(key) + } + + fn insert( + &self, + key: &str, + write: EcKvWrite<'_>, + ) -> Result> { + let _span = self.timings.span(Phase::EcKv); + self.inner.insert(key, write) + } + + fn list_keys_with_prefix( + &self, + prefix: &str, + limit: u32, + ) -> Result, Report> { + let _span = self.timings.span(Phase::EcKv); + self.inner.list_keys_with_prefix(prefix, limit) + } + + fn delete(&self, key: &str) -> Result<(), Report> { + let _span = self.timings.span(Phase::EcKv); + self.inner.delete(key) + } +} + +#[cfg(test)] +mod tests { + use std::time::Duration as StdDuration; + + use super::*; + use crate::ec::kv_backend::test_support::InMemoryEcKv; + + #[test] + fn ec_kv_store_operations_accumulate_into_ec_kv_phase() { + let timings = RequestTimings::new(); + let store = TimedKvStore::new(InMemoryEcKv::new("test-store"), timings.clone()); + + store + .insert( + "key-a", + EcKvWrite { + body: "{}", + metadata: "{}", + ttl: StdDuration::from_secs(60), + mode: crate::ec::kv_backend::EcKvWriteMode::Add, + }, + ) + .expect("should insert into the in-memory store"); + store.lookup("key-a").expect("should read back the entry"); + + timings.mark_headers_ready(); + assert!( + timings.snapshot().kv_ms.is_some(), + "should record Phase::EcKv across both store calls" + ); + } + + #[test] + fn exists_delegates_to_the_inner_store_not_get_bytes() { + // A stub whose `exists` answer contradicts its `get_bytes` answer: + // if the decorator fell back to the trait's get-and-discard default + // body, this would return `false`. + struct ExistsOnlyStore; + + #[async_trait::async_trait(?Send)] + impl PlatformKvStore for ExistsOnlyStore { + async fn get_bytes(&self, _key: &str) -> Result, KvError> { + Ok(None) + } + async fn put_bytes(&self, _key: &str, _value: Bytes) -> Result<(), KvError> { + Ok(()) + } + async fn put_bytes_with_ttl( + &self, + _key: &str, + _value: Bytes, + _ttl: StdDuration, + ) -> Result<(), KvError> { + Ok(()) + } + async fn delete(&self, _key: &str) -> Result<(), KvError> { + Ok(()) + } + async fn list_keys_page( + &self, + _prefix: &str, + _cursor: Option<&str>, + _limit: usize, + ) -> Result { + Ok(KvPage { + keys: Vec::new(), + cursor: None, + }) + } + async fn exists(&self, _key: &str) -> Result { + Ok(true) + } + } + + let timings = RequestTimings::new(); + let inner: Arc = Arc::new(ExistsOnlyStore); + let store = TimedKvStore::new(inner, timings.clone()); + + let exists = futures::executor::block_on(store.exists("key")) + .expect("should forward the existence probe"); + assert!( + exists, + "should delegate to the inner exists, not the get_bytes default body" + ); + timings.mark_headers_ready(); + assert!( + timings.snapshot().kv_ms.is_some(), + "should time the existence probe like any other store operation" + ); + } + + #[test] + fn store_name_is_not_timed() { + let timings = RequestTimings::new(); + let store = TimedKvStore::new(InMemoryEcKv::new("test-store"), timings.clone()); + + assert_eq!(store.store_name(), "test-store"); + timings.mark_headers_ready(); + assert!( + timings.snapshot().kv_ms.is_none(), + "store_name is a metadata accessor, not a store operation" + ); + } + + #[test] + fn platform_kv_store_operations_accumulate_into_ec_kv_phase() { + let timings = RequestTimings::new(); + let inner: Arc = Arc::new(crate::platform::UnavailableKvStore); + let store = TimedKvStore::new(inner, timings.clone()); + + // UnavailableKvStore errors on every call; the decorator still times + // the attempt regardless of outcome. + futures::executor::block_on(async { + let _ = store.get_bytes("key").await; + let _ = store.put_bytes("key", Bytes::from_static(b"value")).await; + }); + + timings.mark_headers_ready(); + assert!( + timings.snapshot().kv_ms.is_some(), + "should record Phase::EcKv even when the inner store errors" + ); + } +} diff --git a/crates/trusted-server-core/src/publisher.rs b/crates/trusted-server-core/src/publisher.rs index 051eb8048..b93bf13d7 100644 --- a/crates/trusted-server-core/src/publisher.rs +++ b/crates/trusted-server-core/src/publisher.rs @@ -75,6 +75,7 @@ use crate::platform::{ reader_url_surrogate_key, }; use crate::price_bucket::{PriceGranularity, price_bucket}; +use crate::request_timing::{AuctionWaitPlacement, Phase, RequestTimings}; use crate::response_privacy::{ apply_inactive_ad_stack_browser_cache_policy, cache_control_forbids_shared_storage, enforce_synthesized_html_cache_privacy, enforce_terminal_private_cache_privacy, @@ -95,21 +96,44 @@ const DEFAULT_PUBLISHER_FIRST_BYTE_TIMEOUT: Duration = Duration::from_secs(15); const HEADER_X_TS_TEMPLATE_CACHE: &str = "x-ts-template-cache"; const HEADER_X_TS_ASSEMBLY: &str = "x-ts-assembly"; -#[derive(Clone, Copy, PartialEq, Eq)] -enum TemplateCacheResponseState { +/// Outcome of a template-cache lookup/store attempt for one response. +/// +/// Set on every response that passes through the assembly pipeline via +/// [`set_template_cache_response_state`], which writes both the +/// `x-ts-template-cache` response header and this same value as a typed +/// response extension, so the two can never drift. Access telemetry reads +/// the extension rather than the header, since operator-configured response +/// headers can override a managed header but cannot touch extensions. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum TemplateCacheResponseState { + /// The cached template was found and reused. Hit, + /// No cached template existed; the cache store is reserved for this + /// content type. MissReserved, + /// No cached template existed; one was stored after assembly. MissStored, + /// No cached template existed; storing the freshly assembled template + /// failed. MissStoreError, + /// The request bypassed the cache lookup. BypassRequest, + /// The response bypassed the cache store. BypassResponse, + /// The response's content type is not supported by the template cache. Unsupported, + /// The cached template entry was invalid and could not be reused. Invalid, + /// A backend error prevented the cache lookup or store. BackendError, } impl TemplateCacheResponseState { - const fn as_str(self) -> &'static str { + /// Renders this variant as the string written to the + /// `x-ts-template-cache` header and the `template_cache_state` access + /// telemetry column. + #[must_use] + pub const fn as_str(self) -> &'static str { match self { Self::Hit => "hit", Self::MissReserved => "miss-reserved", @@ -132,6 +156,7 @@ fn set_template_cache_response_state( HEADER_X_TS_TEMPLATE_CACHE, HeaderValue::from_static(state.as_str()), ); + response.extensions_mut().insert(state); } #[derive(Clone, Copy, PartialEq, Eq)] @@ -613,7 +638,8 @@ fn parse_single_module_filename(filename: &str) -> Option<&'static str> { .and_then(|s| s.strip_suffix(".min.js").or_else(|| s.strip_suffix(".js")))?; trusted_server_js::all_module_ids() - .into_iter() + .iter() + .copied() .find(|&id| id == stem) } @@ -1725,6 +1751,12 @@ pub struct OwnedProcessResponseParams { /// rescanned from the output, which cannot tell a `nonce` attribute from the same /// word inside a script. pub(crate) csp_nonce_observed: Option>, + /// Per-request phase-timing handle, carried into the streaming/buffered + /// finalizers so the `` seam wait can be recorded with the right + /// [`AuctionWaitPlacement`]. Cheap to clone (an `Arc` handle); a request that + /// never attached one to its extensions gets a fresh, unattached collector + /// that nothing ever renders. + pub(crate) timings: RequestTimings, } /// Response-authorized template cache insert inputs. The key is built before origin lookup; the @@ -1968,6 +2000,8 @@ pub async fn buffer_publisher_response_async( ¶ms.request_scheme, ¶ms.request_host, ), + timings: params.timings.clone(), + placement: AuctionWaitPlacement::PreHeader, }, ) .await; @@ -2076,7 +2110,7 @@ fn template_fingerprint(settings: &Settings) -> String { let mut hasher = sha2::Sha256::new(); hasher.update( - trusted_server_js::concatenated_hash(&trusted_server_js::all_module_ids()).as_bytes(), + trusted_server_js::concatenated_hash(trusted_server_js::all_module_ids()).as_bytes(), ); // EdgeZero's canonical form sorts object keys itself, so independently deserialized // HashMaps hash identically even when a dependency enables `serde_json/preserve_order`. @@ -2129,6 +2163,7 @@ fn build_template_assembly_params( request_scheme: &str, price_granularity: PriceGranularity, ad_bids_state: AdBidsState, + timings: RequestTimings, ) -> OwnedProcessResponseParams { OwnedProcessResponseParams { csp_nonce_observed: None, @@ -2151,6 +2186,7 @@ fn build_template_assembly_params( price_granularity, gpt_diagnostics: None, suppress_datadome_client_side_tag: false, + timings, } } @@ -2493,6 +2529,8 @@ pub async fn publisher_response_into_streaming_response( ¶ms.request_scheme, ¶ms.request_host, ), + timings: params.timings.clone(), + placement: AuctionWaitPlacement::InStream, }, ) .await; @@ -2611,6 +2649,7 @@ pub async fn publisher_response_into_streaming_response( &orchestrator, &services, &settings, + AuctionWaitPlacement::InStream, ) .await; // Collection reached a terminal result; disarm only now @@ -2636,6 +2675,8 @@ pub async fn publisher_response_into_streaming_response( ¶ms.request_scheme, ¶ms.request_host, ), + timings: params.timings.clone(), + placement: AuctionWaitPlacement::InStream, }; while let Some(step) = hold_step_next_chunk( @@ -2980,6 +3021,7 @@ pub async fn stream_publisher_body_async( orchestrator, services, settings, + AuctionWaitPlacement::PreHeader, ) .await; if body.is_stream() { @@ -3054,6 +3096,8 @@ pub async fn stream_publisher_body_async( services, settings, request_origin: request_origin(¶ms.request_scheme, ¶ms.request_host), + timings: params.timings.clone(), + placement: AuctionWaitPlacement::PreHeader, }, }, ) @@ -3094,6 +3138,7 @@ struct EcSnapshotPreloadInput { auction_needs_row: bool, eid_cookie_may_need_persistence: bool, privacy_needs_row: bool, + snapshot_already_read: bool, } fn should_preload_ec_snapshot(input: &EcSnapshotPreloadInput) -> bool { @@ -3101,7 +3146,8 @@ fn should_preload_ec_snapshot(input: &EcSnapshotPreloadInput) -> bool { let marker_can_skip = input.marker_valid && !input.auction_needs_row && !input.eid_cookie_may_need_persistence - && !input.privacy_needs_row; + && !input.privacy_needs_row + && !input.snapshot_already_read; eligible && !marker_can_skip } @@ -3239,6 +3285,44 @@ fn request_origin(scheme: &str, host: &str) -> String { /// JSON for every non-empty map; `serde_json::from_str` failed and `unwrap_or_default()` /// turned the failure into `{}`. Shared modes therefore served **zero bids**, silently, /// on every request that had any. Every fixture had empty bids, so nothing caught it. +#[derive(Clone, Debug, serde::Serialize)] +#[serde(rename_all = "camelCase")] +struct BrowserAuctionDiagnostics { + #[serde(skip_serializing_if = "Option::is_none")] + auction_dispatched_ms: Option, + #[serde(skip_serializing_if = "Option::is_none")] + auction_resolved_ms: Option, + #[serde(skip_serializing_if = "Option::is_none")] + auction_committed_ms: Option, + #[serde(skip_serializing_if = "Option::is_none")] + auction_wait_ms: Option, + #[serde(skip_serializing_if = "Option::is_none")] + auction_wait_placement: Option<&'static str>, +} + +const fn auction_wait_placement_wire(placement: AuctionWaitPlacement) -> &'static str { + match placement { + AuctionWaitPlacement::PreHeader => "pre_header", + AuctionWaitPlacement::InStream => "in_stream", + } +} + +impl BrowserAuctionDiagnostics { + fn from_request_timings(timings: &RequestTimings) -> Option { + let snapshot = timings.snapshot(); + snapshot.auction_dispatched_ms?; + Some(Self { + auction_dispatched_ms: snapshot.auction_dispatched_ms, + auction_resolved_ms: snapshot.auction_resolved_ms, + auction_committed_ms: snapshot.auction_committed_ms, + auction_wait_ms: snapshot.auction_wait_ms, + auction_wait_placement: snapshot + .auction_wait_placement + .map(auction_wait_placement_wire), + }) + } +} + #[derive(Clone, Default)] pub(crate) struct AdBidsState { /// Rendered bids `` sequences inside the string. pub(crate) fn build_bids_script(bid_map: &serde_json::Map) -> String { + build_bids_script_with_diagnostics(bid_map, None) +} + +fn build_bids_script_with_diagnostics( + bid_map: &serde_json::Map, + auction_diagnostics: Option<&BrowserAuctionDiagnostics>, +) -> String { let json = serde_json::to_string(bid_map) .expect("serde_json::to_string of Map should be infallible"); let escaped = html_escape_for_script(&json); @@ -5959,6 +6177,23 @@ pub(crate) fn build_bids_script(bid_map: &serde_json::Map(function(){{\ +var t=window.tsjs=window.tsjs||{{}};\ +var b=JSON.parse(\"{}\");\ +var d=JSON.parse(\"{}\");\ +var s=t.scheduleInitialAdInit;\ +if(typeof s===\"function\")s(b,void 0,d);\ +else{{t.bids=b;t.auctionDiagnostics=d;}}\ +}})();", + escaped, + html_escape_for_script(&diagnostics) + ); + } + format!( "", + html_escape_for_script(slots_json), + html_escape_for_script(&bids), + html_escape_for_script(&diagnostics) + ); + } + format!( "