diff --git a/crates/trusted-server-adapter-fastly/src/ec_kv.rs b/crates/trusted-server-adapter-fastly/src/ec_kv.rs index c64fb5131..c427ede96 100644 --- a/crates/trusted-server-adapter-fastly/src/ec_kv.rs +++ b/crates/trusted-server-adapter-fastly/src/ec_kv.rs @@ -160,25 +160,22 @@ impl EcKvStore for FastlyEcKvStore { } } - fn count_keys_with_prefix( + fn list_keys_with_prefix( &self, prefix: &str, limit: u32, - ) -> Result> { + ) -> Result, Report> { let store = self.open_store()?; - let page = store + store .build_list() .prefix(prefix) .limit(limit) .execute() + .map(ListPage::into_keys) .change_context(TrustedServerError::KvStore { store_name: self.store_name.clone(), message: format!("Failed to list keys with prefix '{}'", log_id(prefix)), - })?; - - #[allow(clippy::cast_possible_truncation)] - let count = page.keys().len() as u32; - Ok(count) + }) } fn delete(&self, key: &str) -> Result<(), Report> { diff --git a/crates/trusted-server-core/src/ec/finalize.rs b/crates/trusted-server-core/src/ec/finalize.rs index 90e713aff..7cd0cb6e3 100644 --- a/crates/trusted-server-core/src/ec/finalize.rs +++ b/crates/trusted-server-core/src/ec/finalize.rs @@ -6,19 +6,15 @@ use std::collections::HashSet; use edgezero_core::body::Body as EdgeBody; -use error_stack::Report; use http::Response; use super::consent::{ec_consent_granted, ec_consent_withdrawn}; -use crate::error::TrustedServerError; use crate::settings::Settings; use super::EcContext; use super::cookies::{expire_ec_cookie, set_ec_cookie}; use super::generation::{generate_ec_id, is_valid_ec_id}; -use super::kv::{ - CreateIfAbsentOutcome, KvIdentityGraph, TombstoneOutcome, apply_partner_id_updates, -}; +use super::kv::{CreateIfAbsentOutcome, KvIdentityGraph, apply_partner_id_updates}; use super::kv_types::KvEntry; use super::prebid_eids::collect_eid_cookie_updates; use super::pull_sync_marker::{expire_marker, reconcile_marker}; @@ -349,49 +345,30 @@ fn finalize_unusable_consent( // the context — pull sync discloses the raw EC ID to partners — // sees the tombstone that was just written. Only the active ID has // a snapshot in the context to correct. - let outcome = graph.write_withdrawal_tombstone(ec_id, |snapshot| { - if ec_context.ec_value() == Some(ec_id) { - ec_context.set_kv_snapshot(snapshot); - } - }); - log_tombstone_outcome(ec_id, outcome); + let initial = if ec_context.kv_snapshot().belongs_to(ec_id) { + ec_context.kv_snapshot().clone() + } else { + EcKvSnapshot::NotRead + }; + let outcome = graph.tombstone_existing_from_snapshot(ec_id, initial); + // The browser cookie is already cleared, so a failed tombstone + // leaves a live row that server-side consumers still read as + // consented. Report every failure, including the non-active cookie + // ID whose outcome is not retained on the request context. + if matches!(outcome, EcKvSnapshot::Failed { .. }) { + log::warn!( + "EC withdrawal tombstone failed for '{}': the identity-graph row may \ + still be live with consent granted", + log_id(ec_id) + ); + } + if ec_context.ec_value() == Some(ec_id) { + ec_context.set_kv_snapshot(outcome); + } }); } } -/// Records what happened to one withdrawal tombstone. -/// -/// An unknown identity is expected traffic rather than a fault: the identifier -/// comes from a client-supplied cookie, so it may name something this -/// deployment never issued. An error is different: nothing was recorded, so a -/// real row may have gone unmarked, and that is logged as a fault. The browser -/// cookie is expired in every case, and that is the primary enforcement. -fn log_tombstone_outcome( - ec_id: &str, - outcome: Result>, -) { - match outcome { - Ok(TombstoneOutcome::Written) => {} - Ok(TombstoneOutcome::UnknownIdentity) => { - log::debug!( - "Skipping withdrawal tombstone for unknown EC ID '{}'", - log_id(ec_id), - ); - } - Err(err) => { - // Covers both a failed write and a check that could not determine - // whether the identity exists. Either way no marker was recorded, - // so a withdrawal may go unrecorded for the batch-sync window; the - // browser cookie is expired regardless. - log::error!( - "Could not record the withdrawal of EC ID '{}', so it may go unrecorded \ - for the batch-sync window; the browser cookie is still expired: {err:?}", - log_id(ec_id), - ); - } - } -} - fn withdrawal_ec_ids(ec_context: &EcContext) -> HashSet { let mut hashes = HashSet::new(); @@ -1606,6 +1583,154 @@ mod tests { ); } + #[test] + fn finalize_withdrawal_tombstones_both_present_ids_once() { + let settings = create_test_settings(); + let active_ec = sample_ec_id("activ3"); + let cookie_ec = sample_ec_id("cook3e"); + let consent = ConsentContext { + jurisdiction: Jurisdiction::UsState("CA".to_owned()), + gpc: true, + source: ConsentSource::Cookie, + ..Default::default() + }; + let mut ec_context = + make_context_with_consent(Some(&active_ec), Some(&cookie_ec), true, false, consent); + let graph = KvIdentityGraph::in_memory("test_store"); + graph + .create( + &active_ec, + &KvEntry::minimal("active.example.com", "active-uid", 1_000), + ) + .expect("should seed active row"); + graph + .create( + &cookie_ec, + &KvEntry::minimal("cookie.example.com", "cookie-uid", 1_000), + ) + .expect("should seed cookie row"); + ec_context.set_kv_snapshot(graph.load_snapshot(&active_ec)); + let mut response = empty_response(); + + ec_finalize_response( + &settings, + &mut ec_context, + Some(&graph), + &PartnerRegistry::empty(), + None, + None, + &mut response, + ); + + let (active_tombstone, active_generation) = graph + .get(&active_ec) + .expect("should read active row") + .expect("should retain active tombstone"); + let (cookie_tombstone, cookie_generation) = graph + .get(&cookie_ec) + .expect("should read cookie row") + .expect("should retain cookie tombstone"); + assert!( + !active_tombstone.consent.ok, + "active row should be withdrawn" + ); + assert!( + active_tombstone.ids.is_empty(), + "active IDs should be cleared" + ); + assert!( + !cookie_tombstone.consent.ok, + "cookie row should be withdrawn" + ); + assert!( + cookie_tombstone.ids.is_empty(), + "cookie IDs should be cleared" + ); + + let mut repeated_response = empty_response(); + ec_finalize_response( + &settings, + &mut ec_context, + Some(&graph), + &PartnerRegistry::empty(), + None, + None, + &mut repeated_response, + ); + + assert_eq!( + graph + .get(&active_ec) + .expect("should read active row") + .expect("should retain active tombstone") + .1, + active_generation, + "repeated finalization should not rewrite active tombstone" + ); + assert_eq!( + graph + .get(&cookie_ec) + .expect("should read cookie row") + .expect("should retain cookie tombstone") + .1, + cookie_generation, + "repeated finalization should not rewrite cookie tombstone" + ); + } + + #[test] + fn finalize_withdrawal_keeps_cookie_deletion_on_kv_failure() { + let settings = create_test_settings(); + let ec_id = sample_ec_id("failw1"); + let consent = ConsentContext { + jurisdiction: Jurisdiction::UsState("CA".to_owned()), + gpc: true, + source: ConsentSource::Cookie, + ..Default::default() + }; + let mut ec_context = + make_context_with_consent(Some(&ec_id), Some(&ec_id), true, false, consent); + ec_context.set_pull_sync_marker_for_test( + crate::ec::pull_sync_marker::PullSyncMarkerState::Invalid, + ); + let graph = KvIdentityGraph::failing("unavailable-store"); + let mut response = empty_response(); + + ec_finalize_response( + &settings, + &mut ec_context, + Some(&graph), + &PartnerRegistry::empty(), + None, + None, + &mut response, + ); + + let cookies = response + .headers() + .get_all(http::header::SET_COOKIE) + .iter() + .filter_map(|value| value.to_str().ok()) + .collect::>(); + assert_eq!( + response.status(), + 200, + "KV failure should not change response status" + ); + assert!( + cookies + .iter() + .any(|cookie| { cookie.starts_with("ts-ec=;") && cookie.contains("Max-Age=0") }), + "KV failure should not prevent EC cookie deletion" + ); + assert!( + cookies.iter().any(|cookie| { + cookie.starts_with("ts-ec-pull-complete=;") && cookie.contains("Max-Age=0") + }), + "KV failure should not prevent marker deletion" + ); + } + #[test] fn finalize_sets_marker_for_complete_pull_partner_snapshot() { let settings = create_test_settings(); diff --git a/crates/trusted-server-core/src/ec/kv.rs b/crates/trusted-server-core/src/ec/kv.rs index fd9aa5773..d244a15d1 100644 --- a/crates/trusted-server-core/src/ec/kv.rs +++ b/crates/trusted-server-core/src/ec/kv.rs @@ -19,11 +19,10 @@ use error_stack::{Report, ResultExt}; use crate::error::TrustedServerError; -use super::current_timestamp; use super::generation::ec_hash; use super::kv_backend::{EcKvLookup, EcKvStore, EcKvWrite, EcKvWriteMode, EcKvWriteOutcome}; use super::kv_types::{KvEntry, KvMetadata, KvNetwork}; -use super::{EcKvSnapshot, log_id}; +use super::{EcKvSnapshot, checked_current_timestamp, current_timestamp, log_id}; /// Maximum number of CAS retry attempts before giving up. const MAX_CAS_RETRIES: u32 = 5; @@ -39,6 +38,12 @@ const ENTRY_TTL: Duration = Duration::from_secs(365 * 24 * 60 * 60); /// TTL for withdrawal tombstones (24 hours). const TOMBSTONE_TTL: Duration = Duration::from_secs(24 * 60 * 60); +/// Namespace for completion markers written after a withdrawal tombstone. +const WITHDRAWAL_MARKER_PREFIX: &str = "__ts_ec_withdrawal_complete__:"; + +/// Maximum completion markers inspected for one EC ID. +const WITHDRAWAL_MARKER_LIST_LIMIT: u32 = 100; + /// Outcome of an [`KvIdentityGraph::upsert_partner_id_if_exists`] call. /// /// Like [`KvIdentityGraph::upsert_partner_id`], this method fails closed when @@ -133,9 +138,10 @@ impl fmt::Debug for KvIdentityGraph { } /// Result of [`KvIdentityGraph::write_withdrawal_tombstone`]. +#[cfg(test)] #[derive(Debug, Clone, Copy, PartialEq, Eq)] #[must_use] -pub enum TombstoneOutcome { +pub(crate) enum TombstoneOutcome { /// The identity was found and is now tombstoned. Written, /// No such identity is held, so there was nothing to mark withdrawn. @@ -364,9 +370,8 @@ impl KvIdentityGraph { /// - **Existing tombstone** (`consent.ok = false`) — CAS overwrite with /// the new entry. Retries up to [`MAX_CAS_RETRIES`] on conflict. /// - /// Called by `generate_if_needed()` instead of `create()` so that a - /// user who re-consents within the 24-hour tombstone window recovers - /// immediately. + /// This method is reserved for explicit same-key revival. Production EC + /// generation uses [`Self::create_if_absent`] with a freshly generated ID. /// /// # Errors /// @@ -411,6 +416,11 @@ impl KvIdentityGraph { let mut current_gen = generation; for attempt in 0..MAX_CAS_RETRIES { + // A completion marker belongs to the tombstone generation. Remove + // it before making this key live so a later withdrawal cannot be + // suppressed by stale fallback state. + self.clear_withdrawal_marker(ec_id)?; + match self.write_entry( ec_id, &body, @@ -862,12 +872,107 @@ impl KvIdentityGraph { self.store.key_exists(ec_id) } + fn withdrawal_marker_prefix(ec_id: &str) -> String { + format!("{WITHDRAWAL_MARKER_PREFIX}{ec_id}:") + } + + fn withdrawal_marker_key(ec_id: &str, valid_until: u64) -> String { + format!("{}{valid_until}", Self::withdrawal_marker_prefix(ec_id)) + } + + fn withdrawal_marker_keys( + &self, + ec_id: &str, + ) -> Result, Report> { + self.store.list_keys_with_prefix( + &Self::withdrawal_marker_prefix(ec_id), + WITHDRAWAL_MARKER_LIST_LIMIT, + ) + } + + fn withdrawal_marker_exists( + &self, + ec_id: &str, + now: Option, + ) -> Result> { + // An unusable clock cannot prove the tombstone is still valid. Ignore + // completion markers and let withdrawal strongly check the root instead. + let Some(now) = now else { + return Ok(false); + }; + let marker_prefix = Self::withdrawal_marker_prefix(ec_id); + Ok(self + .store + .list_keys_with_prefix(&marker_prefix, WITHDRAWAL_MARKER_LIST_LIMIT)? + .iter() + .filter_map(|key| key.strip_prefix(&marker_prefix)) + .filter_map(|valid_until| valid_until.parse::().ok()) + .any(|valid_until| valid_until > now)) + } + + fn write_withdrawal_marker( + &self, + ec_id: &str, + tombstone_updated: u64, + ) -> Result<(), Report> { + let valid_until = tombstone_updated.saturating_add(TOMBSTONE_TTL.as_secs()); + let marker_key = Self::withdrawal_marker_key(ec_id, valid_until); + match self.store.insert( + &marker_key, + EcKvWrite { + body: "1", + metadata: "{}", + ttl: TOMBSTONE_TTL, + mode: EcKvWriteMode::Add, + }, + )? { + EcKvWriteOutcome::Written | EcKvWriteOutcome::PreconditionFailed => Ok(()), + } + } + + fn clear_withdrawal_marker(&self, ec_id: &str) -> Result<(), Report> { + let marker_keys = self.withdrawal_marker_keys(ec_id)?; + if marker_keys.len() >= WITHDRAWAL_MARKER_LIST_LIMIT as usize { + return Err(self.kv_error(format!( + "Withdrawal marker cleanup exceeds list budget for '{}'", + log_id(ec_id) + ))); + } + + for marker_key in marker_keys { + if let Err(delete_err) = self.store.delete(&marker_key) { + match self.store.key_exists(&marker_key) { + // Another request removed the marker first. + Ok(false) => {} + Ok(true) | Err(_) => return Err(delete_err), + } + } + } + Ok(()) + } + + fn record_withdrawal_completion(&self, ec_id: &str, tombstone_updated: u64) { + if let Err(err) = self.write_withdrawal_marker(ec_id, tombstone_updated) { + // The root is already tombstoned. Preserve that successful privacy + // write even if the cost-control marker cannot be recorded. + log::warn!( + "withdrawal completion marker failed for '{}': {err:?}", + log_id(ec_id) + ); + } + } + /// Writes a withdrawal tombstone for consent enforcement. /// /// Overwrites the entry with `consent.ok = false`, empty partner IDs, /// and a 24-hour TTL. Uses unconditional overwrite (no CAS) since the /// entry is being withdrawn regardless of concurrent state. /// + /// A successful write records a completion marker whose key carries the + /// tombstone's absolute validity bound. Stale misses trust the marker only + /// while the original tombstone should still exist, so a marker that + /// outlives its root cannot suppress withdrawal of a recreated live row. + /// /// The tombstone preserves consent enforcement for batch sync clients /// (`POST /_ts/api/v1/batch-sync`) during the 24-hour revocation window. /// @@ -910,7 +1015,8 @@ impl KvIdentityGraph { /// [`EcKvSnapshot::Failed`]. Callers on the browser path should log at /// `error` level and continue: cookie deletion is the primary enforcement /// mechanism. - pub fn write_withdrawal_tombstone( + #[cfg(test)] + pub(crate) fn write_withdrawal_tombstone( &self, ec_id: &str, record_snapshot: impl FnOnce(EcKvSnapshot), @@ -931,6 +1037,9 @@ impl KvIdentityGraph { }, }); + if let Ok(Some(entry)) = &written { + self.record_withdrawal_completion(ec_id, entry.consent.updated); + } written.map(|entry| { if entry.is_some() { TombstoneOutcome::Written @@ -944,6 +1053,7 @@ impl KvIdentityGraph { /// /// `Ok(None)` means the store does not hold the identity, so nothing was /// written. + #[cfg(test)] fn tombstone_held_identity( &self, ec_id: &str, @@ -956,6 +1066,13 @@ impl KvIdentityGraph { return Ok(None); } + self.overwrite_withdrawal_tombstone(ec_id).map(Some) + } + + fn overwrite_withdrawal_tombstone( + &self, + ec_id: &str, + ) -> Result> { let entry = KvEntry::tombstone(current_timestamp()); let (body, meta_str) = Self::serialize_entry(&entry, self.store_name())?; @@ -966,7 +1083,7 @@ impl KvIdentityGraph { TOMBSTONE_TTL, EcKvWriteMode::Overwrite, ) - .map(|_| Some(entry)) + .map(|_| entry) .map_err(|report| { report.change_context(TrustedServerError::KvStore { store_name: self.store_name().to_owned(), @@ -975,6 +1092,200 @@ impl KvIdentityGraph { }) } + /// Resolves a tombstone attempt whose point read reported the row absent. + /// + /// A completion marker proves an earlier withdrawal finished and avoids a + /// second root existence check. Otherwise, a proven-absent key is a no-op: + /// there is nothing to withdraw, and a forged cookie must not mint a row. A + /// key that provably exists is tombstoned unconditionally because no CAS + /// generation is available after a missed read. Marker-check failure falls + /// back to that privacy write, as does an unusable clock; root-existence + /// failure leaves withdrawal unresolved rather than silently dropped. + fn tombstone_unproven_missing( + &self, + ec_id: &str, + missing: EcKvSnapshot, + now: Option, + ) -> EcKvSnapshot { + match self.withdrawal_marker_exists(ec_id, now) { + Ok(true) => { + log::debug!( + "withdrawal tombstone for '{}': completion marker already exists", + log_id(ec_id) + ); + return missing; + } + Ok(false) => {} + Err(err) => { + // Marker failure must not weaken withdrawal. Fall back to the + // existing root existence check and unconditional privacy write. + log::warn!( + "withdrawal completion marker lookup failed for '{}': {err:?}", + log_id(ec_id) + ); + } + } + + match self.key_exists_confirmed(ec_id) { + Ok(false) => missing, + Ok(true) => { + log::warn!( + "withdrawal tombstone for '{}': point read missed a row the store still \ + lists; writing an unconditional tombstone", + log_id(ec_id) + ); + match self.overwrite_withdrawal_tombstone(ec_id) { + Ok(tombstone) => { + self.record_withdrawal_completion(ec_id, tombstone.consent.updated); + EcKvSnapshot::Present { + ec_id: ec_id.to_owned(), + entry: Box::new(tombstone), + generation: None, + } + } + Err(err) => { + log::warn!( + "unconditional withdrawal tombstone failed for '{}': {err:?}", + log_id(ec_id) + ); + EcKvSnapshot::Failed { + ec_id: ec_id.to_owned(), + } + } + } + } + Err(err) => { + log::warn!( + "withdrawal tombstone for '{}': existence check failed, cannot confirm \ + absence: {err:?}", + log_id(ec_id) + ); + EcKvSnapshot::Failed { + ec_id: ec_id.to_owned(), + } + } + } + } + + /// Writes a tombstone only when an existing row can be confirmed. + /// + /// Existing-key-only behavior is deliberate: a forged or expired `ts-ec` + /// cookie must not mint a row. But a *point read* cannot prove absence on + /// an eventually-consistent store, and dropping a withdrawal is worse than + /// a redundant read, so absence is established in two stages: + /// + /// 1. Any snapshot that is not a usable `Present` for this EC ID — a + /// publisher preload that read `Missing`, a read that `Failed`, or one + /// lacking a CAS generation — is re-read. On the publisher path that + /// re-read is separated from the preload by the full origin round trip, + /// which gives replication time to converge. + /// 2. A re-read that still reports the row absent is checked against + /// [`key_exists_confirmed`](Self::key_exists_confirmed), which reads + /// the primary data source. + /// + /// Resolving the initial snapshot happens outside the retry counter, so all + /// [`MAX_CAS_RETRIES`] iterations stay available for the tombstone write. + pub(crate) fn tombstone_existing_from_snapshot( + &self, + ec_id: &str, + snapshot: EcKvSnapshot, + ) -> EcKvSnapshot { + let mut current = match snapshot { + EcKvSnapshot::Present { + ec_id: ref snapshot_id, + ref entry, + .. + } if snapshot_id == ec_id && !entry.consent.ok => return snapshot, + EcKvSnapshot::Present { + ec_id: ref snapshot_id, + generation: Some(_), + .. + } if snapshot_id == ec_id => snapshot, + _ => self.load_snapshot(ec_id), + }; + + for _attempt in 0..MAX_CAS_RETRIES { + let generation = match current { + EcKvSnapshot::Present { + ec_id: ref snapshot_id, + ref entry, + .. + } if snapshot_id == ec_id && !entry.consent.ok => return current, + EcKvSnapshot::Present { + ec_id: ref snapshot_id, + generation: Some(generation), + .. + } if snapshot_id == ec_id => generation, + // A missing row (including one that disappeared mid-retry) is + // only a no-op once absence is proven against the primary data + // source. + EcKvSnapshot::Missing { + ec_id: ref snapshot_id, + } if snapshot_id == ec_id => { + return self.tombstone_unproven_missing( + ec_id, + current, + checked_current_timestamp(), + ); + } + // A refreshed read that failed (or any other unusable state) + // fails closed rather than silently dropping the withdrawal. + _ => { + return EcKvSnapshot::Failed { + ec_id: ec_id.to_owned(), + }; + } + }; + + let tombstone = KvEntry::tombstone(current_timestamp()); + let Ok((body, meta_str)) = Self::serialize_entry(&tombstone, self.store_name()) else { + return EcKvSnapshot::Failed { + ec_id: ec_id.to_owned(), + }; + }; + match self.write_entry( + ec_id, + &body, + &meta_str, + TOMBSTONE_TTL, + EcKvWriteMode::IfGenerationMatch(generation), + ) { + Ok(EcKvWriteOutcome::Written) => { + self.record_withdrawal_completion(ec_id, tombstone.consent.updated); + return EcKvSnapshot::Present { + ec_id: ec_id.to_owned(), + entry: Box::new(tombstone), + generation: None, + }; + } + Ok(EcKvWriteOutcome::PreconditionFailed) => { + current = self.load_snapshot(ec_id); + } + Err(err) => { + log::warn!( + "conditional withdrawal tombstone failed for '{}': {err:?}", + log_id(ec_id) + ); + return EcKvSnapshot::Failed { + ec_id: ec_id.to_owned(), + }; + } + } + } + + // Withdrawal enforcement lost every CAS race, so the row can still be + // live with consent granted while the browser cookie is cleared. That + // divergence is only visible to operators if it is logged here. + log::warn!( + "withdrawal tombstone for '{}': CAS conflict after {MAX_CAS_RETRIES} retries; the \ + identity-graph row may still be live with consent granted", + log_id(ec_id) + ); + EcKvSnapshot::Failed { + ec_id: ec_id.to_owned(), + } + } + /// Counts the number of keys sharing the same EC hash prefix. /// /// Uses the platform KV list API with a prefix filter, limited to @@ -1076,10 +1387,10 @@ impl KvIdentityGraph { Ok(Some(cluster_size)) } - /// Hard-deletes the entry. + /// Hard-deletes the entry and any withdrawal completion marker. /// - /// Reserved for the IAB data deletion framework (deferred). For consent - /// withdrawal, use [`write_withdrawal_tombstone`](Self::write_withdrawal_tombstone). + /// Reserved for the IAB data deletion framework (deferred). Consent + /// withdrawal uses [`Self::tombstone_existing_from_snapshot`] instead. /// /// # Errors /// @@ -1087,6 +1398,7 @@ impl KvIdentityGraph { pub fn delete(&self, ec_id: &str) -> Result<(), Report> { // The backend's delete already attaches store context, so propagate // without re-wrapping the same message. + self.clear_withdrawal_marker(ec_id)?; self.store.delete(ec_id) } } @@ -1147,6 +1459,85 @@ mod tests { use super::*; use crate::ec::kv_backend::test_support::InMemoryEcKv; + /// [`EcKvStore`] wrapper whose first CAS write both fails the precondition + /// and deletes the key, simulating a concurrent withdrawal that removes the + /// row between this writer's read and its write. + struct DisappearOnConflictEcKv { + inner: InMemoryEcKv, + conflicts_remaining: std::sync::Mutex, + } + + impl DisappearOnConflictEcKv { + fn new(conflicts: u32) -> Self { + Self { + inner: InMemoryEcKv::new("disappear-store"), + conflicts_remaining: std::sync::Mutex::new(conflicts), + } + } + + fn seed_live(&self, ec_id: &str) { + let (body, meta) = + KvIdentityGraph::serialize_entry(&live_entry(), self.inner.store_name()) + .expect("should serialize seeded entry"); + self.inner + .insert( + ec_id, + EcKvWrite { + body: &body, + metadata: &meta, + ttl: ENTRY_TTL, + mode: EcKvWriteMode::Add, + }, + ) + .expect("should seed live entry"); + } + } + + impl EcKvStore for DisappearOnConflictEcKv { + fn store_name(&self) -> &str { + self.inner.store_name() + } + + fn lookup(&self, key: &str) -> Result, Report> { + self.inner.lookup(key) + } + + fn key_exists(&self, key: &str) -> Result> { + self.inner.key_exists(key) + } + + fn insert( + &self, + key: &str, + write: EcKvWrite<'_>, + ) -> Result> { + if matches!(write.mode, EcKvWriteMode::IfGenerationMatch(_)) { + let mut remaining = self + .conflicts_remaining + .lock() + .expect("should lock conflict counter"); + if *remaining > 0 { + *remaining -= 1; + self.inner.delete(key).expect("should delete on conflict"); + return Ok(EcKvWriteOutcome::PreconditionFailed); + } + } + self.inner.insert(key, write) + } + + fn list_keys_with_prefix( + &self, + prefix: &str, + limit: u32, + ) -> Result, Report> { + self.inner.list_keys_with_prefix(prefix, limit) + } + + fn delete(&self, key: &str) -> Result<(), Report> { + self.inner.delete(key) + } + } + fn snapshot_ec_id() -> String { format!("{}.ABC123", "a".repeat(64)) } @@ -1240,6 +1631,21 @@ mod tests { entry } + fn concurrent_live_entry() -> KvEntry { + let mut entry = live_entry(); + entry.ids.insert( + "concurrent.example.com".to_owned(), + crate::ec::kv_types::KvPartnerId { + uid: "concurrent-uid".to_owned(), + }, + ); + entry + } + + // ----------------------------------------------------------------------- + // CAS-conflict injection tests + // ----------------------------------------------------------------------- + /// [`EcKvStore`] wrapper that injects generation conflicts: the first /// `conflicts_remaining` `IfGenerationMatch` inserts return /// [`EcKvWriteOutcome::PreconditionFailed`] without writing, optionally @@ -1248,6 +1654,7 @@ mod tests { inner: InMemoryEcKv, conflicts_remaining: std::sync::Mutex, revive_on_conflict: bool, + partner_update_on_conflict: bool, } impl ConflictInjectingEcKv { @@ -1256,6 +1663,16 @@ mod tests { inner: InMemoryEcKv::new("conflict-store"), conflicts_remaining: std::sync::Mutex::new(conflicts), revive_on_conflict, + partner_update_on_conflict: false, + } + } + + fn with_partner_update_on_conflict(conflicts: u32) -> Self { + Self { + inner: InMemoryEcKv::new("partner-conflict-store"), + conflicts_remaining: std::sync::Mutex::new(conflicts), + revive_on_conflict: true, + partner_update_on_conflict: true, } } @@ -1324,8 +1741,13 @@ mod tests { if self.revive_on_conflict { // Simulate a concurrent writer reviving the entry // between this writer's read and its CAS write. + let concurrent_entry = if self.partner_update_on_conflict { + concurrent_live_entry() + } else { + live_entry() + }; let (body, meta) = KvIdentityGraph::serialize_entry( - &live_entry(), + &concurrent_entry, self.inner.store_name(), ) .expect("should serialize concurrent live entry"); @@ -1347,12 +1769,12 @@ mod tests { self.inner.insert(key, write) } - fn count_keys_with_prefix( + fn list_keys_with_prefix( &self, prefix: &str, limit: u32, - ) -> Result> { - self.inner.count_keys_with_prefix(prefix, limit) + ) -> Result, Report> { + self.inner.list_keys_with_prefix(prefix, limit) } fn delete(&self, key: &str) -> Result<(), Report> { @@ -1615,31 +2037,521 @@ mod tests { } #[test] - fn upsert_partner_id_if_exists_reports_missing_key() { - let kv = KvIdentityGraph::in_memory("test_store"); + fn create_or_revive_fresh_entry_does_not_access_marker_store() { + let operations = Arc::new(RecordedEcKvOperations::default()); + let kv = KvIdentityGraph::new(RecordingEcKv::with_marker_failures( + Arc::clone(&operations), + 0, + )); let ec_id = format!("{}.ABC123", "a".repeat(64)); - let result = kv - .upsert_partner_id_if_exists(&ec_id, "ssp_x", "uid-1") - .expect("should not error on missing key"); - assert_eq!(result, UpsertResult::NotFound); + kv.create_or_revive(&ec_id, &live_entry()) + .expect("should create without reading withdrawal markers"); + + assert_eq!( + operations.list_count(), + 0, + "fresh create should not list withdrawal markers" + ); + assert_eq!( + operations.inserts().len(), + 1, + "fresh create should only insert the live root" + ); + assert!( + kv.get(&ec_id) + .expect("should read entry") + .is_some_and(|(entry, _)| entry.consent.ok), + "fresh create should persist a live entry" + ); } #[test] - fn upsert_partner_id_if_exists_writes_and_detects_unchanged() { + fn create_or_revive_clears_withdrawal_marker() { let kv = KvIdentityGraph::in_memory("test_store"); let ec_id = format!("{}.ABC123", "a".repeat(64)); - kv.create(&ec_id, &live_entry()).expect("should create"); + kv.create(&ec_id, &live_entry()) + .expect("should create live entry"); + let snapshot = kv.load_snapshot(&ec_id); + kv.tombstone_existing_from_snapshot(&ec_id, snapshot); + assert!( + kv.withdrawal_marker_exists(&ec_id, checked_current_timestamp()) + .expect("should read withdrawal marker"), + "withdrawal should record completion" + ); - let first = kv + kv.create_or_revive(&ec_id, &live_entry()) + .expect("should revive tombstone"); + + let (loaded, _) = kv + .get(&ec_id) + .expect("should read revived entry") + .expect("should find revived entry"); + assert!(loaded.consent.ok, "should be live after revive"); + assert!( + !kv.withdrawal_marker_exists(&ec_id, checked_current_timestamp()) + .expect("should read withdrawal marker"), + "revival should clear stale withdrawal completion" + ); + } + + #[test] + fn stale_completion_marker_does_not_suppress_withdrawal_after_root_recreation() { + let operations = Arc::new(RecordedEcKvOperations::default()); + let graph = KvIdentityGraph::new(RecordingEcKv::with_stale_lookups( + Arc::clone(&operations), + 1, + )); + let ec_id = snapshot_ec_id(); + let expired_updated = current_timestamp().saturating_sub(TOMBSTONE_TTL.as_secs() + 1); + let expired_tombstone = KvEntry::tombstone(expired_updated); + graph + .create(&ec_id, &expired_tombstone) + .expect("should seed old tombstone"); + graph + .write_withdrawal_marker(&ec_id, expired_updated) + .expect("should seed old completion marker"); + + // Model the root expiring just before its later-written completion + // marker, followed by recreation of the same key. + graph.store.delete(&ec_id).expect("should expire root row"); + assert_eq!( + graph + .create_if_absent(&ec_id, &live_entry()) + .expect("should recreate expired root"), + CreateIfAbsentOutcome::Written, + "should make the same key live again after root expiry" + ); + + let outcome = graph.tombstone_existing_from_snapshot( + &ec_id, + EcKvSnapshot::Missing { + ec_id: ec_id.clone(), + }, + ); + + assert!( + outcome + .entry_for(&ec_id) + .is_some_and(|entry| !entry.consent.ok), + "an expired completion marker must not suppress withdrawal of a recreated live row" + ); + let (stored, _) = graph + .get(&ec_id) + .expect("should read recreated row") + .expect("should retain recreated row"); + assert!(!stored.consent.ok, "recreated row should be tombstoned"); + } + + #[test] + fn withdrawal_failed_clock_stale_miss_tombstones_live_root() { + let operations = Arc::new(RecordedEcKvOperations::default()); + let graph = KvIdentityGraph::new(RecordingEcKv::with_stale_lookups( + Arc::clone(&operations), + 1, + )); + let ec_id = snapshot_ec_id(); + graph + .create(&ec_id, &concurrent_live_entry()) + .expect("should seed a live row with partner IDs"); + graph + .write_withdrawal_marker(&ec_id, 1000) + .expect("should seed a completion marker left after root recreation"); + operations.reset(); + let missing = graph.load_snapshot(&ec_id); + assert!( + matches!(missing, EcKvSnapshot::Missing { .. }), + "should reproduce a stale point-read miss for the live row" + ); + + // None is the checked clock's failure result, not Unix epoch zero. + let outcome = graph.tombstone_unproven_missing(&ec_id, missing, None); + + let (stored, generation) = graph + .get(&ec_id) + .expect("should read back persisted withdrawal") + .expect("should retain the root"); + assert!( + !stored.consent.ok, + "an unusable clock must not let a marker suppress withdrawal" + ); + assert!( + stored.ids.is_empty(), + "withdrawal should clear stored partner IDs" + ); + assert_eq!( + generation, 2, + "withdrawal should write the existing root once" + ); + assert!( + outcome + .entry_for(&ec_id) + .is_some_and(|entry| !entry.consent.ok), + "clock failure should not fail a strongly confirmed withdrawal" + ); + assert_eq!( + operations.exact_check_count(), + 1, + "should strongly check the root" + ); + assert_eq!( + operations.inserts(), + vec![ + RecordedEcKvInsert { + mode: EcKvWriteMode::Overwrite, + ttl: TOMBSTONE_TTL + }, + RecordedEcKvInsert { + mode: EcKvWriteMode::Add, + ttl: TOMBSTONE_TTL + }, + ], + "should write only the root tombstone and its completion marker" + ); + } + + #[test] + fn withdrawal_failed_clock_does_not_create_absent_root() { + let operations = Arc::new(RecordedEcKvOperations::default()); + let graph = KvIdentityGraph::new(RecordingEcKv::new(Arc::clone(&operations))); + let ec_id = snapshot_ec_id(); + graph + .write_withdrawal_marker(&ec_id, 1000) + .expect("should seed a marker without a root"); + operations.reset(); + let missing = graph.load_snapshot(&ec_id); + + let outcome = graph.tombstone_unproven_missing(&ec_id, missing, None); + + assert!( + matches!(outcome, EcKvSnapshot::Missing { .. }), + "absent root should remain missing" + ); + assert!( + graph + .get(&ec_id) + .expect("should read absent root") + .is_none(), + "clock failure must not mint an unknown root" + ); + assert_eq!( + operations.exact_check_count(), + 1, + "should prove root absence despite the marker" + ); + assert!( + operations.inserts().is_empty(), + "absent identity should cause no writes" + ); + } + + #[test] + fn withdrawal_marker_validity_requires_a_usable_clock_before_expiry() { + let graph = KvIdentityGraph::in_memory("test_store"); + let ec_id = snapshot_ec_id(); + graph + .write_withdrawal_marker(&ec_id, 1000) + .expect("should seed marker"); + let valid_until = 1000 + TOMBSTONE_TTL.as_secs(); + + for (now, expected) in [ + (Some(valid_until - 1), true), + (Some(valid_until), false), + (Some(valid_until + 1), false), + (None, false), + ] { + assert_eq!( + graph + .withdrawal_marker_exists(&ec_id, now) + .expect("should check marker validity"), + expected, + "marker validity should honor clock availability and exclusive expiry at {now:?}" + ); + } + } + + #[test] + fn withdrawal_marker_existence_rejects_malformed_expiry_suffix() { + let ec_id = format!("{}.ABC123", "a".repeat(64)); + let marker_key = KvIdentityGraph::withdrawal_marker_key( + &ec_id, + current_timestamp() + TOMBSTONE_TTL.as_secs(), + ); + let longer_key = format!("{marker_key}-longer"); + let store = InMemoryEcKv::new("test_store"); + store + .insert( + &longer_key, + EcKvWrite { + body: "1", + metadata: "{}", + ttl: TOMBSTONE_TTL, + mode: EcKvWriteMode::Add, + }, + ) + .expect("should seed longer marker key"); + let kv = KvIdentityGraph::new(store); + + assert!( + !kv.withdrawal_marker_exists(&ec_id, checked_current_timestamp()) + .expect("should validate withdrawal marker expiry"), + "a malformed expiry suffix must not prove withdrawal completion" + ); + } + + #[test] + fn clear_absent_withdrawal_marker_lists_once_without_delete() { + let operations = Arc::new(RecordedEcKvOperations::default()); + let kv = KvIdentityGraph::new(RecordingEcKv::new(Arc::clone(&operations))); + let ec_id = snapshot_ec_id(); + + kv.clear_withdrawal_marker(&ec_id) + .expect("should accept an absent marker"); + + assert_eq!( + operations.list_count(), + 1, + "absent marker guard should list once" + ); + assert_eq!( + operations.delete_count(), + 0, + "absent marker guard should avoid a delete" + ); + } + + #[test] + fn writing_an_existing_withdrawal_marker_is_idempotent() { + let operations = Arc::new(RecordedEcKvOperations::default()); + let kv = KvIdentityGraph::new(RecordingEcKv::new(Arc::clone(&operations))); + let ec_id = snapshot_ec_id(); + let updated = current_timestamp(); + + kv.write_withdrawal_marker(&ec_id, updated) + .expect("should write first completion marker"); + kv.write_withdrawal_marker(&ec_id, updated) + .expect("should accept completion marker add collision"); + + assert_eq!( + operations.inserts(), + vec![ + RecordedEcKvInsert { + mode: EcKvWriteMode::Add, + ttl: TOMBSTONE_TTL, + }, + RecordedEcKvInsert { + mode: EcKvWriteMode::Add, + ttl: TOMBSTONE_TTL, + }, + ], + "both marker add attempts should reach the store" + ); + assert!( + kv.withdrawal_marker_exists(&ec_id, checked_current_timestamp()) + .expect("should read completion marker"), + "completion marker should remain valid" + ); + } + + #[test] + fn marker_list_saturation_blocks_revival_and_hard_delete_before_mutation() { + let operations = Arc::new(RecordedEcKvOperations::default()); + let graph = KvIdentityGraph::new(RecordingEcKv::new(Arc::clone(&operations))); + let ec_id = snapshot_ec_id(); + graph + .create(&ec_id, &KvEntry::tombstone(current_timestamp())) + .expect("should seed tombstone"); + let first_valid_until = current_timestamp() + TOMBSTONE_TTL.as_secs(); + for offset in 0..WITHDRAWAL_MARKER_LIST_LIMIT { + let marker_key = KvIdentityGraph::withdrawal_marker_key( + &ec_id, + first_valid_until + u64::from(offset), + ); + graph + .store + .insert( + &marker_key, + EcKvWrite { + body: "1", + metadata: "{}", + ttl: TOMBSTONE_TTL, + mode: EcKvWriteMode::Add, + }, + ) + .expect("should seed completion marker"); + } + operations.reset(); + + let result = graph.create_or_revive(&ec_id, &live_entry()); + + assert!( + result.is_err(), + "revival should fail when marker cleanup cannot prove the prefix is exhausted" + ); + let (stored, _) = graph + .get(&ec_id) + .expect("should read root") + .expect("should retain root"); + assert!(!stored.consent.ok, "root should remain tombstoned"); + assert_eq!( + operations.delete_count(), + 0, + "saturated marker cleanup should not make partial progress" + ); + + operations.reset(); + let delete_result = graph.delete(&ec_id); + + assert!( + delete_result.is_err(), + "hard deletion should fail when marker cleanup cannot prove the prefix is exhausted" + ); + let (stored, _) = graph + .get(&ec_id) + .expect("should read root after rejected hard delete") + .expect("rejected hard delete should retain root"); + assert!( + !stored.consent.ok, + "rejected hard delete should leave the tombstoned root in place" + ); + assert_eq!( + operations.delete_count(), + 0, + "hard deletion should not remove the root before saturated marker cleanup fails" + ); + } + + #[test] + fn clear_withdrawal_marker_accepts_concurrent_removal_after_delete_error() { + let ec_id = format!("{}.ABC123", "a".repeat(64)); + let operations = Arc::new(RecordedEcKvOperations::default()); + let kv = KvIdentityGraph::new(RecordingEcKv::with_marker_delete_failure(operations, true)); + kv.create(&ec_id, &live_entry()) + .expect("should create live entry"); + let snapshot = kv.load_snapshot(&ec_id); + kv.tombstone_existing_from_snapshot(&ec_id, snapshot); + + kv.clear_withdrawal_marker(&ec_id) + .expect("should accept a marker removed by another request"); + + assert!( + !kv.withdrawal_marker_exists(&ec_id, checked_current_timestamp()) + .expect("should confirm marker removal"), + "completion marker should remain absent" + ); + } + + #[test] + fn clear_withdrawal_marker_preserves_delete_error_when_marker_remains() { + let ec_id = format!("{}.ABC123", "a".repeat(64)); + let operations = Arc::new(RecordedEcKvOperations::default()); + let kv = KvIdentityGraph::new(RecordingEcKv::with_marker_delete_failure(operations, false)); + kv.create(&ec_id, &live_entry()) + .expect("should create live entry"); + let snapshot = kv.load_snapshot(&ec_id); + kv.tombstone_existing_from_snapshot(&ec_id, snapshot); + + assert!( + kv.clear_withdrawal_marker(&ec_id).is_err(), + "a marker that remains after delete failure should block revival" + ); + assert!( + kv.withdrawal_marker_exists(&ec_id, checked_current_timestamp()) + .expect("should confirm marker remains"), + "completion marker should remain present" + ); + } + + #[test] + fn withdrawal_marker_is_excluded_from_hash_prefix_count() { + let hash = "a".repeat(64); + let ec_id = format!("{hash}.ABC123"); + let kv = KvIdentityGraph::in_memory("test_store"); + kv.create(&ec_id, &live_entry()) + .expect("should create live entry"); + assert_eq!( + kv.count_hash_prefix_keys(&hash) + .expect("should count live root"), + 1, + "the root should be the only hash-prefix key" + ); + + let snapshot = kv.load_snapshot(&ec_id); + kv.tombstone_existing_from_snapshot(&ec_id, snapshot); + + assert!( + kv.withdrawal_marker_exists(&ec_id, checked_current_timestamp()) + .expect("should read withdrawal marker"), + "withdrawal should record completion" + ); + assert_eq!( + kv.count_hash_prefix_keys(&hash) + .expect("should count tombstoned root"), + 1, + "the marker namespace must not affect cluster counts" + ); + } + + #[test] + fn delete_removes_withdrawal_marker() { + let kv = KvIdentityGraph::in_memory("test_store"); + let ec_id = format!("{}.ABC123", "a".repeat(64)); + kv.create(&ec_id, &live_entry()) + .expect("should create live entry"); + let snapshot = kv.load_snapshot(&ec_id); + kv.tombstone_existing_from_snapshot(&ec_id, snapshot); + + kv.delete(&ec_id).expect("should delete entry and marker"); + + assert!( + kv.get(&ec_id).expect("should read store").is_none(), + "hard delete should remove the root" + ); + assert!( + !kv.withdrawal_marker_exists(&ec_id, checked_current_timestamp()) + .expect("should read withdrawal marker"), + "hard delete should remove withdrawal completion" + ); + } + + #[test] + fn upsert_partner_id_if_exists_reports_missing_key() { + let kv = KvIdentityGraph::in_memory("test_store"); + let ec_id = format!("{}.ABC123", "a".repeat(64)); + + let result = kv + .upsert_partner_id_if_exists(&ec_id, "ssp_x", "uid-1") + .expect("should not error on missing key"); + assert_eq!( + result, + UpsertResult::NotFound, + "should reject a partner upsert for an EC ID the store does not hold" + ); + } + + #[test] + fn upsert_partner_id_if_exists_writes_and_detects_unchanged() { + let kv = KvIdentityGraph::in_memory("test_store"); + let ec_id = format!("{}.ABC123", "a".repeat(64)); + kv.create(&ec_id, &live_entry()).expect("should create"); + + let first = kv .upsert_partner_id_if_exists(&ec_id, "ssp_x", "uid-1") .expect("should write partner id"); - assert_eq!(first, UpsertResult::Written); + assert_eq!( + first, + UpsertResult::Written, + "should report the first partner upsert as written" + ); let second = kv .upsert_partner_id_if_exists(&ec_id, "ssp_x", "uid-1") .expect("should detect unchanged uid"); - assert_eq!(second, UpsertResult::Unchanged); + assert_eq!( + second, + UpsertResult::Unchanged, + "should report an identical partner upsert as unchanged" + ); } #[test] @@ -1652,7 +2564,11 @@ mod tests { let result = kv .upsert_partner_id_if_exists(&ec_id, "ssp_x", "uid-1") .expect("should not error on tombstone"); - assert_eq!(result, UpsertResult::ConsentWithdrawn); + assert_eq!( + result, + UpsertResult::ConsentWithdrawn, + "should reject a partner upsert for a withdrawn identity" + ); } #[test] @@ -1693,30 +2609,342 @@ mod tests { }, ); - assert!(matches!(outcome, EcKvSnapshot::Missing { .. })); + assert!( + matches!(outcome, EcKvSnapshot::Missing { .. }), + "missing snapshot should remain missing" + ); assert!(kv.get(&ec_id).expect("should read store").is_none()); } #[test] - fn write_withdrawal_tombstone_overwrites_live_entry() { + fn tombstone_existing_from_snapshot_never_creates_missing_key() { let kv = KvIdentityGraph::in_memory("test_store"); let ec_id = format!("{}.ABC123", "a".repeat(64)); - kv.create(&ec_id, &live_entry()).expect("should create"); + let snapshot = EcKvSnapshot::Missing { + ec_id: ec_id.clone(), + }; - assert_eq!( - kv.write_withdrawal_tombstone(&ec_id, drop) - .expect("should write tombstone"), - TombstoneOutcome::Written, - "should tombstone an identity the store holds" - ); + let outcome = kv.tombstone_existing_from_snapshot(&ec_id, snapshot); - let (loaded, _) = kv - .get(&ec_id) - .expect("should read entry back") + assert!( + matches!(outcome, EcKvSnapshot::Missing { .. }), + "withdrawal should preserve the missing snapshot for an absent identity" + ); + assert!( + kv.get(&ec_id).expect("should read store").is_none(), + "withdrawal must not create a tombstone for an absent key" + ); + } + + #[test] + fn tombstone_existing_from_snapshot_uses_existing_generation() { + let kv = KvIdentityGraph::in_memory("test_store"); + let ec_id = format!("{}.ABC123", "a".repeat(64)); + kv.create(&ec_id, &live_entry()).expect("should create"); + let snapshot = kv.load_snapshot(&ec_id); + + let outcome = kv.tombstone_existing_from_snapshot(&ec_id, snapshot); + + assert!( + outcome + .entry_for(&ec_id) + .is_some_and(|entry| !entry.consent.ok), + "should return the persisted tombstone" + ); + let (stored, _) = kv + .get(&ec_id) + .expect("should read store") + .expect("should preserve existing key"); + assert!(!stored.consent.ok, "should persist withdrawal state"); + } + + #[test] + fn write_withdrawal_tombstone_overwrites_live_entry() { + let kv = KvIdentityGraph::in_memory("test_store"); + let ec_id = format!("{}.ABC123", "a".repeat(64)); + kv.create(&ec_id, &live_entry()).expect("should create"); + + assert_eq!( + kv.write_withdrawal_tombstone(&ec_id, drop) + .expect("should write tombstone"), + TombstoneOutcome::Written, + "should tombstone an identity the store holds" + ); + + let (loaded, _) = kv + .get(&ec_id) + .expect("should read entry back") .expect("should find tombstone entry"); assert!(!loaded.consent.ok, "should be withdrawn after tombstone"); } + // ----------------------------------------------------------------------- + // Snapshot-aware mutation stores and tests + // ----------------------------------------------------------------------- + + #[derive(Debug, Clone, Copy, PartialEq, Eq)] + struct RecordedEcKvInsert { + mode: EcKvWriteMode, + ttl: Duration, + } + + #[derive(Default)] + struct RecordedEcKvOperations { + lookups: std::sync::atomic::AtomicUsize, + exact_checks: std::sync::atomic::AtomicUsize, + inserts: std::sync::Mutex>, + lists: std::sync::atomic::AtomicUsize, + deletes: std::sync::atomic::AtomicUsize, + } + + impl RecordedEcKvOperations { + fn reset(&self) { + self.lookups.store(0, std::sync::atomic::Ordering::Relaxed); + self.exact_checks + .store(0, std::sync::atomic::Ordering::Relaxed); + self.inserts + .lock() + .expect("should lock recorded inserts") + .clear(); + self.lists.store(0, std::sync::atomic::Ordering::Relaxed); + self.deletes.store(0, std::sync::atomic::Ordering::Relaxed); + } + + fn lookup_count(&self) -> usize { + self.lookups.load(std::sync::atomic::Ordering::Relaxed) + } + + fn exact_check_count(&self) -> usize { + self.exact_checks.load(std::sync::atomic::Ordering::Relaxed) + } + + fn list_count(&self) -> usize { + self.lists.load(std::sync::atomic::Ordering::Relaxed) + } + + fn delete_count(&self) -> usize { + self.deletes.load(std::sync::atomic::Ordering::Relaxed) + } + + fn operation_count(&self) -> usize { + self.lookup_count() + + self.exact_check_count() + + self.inserts().len() + + self.list_count() + + self.delete_count() + } + + fn inserts(&self) -> Vec { + self.inserts + .lock() + .expect("should lock recorded inserts") + .clone() + } + } + + /// In-memory store that records operations and can inject focused failures. + struct RecordingEcKv { + inner: InMemoryEcKv, + operations: Arc, + stale_lookups_remaining: std::sync::Mutex, + lag_live_entries: bool, + marker_operations_fail: bool, + root_check_fails: bool, + marker_delete_failure_removes_key: Option, + } + + impl RecordingEcKv { + fn new(operations: Arc) -> Self { + Self::with_stale_lookups(operations, 0) + } + + fn with_stale_lookups(operations: Arc, stale_lookups: u32) -> Self { + Self { + inner: InMemoryEcKv::new("recording-store"), + operations, + stale_lookups_remaining: std::sync::Mutex::new(stale_lookups), + lag_live_entries: false, + marker_operations_fail: false, + root_check_fails: false, + marker_delete_failure_removes_key: None, + } + } + + fn with_lagging_live_lookups(operations: Arc) -> Self { + Self { + lag_live_entries: true, + ..Self::new(operations) + } + } + + fn with_marker_failures( + operations: Arc, + stale_lookups: u32, + ) -> Self { + Self { + marker_operations_fail: true, + ..Self::with_stale_lookups(operations, stale_lookups) + } + } + + fn completed_with_root_check_failure( + operations: Arc, + ec_id: &str, + ) -> Self { + let store = Self { + root_check_fails: true, + stale_lookups_remaining: std::sync::Mutex::new(u32::MAX), + ..Self::new(operations) + }; + let tombstone = KvEntry::tombstone(current_timestamp()); + let (body, metadata) = + KvIdentityGraph::serialize_entry(&tombstone, store.inner.store_name()) + .expect("should serialize tombstone"); + store + .inner + .insert( + ec_id, + EcKvWrite { + body: &body, + metadata: &metadata, + ttl: TOMBSTONE_TTL, + mode: EcKvWriteMode::Add, + }, + ) + .expect("should seed tombstone"); + store + .inner + .insert( + &KvIdentityGraph::withdrawal_marker_key( + ec_id, + tombstone.consent.updated + TOMBSTONE_TTL.as_secs(), + ), + EcKvWrite { + body: "1", + metadata: "{}", + ttl: TOMBSTONE_TTL, + mode: EcKvWriteMode::Add, + }, + ) + .expect("should seed completion marker"); + store + } + + fn with_marker_delete_failure( + operations: Arc, + remove_before_error: bool, + ) -> Self { + Self { + marker_delete_failure_removes_key: Some(remove_before_error), + ..Self::new(operations) + } + } + + fn marker_error(&self, operation: &str) -> Report { + Report::new(TrustedServerError::KvStore { + store_name: self.inner.store_name().to_owned(), + message: format!("completion marker {operation} failed"), + }) + } + + fn root_check_error(&self) -> Report { + Report::new(TrustedServerError::KvStore { + store_name: self.inner.store_name().to_owned(), + message: "root existence check failed".to_owned(), + }) + } + } + + impl EcKvStore for RecordingEcKv { + fn store_name(&self) -> &str { + self.inner.store_name() + } + + fn lookup(&self, key: &str) -> Result, Report> { + self.operations + .lookups + .fetch_add(1, std::sync::atomic::Ordering::Relaxed); + let mut stale_lookups = self + .stale_lookups_remaining + .lock() + .expect("should lock stale lookup counter"); + if *stale_lookups > 0 { + *stale_lookups -= 1; + return Ok(None); + } + let found = self.inner.lookup(key)?; + if self.lag_live_entries + && found.as_ref().is_some_and(|entry| { + serde_json::from_slice::(&entry.body) + .expect("should decode test entry") + .consent + .ok + }) + { + return Ok(None); + } + Ok(found) + } + + fn key_exists(&self, key: &str) -> Result> { + self.operations + .exact_checks + .fetch_add(1, std::sync::atomic::Ordering::Relaxed); + if !key.starts_with(WITHDRAWAL_MARKER_PREFIX) && self.root_check_fails { + return Err(self.root_check_error()); + } + self.inner.key_exists(key) + } + + fn insert( + &self, + key: &str, + write: EcKvWrite<'_>, + ) -> Result> { + self.operations + .inserts + .lock() + .expect("should lock recorded inserts") + .push(RecordedEcKvInsert { + mode: write.mode, + ttl: write.ttl, + }); + if key.starts_with(WITHDRAWAL_MARKER_PREFIX) && self.marker_operations_fail { + return Err(self.marker_error("write")); + } + self.inner.insert(key, write) + } + + fn list_keys_with_prefix( + &self, + prefix: &str, + limit: u32, + ) -> Result, Report> { + self.operations + .lists + .fetch_add(1, std::sync::atomic::Ordering::Relaxed); + if prefix.starts_with(WITHDRAWAL_MARKER_PREFIX) && self.marker_operations_fail { + return Err(self.marker_error("lookup")); + } + self.inner.list_keys_with_prefix(prefix, limit) + } + + fn delete(&self, key: &str) -> Result<(), Report> { + self.operations + .deletes + .fetch_add(1, std::sync::atomic::Ordering::Relaxed); + if key.starts_with(WITHDRAWAL_MARKER_PREFIX) + && let Some(remove_before_error) = self.marker_delete_failure_removes_key + { + if remove_before_error { + self.inner.delete(key)?; + } + return Err(self.marker_error("delete")); + } + self.inner.delete(key) + } + } + /// [`EcKvStore`] whose reads succeed but every write fails, simulating a /// store that becomes unwritable mid-request. struct WriteFailingEcKv { @@ -1752,12 +2980,12 @@ mod tests { message: "write failing test store".to_owned(), })) } - fn count_keys_with_prefix( + fn list_keys_with_prefix( &self, prefix: &str, limit: u32, - ) -> Result> { - self.inner.count_keys_with_prefix(prefix, limit) + ) -> Result, Report> { + self.inner.list_keys_with_prefix(prefix, limit) } fn delete(&self, key: &str) -> Result<(), Report> { self.inner.delete(key) @@ -1780,77 +3008,667 @@ mod tests { let outcome = graph.upsert_partner_ids_from_snapshot(&ec_id, &updates, snapshot); assert_eq!( - lookups.load(std::sync::atomic::Ordering::Relaxed), - 0, - "a usable generation must avoid the initial read" + lookups.load(std::sync::atomic::Ordering::Relaxed), + 0, + "a usable generation must avoid the initial read" + ); + assert_eq!( + outcome + .entry_for(&ec_id) + .and_then(|entry| entry.ids.get("ssp_x")) + .map(|id| id.uid.as_str()), + Some("uid-1") + ); + assert_eq!( + outcome.generation_for(&ec_id), + None, + "a successful write should clear the stale generation" + ); + } + + #[test] + fn snapshot_upsert_unchanged_updates_preserve_generation() { + let kv = KvIdentityGraph::in_memory("test_store"); + let ec_id = snapshot_ec_id(); + let mut seeded = live_entry(); + apply_partner_id_updates(&mut seeded, &[PartnerIdUpdate::new("ssp_x", "uid-1")]); + kv.create(&ec_id, &seeded).expect("should seed"); + let snapshot = kv.load_snapshot(&ec_id); + assert_eq!( + snapshot.generation_for(&ec_id), + Some(1), + "seeded snapshot should carry its stored generation" + ); + + let outcome = kv.upsert_partner_ids_from_snapshot( + &ec_id, + &[PartnerIdUpdate::new("ssp_x", "uid-1")], + snapshot, + ); + + assert_eq!( + outcome.generation_for(&ec_id), + Some(1), + "an unchanged merge preserves the usable generation and performs no write" + ); + } + + #[test] + fn snapshot_upsert_refreshes_unavailable_generation_exactly_once() { + let lookups = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0)); + let graph = KvIdentityGraph::counting("counting-store", lookups.clone()); + let ec_id = snapshot_ec_id(); + graph.create(&ec_id, &live_entry()).expect("should seed"); + // Finalize-written style snapshot: entry known, generation unavailable. + let snapshot = EcKvSnapshot::Present { + ec_id: ec_id.clone(), + entry: Box::new(live_entry()), + generation: None, + }; + let updates = [PartnerIdUpdate::new("ssp_x", "uid-1")]; + + let outcome = graph.upsert_partner_ids_from_snapshot(&ec_id, &updates, snapshot); + + assert_eq!( + lookups.load(std::sync::atomic::Ordering::Relaxed), + 1, + "an unavailable generation refreshes exactly once before CAS" + ); + assert!( + outcome + .entry_for(&ec_id) + .is_some_and(|e| e.ids.contains_key("ssp_x")), + "refreshing an unavailable generation should persist the partner update" + ); + } + + #[test] + fn snapshot_upsert_gen_unavailable_survives_four_conflicts_then_writes() { + // A generation-unavailable snapshot (finalize-written style) refreshes + // once to obtain a usable generation. That refresh must not consume a + // CAS attempt, so all five write attempts remain: four conflicts + // followed by a successful fifth write still persist the update. + let graph = KvIdentityGraph::new(ConflictInjectingEcKv::new(4, false)); + let ec_id = snapshot_ec_id(); + graph.create(&ec_id, &live_entry()).expect("should seed"); + let snapshot = EcKvSnapshot::Present { + ec_id: ec_id.clone(), + entry: Box::new(live_entry()), + generation: None, + }; + let updates = [PartnerIdUpdate::new("ssp_x", "uid-1")]; + + let outcome = graph.upsert_partner_ids_from_snapshot(&ec_id, &updates, snapshot); + + assert_eq!( + outcome + .entry_for(&ec_id) + .and_then(|entry| entry.ids.get("ssp_x")) + .map(|id| id.uid.as_str()), + Some("uid-1"), + "the fifth CAS attempt must still succeed after a refresh and four conflicts" + ); + } + + #[test] + fn snapshot_upsert_cas_conflict_remerges_concurrent_data() { + let graph = KvIdentityGraph::new(ConflictInjectingEcKv::new(1, true)); + let ec_id = snapshot_ec_id(); + graph.create(&ec_id, &live_entry()).expect("should seed"); + let snapshot = graph.load_snapshot(&ec_id); + let updates = [PartnerIdUpdate::new("ssp_x", "uid-1")]; + + let outcome = graph.upsert_partner_ids_from_snapshot(&ec_id, &updates, snapshot); + + let entry = outcome + .entry_for(&ec_id) + .expect("should persist re-merged entry"); + assert_eq!( + entry.ids.get("ssp_x").map(|id| id.uid.as_str()), + Some("uid-1"), + "conflict must re-merge our update onto the concurrently revived row" + ); + assert!(entry.consent.ok, "concurrent revive keeps the row live"); + } + + #[test] + fn snapshot_upsert_revalidates_transient_missing_and_persists() { + // An eventually-consistent point read earlier in the request missed a + // row that exists. Named routes such as `/auction` never run orphan + // recovery, so this refresh is the request's only chance to persist the + // collected partner IDs. + let kv = KvIdentityGraph::in_memory("test_store"); + let ec_id = snapshot_ec_id(); + kv.create(&ec_id, &live_entry()).expect("should seed live"); + let updates = [PartnerIdUpdate::new("ssp_x", "uid-1")]; + + let outcome = kv.upsert_partner_ids_from_snapshot( + &ec_id, + &updates, + EcKvSnapshot::Missing { + ec_id: ec_id.clone(), + }, + ); + + assert_eq!( + outcome + .entry_for(&ec_id) + .and_then(|entry| entry.ids.get("ssp_x").map(|id| id.uid.clone())), + Some("uid-1".to_owned()), + "a stale miss must be revalidated before the updates are dropped" + ); + let (stored, _) = kv + .get(&ec_id) + .expect("should read store") + .expect("row should remain"); + assert_eq!( + stored.ids.get("ssp_x").map(|id| id.uid.as_str()), + Some("uid-1"), + "the revalidated update must reach the store" + ); + } + + #[test] + fn snapshot_upsert_confirmed_missing_still_never_creates() { + // Revalidation only changes what a *stale* miss does. A row that is + // genuinely absent on the refresh must stay absent. + let kv = KvIdentityGraph::in_memory("test_store"); + let ec_id = snapshot_ec_id(); + let updates = [PartnerIdUpdate::new("ssp_x", "uid-1")]; + + let outcome = kv.upsert_partner_ids_from_snapshot( + &ec_id, + &updates, + EcKvSnapshot::Missing { + ec_id: ec_id.clone(), + }, + ); + + assert!( + matches!(outcome, EcKvSnapshot::Missing { .. }), + "a confirmed miss must stay missing" + ); + assert!( + kv.get(&ec_id).expect("should read store").is_none(), + "must not create a root entry for a missing key" + ); + } + + #[test] + fn snapshot_upsert_failed_snapshot_is_not_revalidated() { + // A lookup that already errored is not retried on the hot path, even + // though the row exists. + let kv = KvIdentityGraph::in_memory("test_store"); + let ec_id = snapshot_ec_id(); + kv.create(&ec_id, &live_entry()).expect("should seed live"); + let updates = [PartnerIdUpdate::new("ssp_x", "uid-1")]; + + let outcome = kv.upsert_partner_ids_from_snapshot( + &ec_id, + &updates, + EcKvSnapshot::Failed { + ec_id: ec_id.clone(), + }, + ); + + assert!( + matches!(outcome, EcKvSnapshot::Failed { .. }), + "a failed lookup must not be retried by partner enrichment" + ); + } + + #[test] + fn snapshot_upsert_rejects_tombstone() { + let kv = KvIdentityGraph::in_memory("test_store"); + let ec_id = snapshot_ec_id(); + kv.create(&ec_id, &KvEntry::tombstone(1000)) + .expect("should seed tombstone"); + let snapshot = kv.load_snapshot(&ec_id); + let updates = [PartnerIdUpdate::new("ssp_x", "uid-1")]; + + let outcome = kv.upsert_partner_ids_from_snapshot(&ec_id, &updates, snapshot); + + assert!( + outcome + .entry_for(&ec_id) + .is_some_and(|entry| entry.ids.is_empty()), + "a tombstone must reject partner enrichment" + ); + let (stored, _) = kv + .get(&ec_id) + .expect("should read store") + .expect("tombstone should remain"); + assert!(stored.ids.is_empty(), "no update should reach the store"); + } + + #[test] + fn snapshot_upsert_store_failure_returns_failed_not_request_local() { + let graph = KvIdentityGraph::new(WriteFailingEcKv::new()); + let ec_id = snapshot_ec_id(); + let snapshot = EcKvSnapshot::Present { + ec_id: ec_id.clone(), + entry: Box::new(live_entry()), + generation: Some(1), + }; + let updates = [PartnerIdUpdate::new("ssp_x", "uid-1")]; + + let outcome = graph.upsert_partner_ids_from_snapshot(&ec_id, &updates, snapshot); + + assert!( + matches!(outcome, EcKvSnapshot::Failed { .. }), + "a store write failure must not claim request-local IDs were persisted" + ); + } + + #[test] + fn tombstone_existing_from_snapshot_skips_backend_for_authoritative_tombstone() { + for generation in [Some(7), None] { + let operations = Arc::new(RecordedEcKvOperations::default()); + let graph = KvIdentityGraph::new(RecordingEcKv::new(Arc::clone(&operations))); + let ec_id = snapshot_ec_id(); + let snapshot = EcKvSnapshot::Present { + ec_id: ec_id.clone(), + entry: Box::new(KvEntry::tombstone(1_000)), + generation, + }; + + let outcome = graph.tombstone_existing_from_snapshot(&ec_id, snapshot.clone()); + + assert_eq!( + outcome, snapshot, + "should preserve authoritative tombstone state" + ); + assert_eq!( + operations.operation_count(), + 0, + "an authoritative tombstone should not access the backend" + ); + } + } + + #[test] + fn tombstone_existing_from_snapshot_repeated_request_preserves_first_write() { + let operations = Arc::new(RecordedEcKvOperations::default()); + let graph = KvIdentityGraph::new(RecordingEcKv::new(Arc::clone(&operations))); + let ec_id = snapshot_ec_id(); + graph + .create(&ec_id, &live_entry()) + .expect("should seed live row"); + let live_snapshot = graph.load_snapshot(&ec_id); + operations.reset(); + + graph.tombstone_existing_from_snapshot(&ec_id, live_snapshot); + + assert_eq!( + operations.lookup_count(), + 0, + "usable generation should avoid a read" + ); + assert_eq!( + operations.inserts(), + vec![ + RecordedEcKvInsert { + mode: EcKvWriteMode::IfGenerationMatch(1), + ttl: TOMBSTONE_TTL, + }, + RecordedEcKvInsert { + mode: EcKvWriteMode::Add, + ttl: TOMBSTONE_TTL, + }, + ], + "first withdrawal should write the root and its completion marker" + ); + let first_snapshot = graph.load_snapshot(&ec_id); + let (first_entry, first_generation) = match &first_snapshot { + EcKvSnapshot::Present { + entry, generation, .. + } => (entry.as_ref().clone(), *generation), + other => panic!("should load first tombstone, got {other:?}"), + }; + operations.reset(); + + let second_outcome = graph.tombstone_existing_from_snapshot(&ec_id, first_snapshot); + + assert_eq!( + operations.operation_count(), + 0, + "repeated withdrawal should not access the backend" + ); + assert_eq!( + second_outcome.generation_for(&ec_id), + first_generation, + "repeated withdrawal should preserve the stored generation" + ); + assert_eq!( + second_outcome + .entry_for(&ec_id) + .map(|entry| entry.consent.updated), + Some(first_entry.consent.updated), + "repeated withdrawal should preserve the first tombstone timestamp" + ); + } + + #[test] + fn tombstone_existing_from_repeated_stale_miss_preserves_first_write() { + let operations = Arc::new(RecordedEcKvOperations::default()); + let graph = KvIdentityGraph::new(RecordingEcKv::with_stale_lookups( + Arc::clone(&operations), + 2, + )); + let ec_id = snapshot_ec_id(); + graph + .create(&ec_id, &live_entry()) + .expect("should seed live row"); + operations.reset(); + + let first_outcome = graph.tombstone_existing_from_snapshot( + &ec_id, + EcKvSnapshot::Missing { + ec_id: ec_id.clone(), + }, + ); + let first_updated = first_outcome + .entry_for(&ec_id) + .expect("should return first tombstone") + .consent + .updated; + assert_eq!( + operations.lookup_count(), + 1, + "first stale miss should retry its point read once" + ); + assert_eq!( + operations.list_count(), + 1, + "first stale miss should list completion markers once" + ); + assert_eq!( + operations.exact_check_count(), + 1, + "first stale miss should confirm root existence once" + ); + assert_eq!( + operations.inserts().len(), + 2, + "first stale miss should write the root and completion marker" + ); + operations.reset(); + + graph.tombstone_existing_from_snapshot( + &ec_id, + EcKvSnapshot::Missing { + ec_id: ec_id.clone(), + }, + ); + + assert_eq!( + operations.lookup_count(), + 1, + "a repeated stale miss should retry its point read once" + ); + assert_eq!( + operations.exact_check_count(), + 0, + "a valid completion marker should avoid a root existence check" + ); + assert_eq!( + operations.list_count(), + 1, + "a repeated stale miss should list completion markers once" + ); + assert_eq!( + operations.delete_count(), + 0, + "a repeated withdrawal should not delete marker state" + ); + assert!( + operations.inserts().is_empty(), + "a repeated stale miss should not rewrite the completed tombstone" + ); + let (stored, generation) = graph + .get(&ec_id) + .expect("should read stored tombstone") + .expect("should preserve tombstone"); + assert_eq!( + generation, 2, + "only the first withdrawal should advance the root generation" + ); + assert_eq!( + stored.consent.updated, first_updated, + "repeated withdrawal should preserve the first tombstone timestamp" + ); + } + + #[test] + fn tombstone_existing_from_stale_parallel_snapshot_stops_after_conflict() { + let operations = Arc::new(RecordedEcKvOperations::default()); + let graph = KvIdentityGraph::new(RecordingEcKv::new(Arc::clone(&operations))); + let ec_id = snapshot_ec_id(); + graph + .create(&ec_id, &live_entry()) + .expect("should seed live row"); + let stale_snapshot = graph.load_snapshot(&ec_id); + operations.reset(); + + let first_outcome = graph.tombstone_existing_from_snapshot(&ec_id, stale_snapshot.clone()); + let second_outcome = graph.tombstone_existing_from_snapshot(&ec_id, stale_snapshot); + + assert_eq!( + operations.inserts(), + vec![ + RecordedEcKvInsert { + mode: EcKvWriteMode::IfGenerationMatch(1), + ttl: TOMBSTONE_TTL, + }, + RecordedEcKvInsert { + mode: EcKvWriteMode::Add, + ttl: TOMBSTONE_TTL, + }, + RecordedEcKvInsert { + mode: EcKvWriteMode::IfGenerationMatch(1), + ttl: TOMBSTONE_TTL, + }, + ], + "parallel loser should attempt stale CAS once and never replace the winner" + ); + assert_eq!( + operations.lookup_count(), + 1, + "parallel loser should reread exactly once after its conflict" + ); + assert_eq!( + second_outcome.generation_for(&ec_id), + Some(2), + "parallel loser should return the winner's stored generation" + ); + assert_eq!( + second_outcome + .entry_for(&ec_id) + .map(|entry| entry.consent.updated), + first_outcome + .entry_for(&ec_id) + .map(|entry| entry.consent.updated), + "parallel loser should preserve the winner's tombstone" + ); + } + + #[test] + fn tombstone_existing_from_snapshot_succeeds_without_backend_for_tombstone() { + let graph = KvIdentityGraph::failing("unavailable-store"); + let ec_id = snapshot_ec_id(); + let snapshot = EcKvSnapshot::Present { + ec_id: ec_id.clone(), + entry: Box::new(KvEntry::tombstone(1_000)), + generation: Some(3), + }; + + let outcome = graph.tombstone_existing_from_snapshot(&ec_id, snapshot.clone()); + + assert_eq!( + outcome, snapshot, + "authoritative tombstone should not touch unavailable backend" + ); + } + + #[test] + fn tombstone_existing_from_snapshot_non_authoritative_states_reread_live_row() { + let ec_id = snapshot_ec_id(); + let states = [ + EcKvSnapshot::NotRead, + EcKvSnapshot::Failed { + ec_id: ec_id.clone(), + }, + EcKvSnapshot::Present { + ec_id: ec_id.clone(), + entry: Box::new(live_entry()), + generation: None, + }, + EcKvSnapshot::Present { + ec_id: "different-ec-id".to_owned(), + entry: Box::new(KvEntry::tombstone(1_000)), + generation: Some(9), + }, + ]; + + for state in states { + let operations = Arc::new(RecordedEcKvOperations::default()); + let graph = KvIdentityGraph::new(RecordingEcKv::new(Arc::clone(&operations))); + graph + .create(&ec_id, &live_entry()) + .expect("should seed live row"); + operations.reset(); + + let outcome = graph.tombstone_existing_from_snapshot(&ec_id, state); + + assert_eq!( + operations.lookup_count(), + 1, + "state should force one reread" + ); + assert_eq!( + operations.inserts().len(), + 2, + "live reread should write the root and completion marker" + ); + assert!( + outcome + .entry_for(&ec_id) + .is_some_and(|entry| !entry.consent.ok), + "reread live row should be tombstoned" + ); + } + } + + #[test] + fn tombstone_stale_miss_valid_marker_skips_root_check_and_writes() { + let ec_id = snapshot_ec_id(); + let operations = Arc::new(RecordedEcKvOperations::default()); + let graph = KvIdentityGraph::new(RecordingEcKv::completed_with_root_check_failure( + Arc::clone(&operations), + &ec_id, + )); + + let outcome = graph.tombstone_existing_from_snapshot( + &ec_id, + EcKvSnapshot::Missing { + ec_id: ec_id.clone(), + }, + ); + + assert!( + matches!(outcome, EcKvSnapshot::Missing { .. }), + "a completion marker should resolve a repeated stale miss" + ); + assert_eq!( + operations.lookup_count(), + 1, + "should retry the stale point read once" + ); + assert_eq!( + operations.list_count(), + 1, + "should strongly list completion markers once" ); assert_eq!( - outcome - .entry_for(&ec_id) - .and_then(|entry| entry.ids.get("ssp_x")) - .map(|id| id.uid.as_str()), - Some("uid-1") + operations.exact_check_count(), + 0, + "valid marker should bypass the failing root check" + ); + assert!( + operations.inserts().is_empty(), + "valid marker should prevent redundant writes" ); - assert_eq!(outcome.generation_for(&ec_id), None); } #[test] - fn snapshot_upsert_unchanged_updates_preserve_generation() { - let kv = KvIdentityGraph::in_memory("test_store"); + fn tombstone_stale_miss_still_writes_when_marker_operations_fail() { + let operations = Arc::new(RecordedEcKvOperations::default()); + let graph = KvIdentityGraph::new(RecordingEcKv::with_marker_failures( + Arc::clone(&operations), + 1, + )); let ec_id = snapshot_ec_id(); - let mut seeded = live_entry(); - apply_partner_id_updates(&mut seeded, &[PartnerIdUpdate::new("ssp_x", "uid-1")]); - kv.create(&ec_id, &seeded).expect("should seed"); - let snapshot = kv.load_snapshot(&ec_id); - assert_eq!(snapshot.generation_for(&ec_id), Some(1)); + graph + .create(&ec_id, &live_entry()) + .expect("should seed live row"); + operations.reset(); - let outcome = kv.upsert_partner_ids_from_snapshot( + let outcome = graph.tombstone_existing_from_snapshot( &ec_id, - &[PartnerIdUpdate::new("ssp_x", "uid-1")], - snapshot, + EcKvSnapshot::Missing { + ec_id: ec_id.clone(), + }, ); + assert!( + outcome + .entry_for(&ec_id) + .is_some_and(|entry| !entry.consent.ok), + "marker failures must not suppress the withdrawal write" + ); assert_eq!( - outcome.generation_for(&ec_id), - Some(1), - "an unchanged merge preserves the usable generation and performs no write" + operations.inserts(), + vec![ + RecordedEcKvInsert { + mode: EcKvWriteMode::Overwrite, + ttl: TOMBSTONE_TTL, + }, + RecordedEcKvInsert { + mode: EcKvWriteMode::Add, + ttl: TOMBSTONE_TTL, + }, + ], + "failed completion recording should still be attempted after the root write" ); + let (stored, _) = graph + .get(&ec_id) + .expect("should read stored row") + .expect("should preserve the root"); + assert!(!stored.consent.ok, "root should remain tombstoned"); } #[test] - fn snapshot_upsert_refreshes_unavailable_generation_exactly_once() { - let lookups = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0)); - let graph = KvIdentityGraph::counting("counting-store", lookups.clone()); + fn tombstone_existing_from_snapshot_retries_cas_conflict() { + let graph = KvIdentityGraph::new(ConflictInjectingEcKv::new(1, false)); let ec_id = snapshot_ec_id(); graph.create(&ec_id, &live_entry()).expect("should seed"); - // Finalize-written style snapshot: entry known, generation unavailable. - let snapshot = EcKvSnapshot::Present { - ec_id: ec_id.clone(), - entry: Box::new(live_entry()), - generation: None, - }; - let updates = [PartnerIdUpdate::new("ssp_x", "uid-1")]; + let snapshot = graph.load_snapshot(&ec_id); - let outcome = graph.upsert_partner_ids_from_snapshot(&ec_id, &updates, snapshot); + let outcome = graph.tombstone_existing_from_snapshot(&ec_id, snapshot); - assert_eq!( - lookups.load(std::sync::atomic::Ordering::Relaxed), - 1, - "an unavailable generation refreshes exactly once before CAS" - ); assert!( outcome .entry_for(&ec_id) - .is_some_and(|e| e.ids.contains_key("ssp_x")) + .is_some_and(|entry| !entry.consent.ok), + "should retry the conflict and persist the tombstone" ); } #[test] - fn snapshot_upsert_gen_unavailable_survives_four_conflicts_then_writes() { - // A generation-unavailable snapshot (finalize-written style) refreshes - // once to obtain a usable generation. That refresh must not consume a - // CAS attempt, so all five write attempts remain: four conflicts - // followed by a successful fifth write still persist the update. + fn tombstone_gen_unavailable_survives_four_conflicts_then_writes() { + // A generation-unavailable snapshot refreshes once before its CAS. That + // refresh must not spend a CAS attempt, so a withdrawal tombstone still + // persists after four conflicts and a successful fifth write. let graph = KvIdentityGraph::new(ConflictInjectingEcKv::new(4, false)); let ec_id = snapshot_ec_id(); graph.create(&ec_id, &live_entry()).expect("should seed"); @@ -1859,168 +3677,151 @@ mod tests { entry: Box::new(live_entry()), generation: None, }; - let updates = [PartnerIdUpdate::new("ssp_x", "uid-1")]; - let outcome = graph.upsert_partner_ids_from_snapshot(&ec_id, &updates, snapshot); + let outcome = graph.tombstone_existing_from_snapshot(&ec_id, snapshot); - assert_eq!( + assert!( outcome .entry_for(&ec_id) - .and_then(|entry| entry.ids.get("ssp_x")) - .map(|id| id.uid.as_str()), - Some("uid-1"), - "the fifth CAS attempt must still succeed after a refresh and four conflicts" + .is_some_and(|entry| !entry.consent.ok), + "the fifth CAS attempt must persist the tombstone after a refresh and four conflicts" ); } #[test] - fn snapshot_upsert_cas_conflict_remerges_concurrent_data() { - let graph = KvIdentityGraph::new(ConflictInjectingEcKv::new(1, true)); + fn tombstone_existing_from_snapshot_returns_failed_after_cas_exhaustion() { + // Every CAS attempt loses its race, so the row stays live with consent + // granted while the browser cookie is already cleared. The caller must + // see a failure it can report rather than a silent no-op. + let graph = KvIdentityGraph::new(ConflictInjectingEcKv::new(MAX_CAS_RETRIES, false)); let ec_id = snapshot_ec_id(); graph.create(&ec_id, &live_entry()).expect("should seed"); let snapshot = graph.load_snapshot(&ec_id); - let updates = [PartnerIdUpdate::new("ssp_x", "uid-1")]; - let outcome = graph.upsert_partner_ids_from_snapshot(&ec_id, &updates, snapshot); + let outcome = graph.tombstone_existing_from_snapshot(&ec_id, snapshot); + + assert!( + matches!(outcome, EcKvSnapshot::Failed { .. }), + "CAS exhaustion must report a failed withdrawal" + ); + let (stored, _) = graph + .get(&ec_id) + .expect("should read store") + .expect("row should remain"); + assert!( + stored.consent.ok, + "the row is still live, which is exactly why the failure must be reported" + ); + } + + #[test] + fn tombstone_existing_from_snapshot_overrides_concurrent_live_update() { + let graph = KvIdentityGraph::new(ConflictInjectingEcKv::with_partner_update_on_conflict(1)); + let ec_id = snapshot_ec_id(); + graph + .create(&ec_id, &live_entry()) + .expect("should seed live row"); + let snapshot = graph.load_snapshot(&ec_id); + + let outcome = graph.tombstone_existing_from_snapshot(&ec_id, snapshot); let entry = outcome .entry_for(&ec_id) - .expect("should persist re-merged entry"); - assert_eq!( - entry.ids.get("ssp_x").map(|id| id.uid.as_str()), - Some("uid-1"), - "conflict must re-merge our update onto the concurrently revived row" + .expect("should return persisted tombstone"); + assert!(!entry.consent.ok, "withdrawal should win after retry"); + assert!( + entry.ids.is_empty(), + "withdrawal should clear concurrent partner IDs" ); - assert!(entry.consent.ok, "concurrent revive keeps the row live"); } #[test] - fn snapshot_upsert_revalidates_transient_missing_and_persists() { - // An eventually-consistent point read earlier in the request missed a - // row that exists. Named routes such as `/auction` never run orphan - // recovery, so this refresh is the request's only chance to persist the - // collected partner IDs. - let kv = KvIdentityGraph::in_memory("test_store"); + fn upsert_partner_id_rejects_tombstone() { + let graph = KvIdentityGraph::in_memory("test-store"); let ec_id = snapshot_ec_id(); - kv.create(&ec_id, &live_entry()).expect("should seed live"); - let updates = [PartnerIdUpdate::new("ssp_x", "uid-1")]; + graph + .create(&ec_id, &KvEntry::tombstone(1_000)) + .expect("should seed tombstone"); - let outcome = kv.upsert_partner_ids_from_snapshot( - &ec_id, - &updates, - EcKvSnapshot::Missing { - ec_id: ec_id.clone(), - }, - ); + let result = graph.upsert_partner_id(&ec_id, "ssp.example.com", "uid-1"); - assert_eq!( - outcome - .entry_for(&ec_id) - .and_then(|entry| entry.ids.get("ssp_x").map(|id| id.uid.clone())), - Some("uid-1".to_owned()), - "a stale miss must be revalidated before the updates are dropped" - ); - let (stored, _) = kv + assert!(result.is_err(), "public upsert should reject a tombstone"); + let (stored, _) = graph .get(&ec_id) .expect("should read store") - .expect("row should remain"); - assert_eq!( - stored.ids.get("ssp_x").map(|id| id.uid.as_str()), - Some("uid-1"), - "the revalidated update must reach the store" + .expect("should preserve tombstone"); + assert!(!stored.consent.ok, "entry should remain withdrawn"); + assert!( + stored.ids.is_empty(), + "upsert should not repopulate partner IDs" ); } #[test] - fn snapshot_upsert_confirmed_missing_still_never_creates() { - // Revalidation only changes what a *stale* miss does. A row that is - // genuinely absent on the refresh must stay absent. - let kv = KvIdentityGraph::in_memory("test_store"); + fn tombstone_existing_from_snapshot_store_failure_returns_failed() { + let graph = KvIdentityGraph::new(WriteFailingEcKv::new()); let ec_id = snapshot_ec_id(); - let updates = [PartnerIdUpdate::new("ssp_x", "uid-1")]; + let snapshot = EcKvSnapshot::Present { + ec_id: ec_id.clone(), + entry: Box::new(live_entry()), + generation: Some(1), + }; - let outcome = kv.upsert_partner_ids_from_snapshot( - &ec_id, - &updates, - EcKvSnapshot::Missing { - ec_id: ec_id.clone(), - }, + let outcome = graph.tombstone_existing_from_snapshot(&ec_id, snapshot); + + assert!( + matches!(outcome, EcKvSnapshot::Failed { .. }), + "a failed tombstone write should return failed snapshot state" ); + } + + #[test] + fn tombstone_existing_from_snapshot_noop_when_row_disappears_on_retry() { + let store = DisappearOnConflictEcKv::new(1); + store.seed_live(&snapshot_ec_id()); + let graph = KvIdentityGraph::new(store); + let ec_id = snapshot_ec_id(); + let snapshot = graph.load_snapshot(&ec_id); + + let outcome = graph.tombstone_existing_from_snapshot(&ec_id, snapshot); assert!( matches!(outcome, EcKvSnapshot::Missing { .. }), - "a confirmed miss must stay missing" + "a row that disappears during retry becomes a no-op" ); assert!( - kv.get(&ec_id).expect("should read store").is_none(), - "must not create a root entry for a missing key" + graph.get(&ec_id).expect("should read store").is_none(), + "must not recreate the disappeared key" ); } #[test] - fn snapshot_upsert_failed_snapshot_is_not_revalidated() { - // A lookup that already errored is not retried on the hot path, even - // though the row exists. + fn tombstone_existing_from_snapshot_reretries_failed_snapshot_read() { + // A prior request-scoped read failed, so the snapshot is `Failed`. A + // withdrawal must not silently drop consent removal: re-read the store + // and tombstone the row if it is authoritatively present. let kv = KvIdentityGraph::in_memory("test_store"); let ec_id = snapshot_ec_id(); kv.create(&ec_id, &live_entry()).expect("should seed live"); - let updates = [PartnerIdUpdate::new("ssp_x", "uid-1")]; - let outcome = kv.upsert_partner_ids_from_snapshot( + let outcome = kv.tombstone_existing_from_snapshot( &ec_id, - &updates, EcKvSnapshot::Failed { ec_id: ec_id.clone(), }, ); - assert!( - matches!(outcome, EcKvSnapshot::Failed { .. }), - "a failed lookup must not be retried by partner enrichment" - ); - } - - #[test] - fn snapshot_upsert_rejects_tombstone() { - let kv = KvIdentityGraph::in_memory("test_store"); - let ec_id = snapshot_ec_id(); - kv.create(&ec_id, &KvEntry::tombstone(1000)) - .expect("should seed tombstone"); - let snapshot = kv.load_snapshot(&ec_id); - let updates = [PartnerIdUpdate::new("ssp_x", "uid-1")]; - - let outcome = kv.upsert_partner_ids_from_snapshot(&ec_id, &updates, snapshot); - assert!( outcome .entry_for(&ec_id) - .is_some_and(|entry| entry.ids.is_empty()), - "a tombstone must reject partner enrichment" + .is_some_and(|entry| !entry.consent.ok), + "a failed snapshot must re-read and persist the tombstone" ); let (stored, _) = kv .get(&ec_id) .expect("should read store") - .expect("tombstone should remain"); - assert!(stored.ids.is_empty(), "no update should reach the store"); - } - - #[test] - fn snapshot_upsert_store_failure_returns_failed_not_request_local() { - let graph = KvIdentityGraph::new(WriteFailingEcKv::new()); - let ec_id = snapshot_ec_id(); - let snapshot = EcKvSnapshot::Present { - ec_id: ec_id.clone(), - entry: Box::new(live_entry()), - generation: Some(1), - }; - let updates = [PartnerIdUpdate::new("ssp_x", "uid-1")]; - - let outcome = graph.upsert_partner_ids_from_snapshot(&ec_id, &updates, snapshot); - - assert!( - matches!(outcome, EcKvSnapshot::Failed { .. }), - "a store write failure must not claim request-local IDs were persisted" - ); + .expect("should preserve existing key"); + assert!(!stored.consent.ok, "withdrawal must reach the store"); } #[test] @@ -2263,9 +4064,8 @@ mod tests { #[test] fn withdrawal_does_not_lose_a_new_identity_to_lookup_lag() { - let (mut store, _) = CountingEcKv::new(); - store.lag_live_entries = true; - let kv = KvIdentityGraph::new(store); + let operations = Arc::new(RecordedEcKvOperations::default()); + let kv = KvIdentityGraph::new(RecordingEcKv::with_lagging_live_lookups(operations)); let ec_id = format!("{}.ABC123", "a".repeat(64)); kv.create(&ec_id, &live_entry()) .expect("should issue an identity"); @@ -2293,10 +4093,8 @@ mod tests { fn a_withdrawal_checks_strong_existence_once_without_eventual_lookup() { // One existence operation may page at the backend, but must not be // followed by an eventually consistent lookup. - let (store, reads) = CountingEcKv::new(); - let lookups = std::sync::Arc::clone(&store.lookups); - let kv = KvIdentityGraph::new(store); - let count = || *reads.lock().expect("should lock the read counter"); + let operations = Arc::new(RecordedEcKvOperations::default()); + let kv = KvIdentityGraph::new(RecordingEcKv::new(Arc::clone(&operations))); let absent = format!("{}.ABC123", "e".repeat(64)); assert_eq!( @@ -2305,11 +4103,15 @@ mod tests { TombstoneOutcome::UnknownIdentity, "an absent identity is not held" ); - assert_eq!(count(), 1, "an unknown identity costs a single read"); + assert_eq!( + operations.exact_check_count(), + 1, + "an unknown identity should cost one strong read" + ); let held = format!("{}.ABC123", "a".repeat(64)); kv.create(&held, &live_entry()).expect("should create"); - let before = count(); + let before = operations.exact_check_count(); assert_eq!( kv.write_withdrawal_tombstone(&held, drop) .expect("should resolve the withdrawal"), @@ -2317,12 +4119,12 @@ mod tests { "a held identity is tombstoned" ); assert_eq!( - count() - before, + operations.exact_check_count() - before, 1, "a held identity is checked once, then written" ); assert_eq!( - *lookups.lock().expect("should lock lookup counter"), + operations.lookup_count(), 0, "should not use an eventual lookup for withdrawal" ); @@ -2350,78 +4152,6 @@ mod tests { ); } - /// Store double that records how many reads reach it. - struct CountingEcKv { - inner: super::super::kv_backend::test_support::InMemoryEcKv, - reads: std::sync::Arc>, - lag_live_entries: bool, - lookups: std::sync::Arc>, - } - - impl CountingEcKv { - /// Returns the store and a handle to its counter, which stays readable - /// after the store moves into the graph. - fn new() -> (Self, std::sync::Arc>) { - let counter = std::sync::Arc::new(std::sync::Mutex::new(0)); - ( - Self { - inner: super::super::kv_backend::test_support::InMemoryEcKv::new("test_store"), - reads: std::sync::Arc::clone(&counter), - lag_live_entries: false, - lookups: std::sync::Arc::new(std::sync::Mutex::new(0)), - }, - counter, - ) - } - } - - impl EcKvStore for CountingEcKv { - fn store_name(&self) -> &str { - self.inner.store_name() - } - - fn lookup(&self, key: &str) -> Result, Report> { - *self.lookups.lock().expect("should lock the lookup counter") += 1; - let found = self.inner.lookup(key)?; - if self.lag_live_entries - && found.as_ref().is_some_and(|entry| { - serde_json::from_slice::(&entry.body) - .expect("should decode test entry") - .consent - .ok - }) - { - return Ok(None); - } - Ok(found) - } - - fn key_exists(&self, key: &str) -> Result> { - *self.reads.lock().expect("should lock the read counter") += 1; - self.inner.key_exists(key) - } - - fn insert( - &self, - key: &str, - write: EcKvWrite<'_>, - ) -> Result> { - self.inner.insert(key, write) - } - - fn count_keys_with_prefix( - &self, - prefix: &str, - limit: u32, - ) -> Result> { - self.inner.count_keys_with_prefix(prefix, limit) - } - - fn delete(&self, key: &str) -> Result<(), Report> { - self.inner.delete(key) - } - } - /// Store double whose reads always fail while writes still work. struct ReadFailingEcKv { inner: super::super::kv_backend::test_support::InMemoryEcKv, @@ -2461,12 +4191,12 @@ mod tests { // Left working so a test can prove no row was created without going // through the read path it just made fail. - fn count_keys_with_prefix( + fn list_keys_with_prefix( &self, prefix: &str, limit: u32, - ) -> Result> { - self.inner.count_keys_with_prefix(prefix, limit) + ) -> Result, Report> { + self.inner.list_keys_with_prefix(prefix, limit) } fn delete(&self, key: &str) -> Result<(), Report> { diff --git a/crates/trusted-server-core/src/ec/kv_backend.rs b/crates/trusted-server-core/src/ec/kv_backend.rs index 43b835972..bfe74ca6f 100644 --- a/crates/trusted-server-core/src/ec/kv_backend.rs +++ b/crates/trusted-server-core/src/ec/kv_backend.rs @@ -7,7 +7,7 @@ //! `trusted-server-adapter-fastly`). //! //! This trait is intentionally narrow: lookup with a generation marker, -//! conditional insert, prefix counting, and delete. Conditional writes are +//! conditional insert, strong prefix listing, and delete. Conditional writes are //! expressed through [`EcKvWriteMode`] so compare-and-swap loops stay in //! core while the platform supplies the actual precondition mechanics. @@ -123,7 +123,32 @@ pub trait EcKvStore { write: EcKvWrite<'_>, ) -> Result>; - /// Counts keys sharing the given prefix, up to `limit`. + /// Lists keys sharing the given prefix from strongly consistent state. + /// + /// Returns one bounded page of at most `limit` keys, not an exhaustive scan. + /// Results may be truncated; a full page does not prove that all matching + /// keys were returned. Callers must handle truncation conservatively (for + /// example, reject a saturated cleanup) and validate the complete key shape + /// before using key contents for a correctness decision. + /// + /// Completed writes must be visible even when [`Self::lookup`] lags. If the + /// backend cannot provide strong listing, it must return an error rather + /// than substitute eventually consistent results. + /// + /// # Errors + /// + /// Returns [`TrustedServerError::KvStore`] on store open or list failure. + fn list_keys_with_prefix( + &self, + prefix: &str, + limit: u32, + ) -> Result, Report>; + + /// Counts keys in one strongly consistent prefix-list page, up to `limit`. + /// + /// Delegates to [`Self::list_keys_with_prefix`] and inherits its bounded, + /// potentially truncated result. This is a capped count, not an exhaustive + /// total; callers needing completeness must account for truncation. /// /// # Errors /// @@ -132,7 +157,11 @@ pub trait EcKvStore { &self, prefix: &str, limit: u32, - ) -> Result>; + ) -> Result> { + let count = self.list_keys_with_prefix(prefix, limit)?.len(); + #[allow(clippy::cast_possible_truncation)] + Ok(count as u32) + } /// Hard-deletes a key. /// @@ -191,12 +220,12 @@ pub(crate) mod test_support { self.inner.insert(key, write) } - fn count_keys_with_prefix( + fn list_keys_with_prefix( &self, prefix: &str, limit: u32, - ) -> Result> { - self.inner.count_keys_with_prefix(prefix, limit) + ) -> Result, Report> { + self.inner.list_keys_with_prefix(prefix, limit) } fn delete(&self, key: &str) -> Result<(), Report> { @@ -207,7 +236,7 @@ pub(crate) mod test_support { /// [`EcKvStore`] wrapper that models an eventually-consistent point read. /// /// The first `stale_lookups` calls to [`EcKvStore::lookup`] report the key - /// absent while [`EcKvStore::count_keys_with_prefix`] — the list API, which + /// absent while [`EcKvStore::list_keys_with_prefix`] — the list API, which /// reads the primary data source — still sees it. Writes reach the inner /// store, so a test can assert what actually persisted. /// @@ -264,18 +293,18 @@ pub(crate) mod test_support { self.inner.insert(key, write) } - fn count_keys_with_prefix( + fn list_keys_with_prefix( &self, prefix: &str, limit: u32, - ) -> Result> { + ) -> Result, Report> { if self.list_fails { return Err(Report::new(TrustedServerError::KvStore { store_name: self.inner.store_name().to_owned(), message: "list unavailable".to_owned(), })); } - self.inner.count_keys_with_prefix(prefix, limit) + self.inner.list_keys_with_prefix(prefix, limit) } fn delete(&self, key: &str) -> Result<(), Report> { @@ -356,19 +385,18 @@ pub(crate) mod test_support { Ok(EcKvWriteOutcome::Written) } - fn count_keys_with_prefix( + fn list_keys_with_prefix( &self, prefix: &str, limit: u32, - ) -> Result> { + ) -> Result, Report> { let entries = self.entries.lock().expect("should lock in-memory store"); - let count = entries + Ok(entries .keys() .filter(|key| key.starts_with(prefix)) .take(limit as usize) - .count(); - #[allow(clippy::cast_possible_truncation)] - Ok(count as u32) + .cloned() + .collect()) } fn delete(&self, key: &str) -> Result<(), Report> { @@ -419,11 +447,11 @@ pub(crate) mod test_support { Err(self.error("insert")) } - fn count_keys_with_prefix( + fn list_keys_with_prefix( &self, _prefix: &str, _limit: u32, - ) -> Result> { + ) -> Result, Report> { Err(self.error("list")) } diff --git a/crates/trusted-server-core/src/ec/mod.rs b/crates/trusted-server-core/src/ec/mod.rs index e3476a2ba..0ce6dd13a 100644 --- a/crates/trusted-server-core/src/ec/mod.rs +++ b/crates/trusted-server-core/src/ec/mod.rs @@ -98,6 +98,9 @@ pub enum EcKvSnapshot { /// The store authoritatively reported that this EC ID does not exist. Missing { ec_id: String }, /// Persisted entry data, optionally with a generation usable for CAS. + /// + /// A generation never authorizes a write by itself. Callers must first + /// enforce entry policy such as rejecting a withdrawal tombstone. Present { ec_id: String, entry: Box, @@ -659,20 +662,28 @@ impl EcContext { } } -/// Returns the current Unix timestamp in seconds. +/// Returns the current Unix timestamp in seconds, falling back to zero on clock failure. /// /// Uses [`web_time::SystemTime`], which maps to `std::time::SystemTime` on /// native and `wasm32-wasip1` targets and to a JS-backed clock on /// `wasm32-unknown-unknown` (Cloudflare Workers), where `std::time` is not /// available. pub(crate) fn current_timestamp() -> u64 { + checked_current_timestamp().unwrap_or(0) +} + +/// Returns the current Unix timestamp, or `None` when the clock precedes the epoch. +/// +/// Use this instead of [`current_timestamp`] when a fallback could authorize +/// a time-bounded correctness decision. +pub(crate) fn checked_current_timestamp() -> Option { web_time::SystemTime::now() .duration_since(web_time::UNIX_EPOCH) - .map(|d| d.as_secs()) - .unwrap_or_else(|err| { - log::error!("SystemTime::now() failed, falling back to epoch 0: {err}"); - 0 + .map(|duration| duration.as_secs()) + .map_err(|err| { + log::error!("SystemTime::now() failed: {err}"); }) + .ok() } #[cfg(test)] @@ -731,12 +742,12 @@ mod tests { } self.inner.insert(key, write) } - fn count_keys_with_prefix( + fn list_keys_with_prefix( &self, prefix: &str, limit: u32, - ) -> Result> { - self.inner.count_keys_with_prefix(prefix, limit) + ) -> Result, Report> { + self.inner.list_keys_with_prefix(prefix, limit) } fn delete(&self, key: &str) -> Result<(), Report> { self.inner.delete(key) diff --git a/crates/trusted-server-core/src/publisher.rs b/crates/trusted-server-core/src/publisher.rs index b2c71c749..12b4f3702 100644 --- a/crates/trusted-server-core/src/publisher.rs +++ b/crates/trusted-server-core/src/publisher.rs @@ -7290,12 +7290,12 @@ mod tests { ) -> Result> { self.inner.insert(key, write) } - fn count_keys_with_prefix( + fn list_keys_with_prefix( &self, prefix: &str, limit: u32, - ) -> Result> { - self.inner.count_keys_with_prefix(prefix, limit) + ) -> Result, Report> { + self.inner.list_keys_with_prefix(prefix, limit) } fn delete(&self, key: &str) -> Result<(), Report> { self.inner.delete(key) diff --git a/docs/guide/edge-cookies.md b/docs/guide/edge-cookies.md index aa1710647..6cf8ae78f 100644 --- a/docs/guide/edge-cookies.md +++ b/docs/guide/edge-cookies.md @@ -124,7 +124,7 @@ flowchart TD - **Non-regulated**: EC always allowed. - **Unknown**: Fail-closed when jurisdiction cannot be determined. -The `ec_identity_store` KV store is the only EC lifecycle store. It holds identity graph state, source-domain keyed partner UIDs, a minimal consent snapshot used for EC entry metadata, and withdrawal tombstones. Consent interpretation for each request remains based on the live request signals listed above. +The `ec_identity_store` KV store is the only EC lifecycle store. It holds identity graph state, source-domain keyed partner UIDs, a minimal consent snapshot used for EC entry metadata, withdrawal tombstones, and completion markers that prevent stale point-read misses from rewriting completed tombstones. A marker key records the original tombstone's validity bound and is ignored at or after that time, even if its KV row has not expired yet. If the clock is unusable, the marker is ignored and withdrawal falls back to a strong root-existence check; only a confirmed existing root can be written. This prevents a stale marker from suppressing withdrawal when an expired EC key is created again. With a healthy store, repeated withdrawal of an already tombstoned EC ID leaves the row unchanged instead of refreshing its 24-hour TTL. Consent interpretation for each request remains based on the live request signals listed above. ## Partner Sync Channels diff --git a/docs/superpowers/plans/2026-07-13-issue-881-idempotent-withdrawal-tombstones.md b/docs/superpowers/plans/2026-07-13-issue-881-idempotent-withdrawal-tombstones.md new file mode 100644 index 000000000..bad834714 --- /dev/null +++ b/docs/superpowers/plans/2026-07-13-issue-881-idempotent-withdrawal-tombstones.md @@ -0,0 +1,423 @@ +# Issue #881: Idempotent EC Withdrawal Tombstones Plan + +- **Date:** 2026-07-13 +- **Status:** Implemented; clock-safety review fixes independently reviewed and full local gates passed +- **Issue:** [#881 — Make EC withdrawal tombstoning idempotent across request bursts](https://github.com/IABTechLab/trusted-server/issues/881) +- **Stack base:** [Draft PR #900 — Avoid no-op EC KV reads in post-send pull sync](https://github.com/IABTechLab/trusted-server/pull/900) +- **Underlying dependency:** [PR #885 — Request-scoped EC KV snapshot and orphan recovery](https://github.com/IABTechLab/trusted-server/pull/885) + +## Goal + +Make explicit EC withdrawal idempotent across repeated and concurrent requests +without weakening the existing-key-only privacy invariant introduced by PR +#885. The first successful withdrawal of a live row writes a CAS-protected +24-hour tombstone and a completion marker whose key records the tombstone's +absolute validity bound. A request that already has authoritative tombstone +state returns without reading or writing KV. When an eventually consistent +point read instead misses the existing tombstone, a strongly listed marker +prevents an unconditional replacement write only while the original tombstone +should still exist. An expired marker cannot suppress withdrawal of a recreated +live row. An unusable clock cannot establish marker validity and must fall back +to strong root existence, never skip withdrawal or create an unknown root. + +Browser-cookie deletion remains synchronous and best-effort KV failure must +never block the response. + +## Clarified Semantics + +- An authoritative missing row is a no-op. A valid-looking but unverified + browser cookie must never create a KV root. +- A matching authoritative tombstone snapshot is returned unchanged before any + lookup, serialization, or write, regardless of whether its generation is + available. +- A matching live snapshot with a generation uses one conditional write. +- A live snapshot without a generation, a failed/not-read snapshot, or a + snapshot for another EC ID rereads the requested row before deciding. +- CAS conflicts reread and retry. If another withdrawal has already written a + tombstone, retry ends without another write. If a concurrent live update or + re-consent changed the generation first, withdrawal retries against that live + row and tombstones it. +- A later re-consent may legitimately win when it linearizes after a completed + or no-op withdrawal. Idempotency does not impose global withdrawal priority. +- Repeated withdrawal preserves the original tombstone expiration because it + performs no second root write. A completion-marker insert is attempted only + after the first successful tombstone write. +- A repeated stale point-read miss uses the strongly consistent completion + marker to avoid another unconditional root write. The marker key records + `consent.updated + TOMBSTONE_TTL` and is ignored at or after that bound. + `checked_current_timestamp()` preserves clock failure as `None`; it cannot + use the legacy `current_timestamp()` epoch-zero sentinel to validate a marker. + Without a usable clock, strong root existence still gates any overwrite. +- Completion-marker failure never suppresses the privacy write. Explicit + same-key CAS revival and hard deletion clear the marker before replacing or + removing the tombstoned generation. +- Production EC creation and re-consent mint fresh IDs through + `create_if_absent`; they do not revive a tombstone or read marker state. +- When cookie and active EC IDs differ, every valid existing row is withdrawn + independently; missing or malformed IDs are never created. +- `ts-ec` and the pull-completeness marker are expired before best-effort KV + work. Store failure is logged/swallowed by finalization. + +## Non-Goals + +- Do not recreate missing roots from browser cookies. +- Do not add a separate deduplication store, cross-request lock, or withdrawal + cookie. The stale-miss completion marker lives in the existing EC store. +- Do not restrict withdrawal handling to document navigations. +- Do not change the 24-hour tombstone duration or root-entry KV schema. +- Do not change EC generation, browser pull-marker behavior, batch sync, pull sync, or + partner-upsert semantics beyond preserving tombstone rejection. +- Do not make withdrawal dominate a re-consent that occurs after withdrawal's + linearization point. +- Do not rewrite archival lifecycle descriptions; add only the completion-key + namespace annotation to the historical technical spec. +- Defer a marker-value redesign: an eventually consistent point read cannot + replace strong listing for completion proof. +- Defer Fastly list `ItemNotFound` normalization until an operation-specific + platform contract confirms that it means an empty key set. SDK propagation + alone does not establish that contract. + +## Original Baseline (Before Issue #881) + +`KvIdentityGraph::tombstone_existing_from_snapshot` already preserves PR #885's +existing-key-only and CAS behavior, but it writes a fresh tombstone whenever the +snapshot is `Present`, including when that entry is already a tombstone. Parallel +withdrawals therefore converge safely but still perform a redundant CAS write, +and repeated requests reset the tombstone's 24-hour TTL. + +An unconditional root overwrite remains necessary when point reads miss a root +that the strong list still sees. Without durable completion state, repeated +stale misses bypass the `Present` fast path and keep rewriting that root. + +Finalization already: + +- expires browser state before KV work; +- collects both valid cookie and active EC IDs; +- uses the carried snapshot only for its matching ID; +- independently resolves the other ID; +- logs/swallows KV failures. + +The implementation should therefore remain concentrated in the core KV method, +with finalization changes limited to acceptance-level integration tests unless a +test exposes a defect. + +## Proposed Design + +### 1. Add an authoritative tombstone fast path + +At the beginning of each `tombstone_existing_from_snapshot` retry iteration, +match a `Present` snapshot only when its `ec_id` equals the requested ID. + +- If `entry.consent.ok == false`, return the exact snapshot immediately. +- Perform this check before requiring a generation or constructing a new + `KvEntry::tombstone`. +- Preserve the snapshot's entry timestamp and generation exactly. + +Then retain the existing state machine: + +| Initial/refreshed state | Action | +| ----------------------------------------- | ------------------------------------- | +| Matching tombstone | Return unchanged; zero backend work | +| Matching live entry + generation | CAS-write a 24-hour tombstone | +| Matching live entry without generation | Reread | +| Refreshed `Missing` | Prove absence or use guarded fallback | +| `Failed`, `NotRead`, or wrong-ID snapshot | Reread requested ID | +| CAS precondition failure | Reread and retry | +| Row disappears during retry | Return `Missing` | +| Store failure or retry exhaustion | Return ID-bound `Failed` | + +A successful tombstone remains `consent.ok = false`, has empty partner IDs, and +uses `TOMBSTONE_TTL`. + +Each successful root tombstone also creates an add-only completion marker in the +same EC store with `TOMBSTONE_TTL`. Its key includes the original tombstone's +absolute validity bound. The marker namespace cannot collide with EC IDs and is +excluded from hash-prefix cluster counts. If a later point read misses and a +strong prefix list returns a correctly shaped marker whose bound is still in +the future, withdrawal returns without a second strong root-existence check or +root rewrite. This is a bounded prefix list plus expiry-suffix validation, not +an exact-key existence check: each completion key includes its own validity +bound. Expired or malformed markers and an unusable clock are ignored. Marker +list failure falls back to strong root existence before any privacy write; +marker write failure preserves an already completed root write. Explicit +same-key CAS revival and hard deletion remove markers for that EC ID. Fresh-ID +creation does not read marker state. + +### 2. Contain the unconditional fallback + +Production withdrawal uses `tombstone_existing_from_snapshot`; its +`tombstone_unproven_missing` fallback checks marker validity, then confirms +strong root existence before calling `overwrite_withdrawal_tombstone`. +`write_withdrawal_tombstone` and `tombstone_held_identity` are test-only helpers +that retain the legacy strong-existence-then-overwrite behavior. Record +completion after a successful root write so another stale miss cannot refresh +the root. The marker is +best-effort after a successful root write; marker failure is logged and cannot +turn a completed privacy write into a reported failure. + +Update hard deletion and same-key revival to remove completion state. Historical +lifecycle descriptions remain unchanged apart from the namespace annotation. + +### 3. Prove operation-level idempotency + +Use a focused recording backend around the existing in-memory store. It must +count every point lookup, exact existence check, insert, prefix list, and delete +attempt, including inserts that return a CAS precondition failure, and record +each write's mode and TTL. Tests must show: + +- a supplied matching tombstone returns with zero lookups and zero insert + attempts; +- generation-unavailable tombstone state also performs no backend operation; +- the first live withdrawal performs one `IfGenerationMatch` root insert and + one add-only completion-marker insert with `TOMBSTONE_TTL`, while a repeated + authoritative tombstone causes zero additional insert attempts and leaves + stored root generation and `consent.updated` unchanged; +- two consecutive stale point-read misses cause one unconditional root write; + the second request observes the strong completion marker and leaves the root + generation and `consent.updated` unchanged; +- an expired root recreated under the same key is tombstoned when the point + read misses, because the older marker's absolute bound is no longer valid; +- an unusable clock plus a marker and stale point-read miss still tombstones a + strongly confirmed live row, clearing partner IDs; an absent root stays absent; +- marker validity is checked immediately before, at, and after its expiry bound; +- two stale live snapshots model parallel requests: the first writes the + tombstone; the second conflicts, rereads the tombstone, and performs no + replacement write; +- a supplied tombstone succeeds even against an always-failing backend, proving + no hidden operation; +- live state with generation avoids an initial read and uses one CAS write; +- missing state remains a no-op; +- `NotRead`, failed, generation-unavailable, and wrong-ID states reread the + requested row before applying the documented write/no-create behavior. + +The recording backend's first-write TTL assertion plus zero additional insert +attempts on repetition is the authoritative proof that the original tombstone +TTL was not refreshed. + +### 4. Preserve race ordering and tombstone authority + +Extend conflict tests to model a concurrent live update/re-consent that changes +the generation before withdrawal's first CAS. Withdrawal must reread the live +row and eventually write the tombstone, clearing any partner IDs. + +Retain the existing batch-sync conditional-upsert and snapshot bulk-upsert +tombstone tests, and add focused coverage for the public single-partner upsert +path so all live enrichment APIs are proven unable to repopulate tombstones: + +- `upsert_partner_id_if_exists` rejects tombstones; +- `upsert_partner_id` returns an error and leaves tombstone IDs empty; +- snapshot bulk upsert cannot repopulate a tombstone; +- a disappeared row is not recreated; +- store errors return failed state rather than claiming persistence. + +### 5. Verify finalization behavior + +Retain the existing finalization coverage for malformed/absent IDs and a +present active ID plus missing secondary ID. Add only the missing integration +cases: + +- differing valid cookie and active EC IDs are both tombstoned when both rows + exist; +- repeated finalization preserves existing tombstone generations; +- a failing KV graph still returns the response and emits both applicable + browser-cookie expiration headers. + +Production finalization should not change unless these tests expose a defect. + +## File Map + +### Modify + +- `crates/trusted-server-core/src/ec/kv.rs` + - Add the matching-tombstone no-op branch. + - Gate the stale-miss fallback with a same-store completion marker. + - Clear completion markers on explicit same-key CAS revival and hard deletion. + - Keep fresh-ID creation independent of marker availability. + - Update withdrawal documentation. + - Add operation-count, repetition, stale-read, and concurrency tests. + - Pass the checked timestamp into the private missing-row/marker helpers for + deterministic clock-failure tests without a global clock override. +- `crates/trusted-server-core/src/ec/mod.rs` + - Add an optional checked timestamp; preserve legacy epoch-zero fallback for + unrelated callers. +- `crates/trusted-server-core/src/ec/kv_backend.rs` + - Document strong single-page listing, truncation, and caller responsibilities. +- `docs/superpowers/specs/2026-03-24-ssc-technical-spec-design.md` + - Annotate the internal completion-marker namespace only. +- `crates/trusted-server-core/src/ec/finalize.rs` + - Add two-ID, repeated-withdrawal, and KV-failure integration coverage. +- `docs/guide/edge-cookies.md` + - Document that repeated withdrawal preserves the first tombstone TTL. + +### Add + +- `docs/superpowers/plans/2026-07-13-issue-881-idempotent-withdrawal-tombstones.md` + - Record the reviewed design and verification contract. + +No dependency, configuration, adapter, JavaScript, or public wire-format change +is included. The EC store gains an internal completion-marker key namespace. + +## Implementation Tasks + +### Task 1 — Establish failing idempotency tests + +- [x] Add a recording withdrawal backend that counts every backend operation + and captures each insert's write mode and TTL. +- [x] Add a supplied-tombstone test proving zero lookups and zero inserts. +- [x] Add repeated and stale-parallel snapshot tests proving the first insert is + `IfGenerationMatch` with `TOMBSTONE_TTL`, then no further insert occurs and + generation/`consent.updated` remain unchanged. +- [x] Add live-generation, unavailable-generation, `NotRead`, failed, wrong-ID, + and missing state tests. +- [x] Run `cargo test-fastly tombstone_existing_from_snapshot` and confirm the + new repeated/no-backend tests fail before implementation. + +### Task 2 — Implement the no-op branch and guard the fallback + +- [x] Return a matching tombstone snapshot before generation lookup, + serialization, or write. +- [x] Keep live/missing/failed/mismatched/CAS behavior unchanged. +- [x] Record completion after successful tombstone writes. +- [x] Strongly list and validate completion markers before checking root + existence so a repeated stale miss performs one strong read. +- [x] Suppress repeated stale-miss overwrites only while a marker's encoded + tombstone-validity bound is still in the future. +- [x] Clear completion state on explicit same-key CAS revival and hard deletion. +- [x] Keep fresh-ID creation independent of marker availability. +- [x] Run focused KV tests until green. + +### Task 3 — Cover concurrent state changes + +- [x] Model another withdrawal winning between read and CAS; prove the loser + rereads and stops without replacing the tombstone. +- [x] Model a concurrent live update/re-consent winning before CAS; prove + withdrawal retries and tombstones the refreshed row. +- [x] Assert final tombstones contain no partner IDs. +- [x] Add direct `upsert_partner_id` tombstone rejection coverage and re-run the + existing conditional and snapshot-bulk rejection tests. + +### Task 4 — Verify finalization + +- [x] Add a test with both differing valid IDs present and assert both become + tombstones. +- [x] Retain the existing missing and invalid ID tests unchanged as regression + coverage. +- [x] Add repeated-finalization generation-stability coverage. +- [x] Add a failing-store test proving EC and marker cookie deletion survives KV + failure. +- [x] Run focused withdrawal/finalization tests. + +### Task 5 — Original review and full verification (historical) + +- [x] Run independent correctness/concurrency and test-quality reviews. +- [x] Apply only fixes required by issue scope. +- [x] Mark this plan implemented only after all checks below pass. + +### Review follow-up — Clock safety + +The original checklist above records the prior implementation, not fresh full +verification of these review fixes. + +- [x] Preserve failed clock reads as `None` at completion-marker validation. +- [x] Reproduce live-row suppression before the guard; prove persisted consent + is false and partner IDs are empty afterward using `RecordingEcKv`. +- [x] Cover absent roots and valid/at-expiry/expired marker boundaries. +- [x] Verify restoring the unsafe epoch-zero fallback fails the regression, + then revert the mutation. +- [x] Add missing assertion messages and precise marker test names/counts. +- [x] Clarify bounded strong listing and the completion-key namespace. +- [x] Rerun every full repository verification gate for the review follow-up. + +Follow-up validation passed all six clippy targets, all four adapter test +commands, parity and CLI tests, Rust/JS/docs formatting, JS tests/build, and the +release Fastly WASM build. Rust suites reported 3,208 passed and 10 ignored; +Vitest reported 959 passed. Independent correctness review found no defects. +Clock failure is injected at the private withdrawal decision; an actual failed +host clock and production Fastly list-error semantics were not exercised. + +## Acceptance Mapping + +| Issue requirement | Planned evidence | +| ------------------------------------------------ | ----------------------------------------------------------------------------- | +| First withdrawal establishes a 24-hour tombstone | Live-entry CAS test and existing `TOMBSTONE_TTL` assertion | +| Repeated requests avoid overwrite writes | Present and stale-miss operation counts plus unchanged root generation/time | +| Concurrent withdrawal cannot restore IDs | Stale-snapshot and concurrent-live-update conflict tests | +| Both differing valid IDs are withdrawn | Finalization test with both rows seeded | +| Missing/unverified IDs create no root | Existing-key-only and invalid-ID tests | +| KV failure remains best-effort | Finalization response/cookie test with failing graph | +| Late partner updates cannot repopulate | Conditional, public single, and snapshot-bulk upsert rejection tests | +| Original TTL is not refreshed | First insert records `TOMBSTONE_TTL`; repetition records zero further inserts | + +## Verification Contract + +Run focused checks during implementation: + +```bash +cargo test-fastly tombstone_existing_from_snapshot +cargo test-fastly withdrawal +cargo test-fastly upsert_partner_id_if_exists_rejects_tombstone +cargo test-fastly upsert_partner_id_rejects_tombstone +cargo test-fastly snapshot_upsert_rejects_tombstone +``` + +Before committing and opening the draft PR, run: + +```bash +cargo fmt --all -- --check +cargo test-fastly +cargo test-axum +cargo test-cloudflare +cargo test-spin +cargo test --manifest-path crates/trusted-server-integration-tests/Cargo.toml --test parity +cargo clippy-fastly +cargo clippy-axum +cargo clippy-cloudflare +cargo clippy-cloudflare-wasm +cargo clippy-spin-native +cargo clippy-spin-wasm +cd crates/trusted-server-js/lib && npx vitest run +cd crates/trusted-server-js/lib && npm run format +cd docs && npm run format +cargo build --package trusted-server-adapter-fastly --release --target wasm32-wasip1 +git diff --check +``` + +## Definition of Done + +- The first live-row withdrawal writes one CAS-protected 24-hour tombstone. +- Repeated and concurrent withdrawals observing that tombstone perform no + replacement write and do not refresh its expiration. +- Repeated stale point-read misses use the completion marker and do not rewrite + the root tombstone. +- Missing IDs remain absent; unconditional overwrite remains confined to the + strongly confirmed stale-miss fallback. +- Concurrent live updates before successful withdrawal are tombstoned on retry. +- Later partner writes cannot repopulate tombstones. +- Both valid differing IDs are handled independently. +- Browser-cookie deletion remains independent of KV success. +- Focused tests, independent review, and every applicable repository gate pass. + +## Risks and Mitigations + +- **Snapshot binding:** The no-op branch must require the snapshot ID to match + the requested EC ID; a tombstone for another ID cannot suppress withdrawal. +- **Linearization:** Returning an observed tombstone linearizes withdrawal at + that observation. A later re-consent may legitimately win. +- **TTL visibility:** Record the first root and completion-marker TTLs and every + subsequent insert attempt at the wrapper boundary; stable root + generation/timestamp alone is not sufficient evidence. +- **Stale-miss completion:** Write the marker only after the root tombstone + succeeds. Encode the original tombstone's absolute validity bound in the key + and ignore the marker at or after that bound, or whenever the clock is + unusable. Marker lookup failure or clock failure must fall back to strong + root existence, not a blind write or a failed withdrawal. Explicit same-key + CAS revival and hard deletion must remove stale completion state. Revival and hard deletion fail closed before changing + the root when a bounded marker list is full, rather than making partial + cleanup progress that could leave a future live row next to an unseen marker. +- **Fallback containment:** Strongly list and validate completion markers first. + Keep unconditional overwrite behind strong root existence when no valid + marker exists so no other path can refresh a completed tombstone. +- **Conflict-test realism:** Inject actual generation changes and persisted + state, not endless synthetic precondition failures. +- **Stack dependency:** Reconcile changes if PR #885 or draft PR #900 modifies + snapshot/finalization contracts before this stack lands. diff --git a/docs/superpowers/specs/2026-03-24-ssc-technical-spec-design.md b/docs/superpowers/specs/2026-03-24-ssc-technical-spec-design.md index 3b356cb33..74ba67769 100644 --- a/docs/superpowers/specs/2026-03-24-ssc-technical-spec-design.md +++ b/docs/superpowers/specs/2026-03-24-ssc-technical-spec-design.md @@ -10,6 +10,12 @@ > Current runtime behavior interprets live consent from request cookies, headers, > geolocation, and policy defaults. `ec_identity_store` is the only KV-backed EC > lifecycle store and holds identity graph state plus withdrawal tombstones. +> +> **Completion-marker namespace (issue #881):** The same store also holds +> `__ts_ec_withdrawal_complete__:{ec_id}:{valid_until}` keys, separate from EC +> roots and excluded from hash-prefix cluster counts. Strong prefix listing +> validates the encoded expiry; an expired/malformed marker or unusable clock +> cannot suppress the existing-key-only withdrawal fallback. ---