Skip to main content

tollgate_client/
snapshot_manager.rs

1//! Background snapshot distribution: initial load, push subscription with
2//! lag recovery, periodic refresh, and revocation (review finding GL-5).
3//!
4//! The manager keeps an admission map stocked for a tracked set of
5//! principals from a [`SnapshotSource`]:
6//!
7//! - **Initial load** — every tracked principal is resolved (installed, or
8//!   negative-cached when the source confirms it unknown) before the
9//!   [`ready`](SnapshotManager::ready) watch becomes true, so readiness never
10//!   precedes admissibility (INVARIANTS.md GL-10). It returns to false when a
11//!   resolution expires or the manager exits. Source errors keep retrying;
12//!   readiness waits.
13//! - **Pushes** — subscribed updates install immediately (the map's
14//!   generation monotonicity discards stale or reordered pushes). A lagged
15//!   subscription triggers a full refetch, so a burst of missed pushes can
16//!   only delay freshness, never lose it.
17//! - **Periodic refresh** — re-fetches every tracked principal, which is
18//!   also how *revocation* propagates: a principal the source no longer
19//!   knows is removed from the map and negative-cached. The revocation
20//!   window is therefore bounded by the refresh interval (plus snapshot
21//!   `valid_until`, which fails closed on its own).
22//!
23//! Lease slots come from a shared [`SlotRegistry`], so every principal of an
24//! account observes the account's one slot — the same instance the account's
25//! `LeaseManager` refills.
26
27use std::collections::{BTreeSet, HashMap, HashSet};
28use std::sync::Arc;
29use std::sync::atomic::{AtomicU64, Ordering};
30
31use jiff::SignedDuration;
32use tokio::sync::watch;
33use tokio::task::JoinSet;
34use tracing::Instrument as _;
35
36use tollgate_admission::{
37    PublicationError, PublishableSnapshotUpdate, RefreshBatch, Refreshed, SnapshotMap, Watermark,
38    accept_positive, accept_revoked, accept_unknown,
39};
40use tollgate_core::{Generation, Locality, Principal};
41use tollgate_store::{Clock, SnapshotResolution, SnapshotSource, StoreError};
42
43pub use crate::registry::SlotRegistry;
44
45/// Which principals an instance serves.
46///
47/// A shape rather than a flag beside a list, so there is no boolean that can
48/// disagree with the data it governs (GL-16's lesson, applied to GL-48).
49#[derive(Debug, Clone, PartialEq, Eq)]
50pub enum TrackedPrincipals {
51    /// Exactly these, fixed for the process's life. Onboarding a principal
52    /// needs a restart, and readiness means every one of them is resolved.
53    Fixed(Vec<Principal>),
54    /// Every principal the source knows, re-enumerated each refresh — the
55    /// stateless topology, where any instance may serve any customer.
56    ///
57    /// Falls back to `seed` for a source that cannot enumerate, so an
58    /// embedder whose adapter predates
59    /// [`tollgate_store::SnapshotSource::principals`]
60    /// keeps the old behaviour rather than silently tracking nothing.
61    All {
62        /// Tracked until the first successful enumeration, and the permanent
63        /// set if the source cannot enumerate at all. Usually empty.
64        seed: Vec<Principal>,
65    },
66}
67
68impl TrackedPrincipals {
69    /// The set to start from, before any discovery has happened.
70    fn initial(&self) -> &[Principal] {
71        match self {
72            TrackedPrincipals::Fixed(principals) => principals,
73            TrackedPrincipals::All { seed } => seed,
74        }
75    }
76
77    fn discovers(&self) -> bool {
78        matches!(self, TrackedPrincipals::All { .. })
79    }
80}
81
82/// Which principals to serve and how to keep their snapshots fresh.
83///
84/// There are no defaults; every value is a deployment decision.
85/// [`validate`](Self::validate) runs before the task starts
86/// (INVARIANTS.md 16). `docs/SNAPSHOT_OPERATIONS.md` covers operating it.
87#[derive(Debug, Clone)]
88pub struct SnapshotManagerConfig {
89    /// The principals this instance serves — a fixed list, or everything the
90    /// source knows (GL-48).
91    ///
92    /// Also decides the readiness rule: `Fixed` needs every listed principal
93    /// resolved, `All` needs at least one. Neither a `Fixed` list nor an
94    /// `All` seed may contain duplicates, and a `Fixed` list must fit the
95    /// snapshot map's generation capacity.
96    pub principals: TrackedPrincipals,
97    /// Full refetch cadence — also the revocation propagation bound.
98    ///
99    /// Each refresh re-fetches every tracked principal (and, with `All`,
100    /// re-enumerates the catalogue), which is how a withdrawn principal is
101    /// removed. Too long lets a revoked principal keep admitting for up to
102    /// this long, bounded also by its snapshot's `valid_until`; too short
103    /// multiplies source load by the tracked set on every instance. Keep it
104    /// well below snapshot validity so a live principal is refreshed before
105    /// its snapshot expires. Must be positive.
106    pub refresh_interval: std::time::Duration,
107    /// How long a principal the source returned no row for stays negative
108    /// before the manager rechecks it.
109    ///
110    /// Short, and it covers every absence rather than only new principals: a
111    /// signup may be in flight, and a source that is rebuilding, failing over,
112    /// or serving a lagging replica reports a principal it has served for
113    /// years as absent too. This is the ceiling on how long such a gap can
114    /// deny a live customer, so it is an availability bound, not just
115    /// onboarding latency.
116    ///
117    /// Too long denies a customer whose row was briefly missing for that
118    /// long; too short rechecks every absent tracked principal more often,
119    /// one source fetch each on every instance. Must be positive.
120    pub unknown_ttl: SignedDuration,
121    /// How long a *published revocation tombstone* stays negative before the
122    /// manager rechecks it (GL-52).
123    ///
124    /// Long: coming back means an operator reinstated the account, which is
125    /// rare, and a catalogue accumulates these forever — every cancelled
126    /// customer is one, and each recheck is a fetch on every instance.
127    ///
128    /// This bounds *reinstatement*, never revocation. A live principal is
129    /// always swept, so withdrawing one still propagates within
130    /// `refresh_interval`.
131    ///
132    /// It applies only when the source *said* "revoked at generation N". An
133    /// absent row is an unknown principal and takes `unknown_ttl`, however
134    /// long this instance has served that principal: a tombstone is a durable
135    /// statement, an absence is not.
136    ///
137    /// Must be positive.
138    pub revoked_ttl: SignedDuration,
139    /// Backoff between initial-load retries while the source is down.
140    ///
141    /// Also the delay before a failed or timed-out fetch is retried. Too
142    /// short hammers a source that is already failing; too long delays
143    /// readiness after the source recovers. Must be positive.
144    pub retry_backoff: std::time::Duration,
145    /// Maximum snapshot fetches in flight during a full refresh.
146    ///
147    /// Too low makes a refresh of a large tracked set take many round trips,
148    /// stretching it toward or past `refresh_interval`; too high bursts
149    /// concurrent load at the source on every instance at once. Must be
150    /// positive.
151    pub max_concurrent_fetches: usize,
152    /// How long one source fetch may run before it is abandoned.
153    ///
154    /// The snapshot plane bounds its source calls by cancellation rather than
155    /// by a wall clock everywhere else — a slow source is answered by
156    /// readiness falling as resolutions expire, not by cutting the call off.
157    /// That is deliberate, and it is why this bound sits at the *fetch*: what
158    /// it protects is the loop's ability to come back, not the freshness of
159    /// any one principal (GL-103). A future that never resolves is never
160    /// joined, so without it one hung fetch stops the sweep from returning
161    /// and no tick, push, or control wakeup is processed again for the life
162    /// of the process.
163    ///
164    /// Set it above the slowest fetch the source legitimately makes, not
165    /// against the fast path: an abandoned fetch keeps the principal's
166    /// previous resolution and retries with backoff, so a value below real
167    /// source latency turns a slow catalogue into one that never refreshes.
168    /// It is independent of `refresh_interval`, which is a freshness cadence
169    /// rather than a statement about call latency. Must be positive.
170    pub fetch_timeout: std::time::Duration,
171    /// How long one principal enumeration may run before it is abandoned.
172    ///
173    /// Separate from [`fetch_timeout`](Self::fetch_timeout) because the two
174    /// calls have different worst cases, not because the mechanism differs:
175    /// `snapshot` returns one principal's row, `principals` returns the whole
176    /// catalogue. A single bound would have to be sized for the enumeration,
177    /// which would leave the per-fetch bound uselessly loose — and a limit is
178    /// justified against the largest legitimate input, so one value cannot
179    /// serve two inputs that differ by orders of magnitude.
180    ///
181    /// Set it above the slowest enumeration this source legitimately
182    /// performs, counted over the whole tracked set rather than a typical one.
183    /// An abandoned enumeration keeps the set it already had, so a value under
184    /// real catalogue latency freezes discovery while everything already
185    /// tracked keeps working — the failure GL-48 exists to make visible.
186    /// Must be positive.
187    pub enumeration_timeout: std::time::Duration,
188}
189
190/// Why a [`SnapshotManagerConfig`] was refused, or why [`SnapshotManager::spawn`]
191/// refused its map and slot registry: the first rule broken.
192#[derive(Debug, Clone, Copy, PartialEq, Eq)]
193pub struct SnapshotManagerConfigError(pub &'static str);
194
195impl std::fmt::Display for SnapshotManagerConfigError {
196    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
197        f.write_str(self.0)
198    }
199}
200
201impl std::error::Error for SnapshotManagerConfigError {}
202
203impl SnapshotManagerConfig {
204    /// Check the configuration on its own: every duration and TTL positive,
205    /// `max_concurrent_fetches` positive, and no duplicate in the starting
206    /// principal set. Checks that need the map, such as capacity, happen in
207    /// [`SnapshotManager::spawn`].
208    ///
209    /// # Errors
210    ///
211    /// The first rule the configuration breaks.
212    pub fn validate(&self) -> Result<(), SnapshotManagerConfigError> {
213        if self.refresh_interval.is_zero() {
214            return Err(SnapshotManagerConfigError(
215                "refresh_interval must be positive",
216            ));
217        }
218        if self.unknown_ttl <= SignedDuration::ZERO {
219            return Err(SnapshotManagerConfigError("unknown_ttl must be positive"));
220        }
221        if self.revoked_ttl <= SignedDuration::ZERO {
222            return Err(SnapshotManagerConfigError("revoked_ttl must be positive"));
223        }
224        if self.retry_backoff.is_zero() {
225            return Err(SnapshotManagerConfigError("retry_backoff must be positive"));
226        }
227        if self.fetch_timeout.is_zero() {
228            return Err(SnapshotManagerConfigError("fetch_timeout must be positive"));
229        }
230        if self.enumeration_timeout.is_zero() {
231            return Err(SnapshotManagerConfigError(
232                "enumeration_timeout must be positive",
233            ));
234        }
235        if self.max_concurrent_fetches == 0 {
236            return Err(SnapshotManagerConfigError(
237                "max_concurrent_fetches must be positive",
238            ));
239        }
240        let initial = self.principals.initial();
241        let distinct: HashSet<_> = initial.iter().copied().collect();
242        if distinct.len() != initial.len() {
243            return Err(SnapshotManagerConfigError(
244                "principals must not contain duplicates",
245            ));
246        }
247        Ok(())
248    }
249}
250
251/// What a snapshot-manager shutdown observed.
252#[derive(Debug, Clone, Copy, PartialEq, Eq)]
253pub struct SnapshotManagerReport {
254    /// The task panicked or was aborted rather than stopping on request. Its
255    /// snapshots stopped refreshing at that moment, whatever the map still
256    /// holds.
257    pub task_died: bool,
258}
259
260/// What the snapshot task has done, readable at any time.
261///
262/// [`ready`](SnapshotManager::ready) answers one bit under the configured
263/// Fixed/All rule. That is the right shape for a readiness probe and the
264/// wrong shape for diagnosis: it cannot say whether one principal is
265/// unresolved or a thousand, nor whether the source has been failing all
266/// morning (GL-4). These counters are the scrapeable half, and the
267/// `unresolved` gauge is computed from the same resolutions used by readiness.
268/// Task liveness is independent: a stopped manager is unready even if its
269/// last resolution pass had no unresolved principals.
270///
271/// Written only by the snapshot task; `Relaxed` throughout.
272#[derive(Debug)]
273pub struct SnapshotCounters {
274    refresh_attempts: AtomicU64,
275    refresh_failures: AtomicU64,
276    refresh_timeouts: AtomicU64,
277    discovery_failures: AtomicU64,
278    refused_updates: AtomicU64,
279    history_evictions: AtomicU64,
280    publication_failures: AtomicU64,
281    unresolved: AtomicU64,
282}
283
284impl SnapshotCounters {
285    /// A counter set that has recorded nothing.
286    #[must_use]
287    pub const fn new() -> Self {
288        SnapshotCounters {
289            refresh_attempts: AtomicU64::new(0),
290            refresh_failures: AtomicU64::new(0),
291            refresh_timeouts: AtomicU64::new(0),
292            discovery_failures: AtomicU64::new(0),
293            refused_updates: AtomicU64::new(0),
294            history_evictions: AtomicU64::new(0),
295            publication_failures: AtomicU64::new(0),
296            unresolved: AtomicU64::new(0),
297        }
298    }
299
300    /// One fetch of one principal, whatever its outcome.
301    fn record_attempt(&self) {
302        self.refresh_attempts.fetch_add(1, Ordering::Relaxed);
303    }
304
305    /// A fetch the source refused or could not answer. The principal keeps
306    /// its previous resolution, so this is not itself a loss of authorization
307    /// — only a loss of freshness, which `unresolved` reports once the old
308    /// resolution lapses.
309    fn record_failure(&self) {
310        self.refresh_failures.fetch_add(1, Ordering::Relaxed);
311    }
312
313    /// A fetch abandoned at `fetch_timeout` rather than answered (GL-103).
314    ///
315    /// Counted apart from `refresh_failures` for the reason the lease
316    /// manager keeps `acquire_timeouts` apart from refusals: a timeout is not
317    /// a domain answer. The source may well have resolved the principal and
318    /// simply not said so in time, so this cannot be read as "the source
319    /// could not answer" — and a rate that climbs here rather than there
320    /// points at latency, not at the catalogue.
321    fn record_timeout(&self) {
322        self.refresh_timeouts.fetch_add(1, Ordering::Relaxed);
323    }
324
325    /// Enumeration failed, so the tracked set is whatever it already was.
326    ///
327    /// Counted apart from `refresh_failures` because it is a different
328    /// failure with a different consequence: fetches failing means known
329    /// principals go stale, while enumeration failing means *new* principals
330    /// never appear at all — and that one is otherwise invisible, since
331    /// everything already tracked keeps working perfectly (GL-48).
332    fn record_discovery_failure(&self) {
333        self.discovery_failures.fetch_add(1, Ordering::Relaxed);
334    }
335
336    fn set_unresolved(&self, principals: u64) {
337        self.unresolved.store(principals, Ordering::Relaxed);
338    }
339
340    /// Read every counter. Each is read independently, so a reading taken
341    /// while the task runs is not one atomic instant across fields.
342    #[must_use]
343    pub fn snapshot(&self) -> SnapshotStats {
344        SnapshotStats {
345            refresh_attempts: self.refresh_attempts.load(Ordering::Relaxed),
346            refresh_failures: self.refresh_failures.load(Ordering::Relaxed),
347            refresh_timeouts: self.refresh_timeouts.load(Ordering::Relaxed),
348            discovery_failures: self.discovery_failures.load(Ordering::Relaxed),
349            refused_updates: self.refused_updates.load(Ordering::Relaxed),
350            history_evictions: self.history_evictions.load(Ordering::Relaxed),
351            publication_failures: self.publication_failures.load(Ordering::Relaxed),
352            unresolved: self.unresolved.load(Ordering::Relaxed),
353        }
354    }
355}
356
357impl Default for SnapshotCounters {
358    fn default() -> Self {
359        Self::new()
360    }
361}
362
363/// A reading of [`SnapshotCounters`], safe to serialise.
364#[derive(Debug, Clone, Copy, PartialEq, Eq)]
365pub struct SnapshotStats {
366    /// Fetches attempted, one per principal per pass.
367    pub refresh_attempts: u64,
368    /// Fetches the source could not answer.
369    pub refresh_failures: u64,
370    /// Fetches abandoned at `fetch_timeout` rather than answered. Apart from
371    /// `refresh_failures` because a timeout is not a domain answer: the
372    /// source may have resolved the principal and not said so in time.
373    pub refresh_timeouts: u64,
374    /// Principal enumerations the source could not answer. Nonzero means the
375    /// tracked set is frozen: everything already known keeps being refreshed,
376    /// and nothing new is ever discovered.
377    pub discovery_failures: u64,
378    /// Pushes or fetched updates refused by generation ordering. Excludes
379    /// an already installed positive at the same generation: that is an
380    /// ordinary unchanged catalogue, not a stale or revoked update.
381    pub refused_updates: u64,
382    /// Histories reclaimed under capacity pressure. Each removal also removes
383    /// visible state and requires a fresh authoritative read to recover.
384    pub history_evictions: u64,
385    /// Reservations or publications refused by the retention/fence boundary.
386    /// Distinct from source failures and ordinary generation rejection.
387    pub publication_failures: u64,
388    /// Principals with no currently valid resolution — a gauge, not a total.
389    /// Reports the last resolution pass, not task liveness. Fixed mode needs
390    /// zero unresolved; All mode needs some resolved (or an empty catalogue).
391    /// Both additionally require a live task.
392    pub unresolved: u64,
393}
394
395/// Handle to the snapshot task.
396pub struct SnapshotManager {
397    shutdown: watch::Sender<bool>,
398    ready: watch::Receiver<bool>,
399    handle: Option<tokio::task::JoinHandle<()>>,
400    counters: Arc<SnapshotCounters>,
401}
402
403impl SnapshotManager {
404    /// Validate `config` and start the snapshot task, which keeps `map`
405    /// stocked for the tracked principals and binds each account to its
406    /// slot in `slots`. Must be called within a Tokio runtime.
407    ///
408    /// Retain the handle and call [`shutdown`](Self::shutdown); dropping it
409    /// aborts the task.
410    ///
411    /// # Errors
412    ///
413    /// [`SnapshotManagerConfigError`] when `config` fails validation, when a
414    /// `Fixed` principal list exceeds the map's generation capacity, or when
415    /// the map and `slots` use different local sharding. Nothing is started.
416    pub fn spawn(
417        source: Arc<dyn SnapshotSource>,
418        map: Arc<dyn SnapshotMap>,
419        slots: Arc<SlotRegistry>,
420        clock: Arc<dyn Clock>,
421        config: SnapshotManagerConfig,
422    ) -> Result<Self, SnapshotManagerConfigError> {
423        config.validate()?;
424        if matches!(&config.principals, TrackedPrincipals::Fixed(principals) if principals.len() > map.generation_capacity().get())
425        {
426            return Err(SnapshotManagerConfigError(
427                "fixed principals exceed snapshot generation capacity",
428            ));
429        }
430        if map.local_sharding() != slots.sharding() {
431            return Err(SnapshotManagerConfigError(
432                "snapshot map and lease slots must use the same local sharding",
433            ));
434        }
435        let (shutdown, shutdown_rx) = watch::channel(false);
436        let (ready_tx, ready) = crate::task_health::TaskHealth::channel(false);
437        let principals = config.principals.initial().len();
438        let counters = Arc::new(SnapshotCounters::new());
439        let task_counters = Arc::clone(&counters);
440        let handle = tokio::spawn(
441            run(
442                source,
443                map,
444                slots,
445                clock,
446                config,
447                shutdown_rx,
448                ready_tx,
449                task_counters,
450            )
451            .instrument(tracing::info_span!("snapshot_manager", principals)),
452        );
453        Ok(SnapshotManager {
454            shutdown,
455            ready,
456            handle: Some(handle),
457            counters,
458        })
459    }
460
461    /// The snapshot task's running counters.
462    ///
463    /// Returns the shared handle for the same reason
464    /// [`LeaseManager::counters`](crate::LeaseManager::counters) does: a
465    /// service keeps it in request state while the manager itself is moved
466    /// into whatever owns shutdown, and the counters outlive the task.
467    #[must_use]
468    pub fn counters(&self) -> Arc<SnapshotCounters> {
469        Arc::clone(&self.counters)
470    }
471
472    /// True while the task is alive and its current resolutions meet the
473    /// configured Fixed/All readiness rule. The task stores false before the
474    /// channel closes on normal return, panic or cancellation; a retained
475    /// receiver's value alone cannot keep advertising an exited task as ready.
476    #[must_use]
477    pub fn ready(&self) -> watch::Receiver<bool> {
478        self.ready.clone()
479    }
480
481    /// Signal the task and wait for it to stop, reporting whether it got
482    /// there on its own.
483    ///
484    /// Its two peers already return what they know at shutdown; this one
485    /// returned nothing, so a manager that panicked mid-refresh was visible
486    /// only as snapshots quietly going stale — indistinguishable from a
487    /// control plane with nothing to say (issue GL-36).
488    pub async fn shutdown(mut self) -> SnapshotManagerReport {
489        crate::signal(&self.shutdown, true, "snapshot-manager shutdown");
490        let Some(handle) = self.handle.as_mut() else {
491            return SnapshotManagerReport { task_died: true };
492        };
493        // Keep ownership in `self` across the await: cancelling shutdown must
494        // run Drop's abort, rather than detach a taken JoinHandle.
495        match handle.await {
496            Ok(()) => SnapshotManagerReport { task_died: false },
497            Err(error) => {
498                tracing::error!(%error, "snapshot manager task died before shutdown completed");
499                SnapshotManagerReport { task_died: true }
500            }
501        }
502    }
503}
504
505impl Drop for SnapshotManager {
506    fn drop(&mut self) {
507        if let Some(handle) = self.handle.take() {
508            handle.abort();
509        }
510    }
511}
512
513#[derive(Debug, Clone, Copy)]
514enum Resolution {
515    Present {
516        deadline: jiff::Timestamp,
517        generation: Generation,
518    },
519    Negative {
520        deadline: jiff::Timestamp,
521        next_refetch: jiff::Timestamp,
522        /// The durable generation this principal carries, and why it is
523        /// durable — the same [`Watermark`] the admission map stores.
524        ///
525        /// It used to be a bare `Option<Generation>`, which could not say
526        /// whether the number came from a published revocation or merely from
527        /// the positive this negative replaced. That conflation is GL-53.
528        watermark: Option<Watermark>,
529    },
530}
531
532#[derive(Clone, Copy)]
533enum UpdateOrigin {
534    Push,
535    Refresh,
536}
537
538impl UpdateOrigin {
539    fn name(self) -> &'static str {
540        match self {
541            Self::Push => "push",
542            Self::Refresh => "refresh",
543        }
544    }
545}
546
547impl Resolution {
548    /// Test-only since GL-22: production code reads deadlines out of the
549    /// ordered indexes rather than out of the resolution, and this survives
550    /// as the naive reference's accessor — the thing the property test checks
551    /// the indexes against.
552    #[cfg(test)]
553    fn deadline(self) -> jiff::Timestamp {
554        match self {
555            Resolution::Present { deadline, .. } | Resolution::Negative { deadline, .. } => {
556                deadline
557            }
558        }
559    }
560
561    /// The durable watermark this resolution carries.
562    ///
563    /// A live snapshot's own generation is a [`Watermark::Positive`]: it orders
564    /// snapshots, but it asserts nothing about that generation being dead, so
565    /// the same generation arriving again is a re-observation rather than a
566    /// resurrection (GL-53).
567    fn watermark(self) -> Option<Watermark> {
568        match self {
569            Resolution::Present { generation, .. } => Some(Watermark::Positive(generation)),
570            Resolution::Negative { watermark, .. } => watermark,
571        }
572    }
573}
574
575/// When nothing is scheduled: re-examine in an hour rather than never.
576const IDLE_WAKEUP: std::time::Duration = std::time::Duration::from_secs(3_600);
577
578enum Publications {
579    Push(Vec<PublishableSnapshotUpdate>),
580    Refreshed(Vec<Refreshed<PublishableSnapshotUpdate>>),
581}
582
583impl Publications {
584    fn principals(&self) -> Vec<Principal> {
585        match self {
586            Self::Push(updates) => updates
587                .iter()
588                .map(PublishableSnapshotUpdate::principal)
589                .collect(),
590            Self::Refreshed(updates) => updates.iter().map(Refreshed::principal).collect(),
591        }
592    }
593}
594
595/// The running manager's sole publication boundary. Runtime membership is
596/// observed only after the map has applied its generation acceptance rule.
597struct Publication {
598    map: Arc<dyn SnapshotMap>,
599    slots: Arc<SlotRegistry>,
600}
601
602impl std::fmt::Debug for Publication {
603    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
604        f.debug_struct("Publication").finish_non_exhaustive()
605    }
606}
607
608impl Publication {
609    fn apply(&self, updates: Publications, now: jiff::Timestamp) -> Result<(), PublicationError> {
610        let principals = if self.slots.observes() {
611            updates.principals()
612        } else {
613            Vec::new()
614        };
615        match updates {
616            Publications::Push(updates) => self.map.apply_publishable_many_at(updates, now)?,
617            Publications::Refreshed(updates) => self.map.apply_refreshed_many_at(updates, now)?,
618        }
619        self.slots
620            .observe_many(principals.into_iter().map(|principal| {
621                // Observing our own publication is control-plane work, so it
622                // reads a fixed affinity rather than claiming one (GL-124).
623                (principal, self.map.get_at(&principal, Locality::OBSERVER))
624            }));
625        Ok(())
626    }
627}
628
629/// The tracked principals' resolutions, with the three ordered questions the
630/// manager asks of them answered by index rather than by rescanning.
631///
632/// Every answer here used to be a full scan of the map, and one of them ran
633/// once per completed fetch inside the refresh sweep — so a sweep of N
634/// principals did N × O(N) work (GL-22). That is invisible at today's static
635/// principal counts and becomes the binding constraint under the dynamic
636/// discovery seam `docs/DESIGN.md` defers.
637///
638/// **The map is private on purpose.** Six call sites mutate resolutions, and
639/// each must keep three indexes in step; maintained by hand that is exactly
640/// the convention-upheld-by-caller-discipline that AGENTS.md's enforcement
641/// ladder says drifts, and a drifted index here is silent — it surfaces only
642/// as readiness that is wrong. Going through methods makes a desync
643/// unrepresentable instead of merely tested.
644///
645/// Indexes hold `(deadline, principal)` so ordering is total: `Principal` is
646/// `Ord`, so two principals sharing an instant cannot collide.
647#[derive(Debug)]
648struct Resolutions {
649    publication: Option<Publication>,
650    /// The principals this instance tracks — the denominator readiness is
651    /// measured against.
652    ///
653    /// A set rather than GL-22's `usize` because discovery can change it (GL-48),
654    /// and a count alone cannot answer "is this one still ours?" when a push
655    /// arrives or an enumeration drops someone.
656    tracked: HashSet<Principal>,
657    /// Resolution and deadline state for retained histories. Expiry keeps the
658    /// watermark; a reported history reclamation removes the whole resolution
659    /// before reconstruction starts. Source-catalogue membership is separate.
660    by_principal: HashMap<Principal, Resolution>,
661    /// Live `Present` deadlines, drained as they pass.
662    present: BTreeSet<(jiff::Timestamp, Principal)>,
663    /// Live `Negative` deadlines, drained as they pass.
664    negative: BTreeSet<(jiff::Timestamp, Principal)>,
665    /// Every `Negative`'s `next_refetch`, including ones already due —
666    /// deliberately *not* drained by time, because a past refetch time is the
667    /// signal to retry now, not something to forget.
668    refetch: BTreeSet<(jiff::Timestamp, Principal)>,
669}
670
671impl Resolutions {
672    fn new(tracked: impl IntoIterator<Item = Principal>) -> Self {
673        let tracked: HashSet<Principal> = tracked.into_iter().collect();
674        Resolutions {
675            publication: None,
676            by_principal: HashMap::new(),
677            tracked,
678            present: BTreeSet::new(),
679            negative: BTreeSet::new(),
680            refetch: BTreeSet::new(),
681        }
682    }
683
684    fn publishing(
685        tracked: impl IntoIterator<Item = Principal>,
686        map: Arc<dyn SnapshotMap>,
687        slots: Arc<SlotRegistry>,
688    ) -> Self {
689        let mut resolutions = Self::new(tracked);
690        slots.retain(&resolutions.tracked);
691        resolutions.publication = Some(Publication { map, slots });
692        resolutions
693    }
694
695    fn publish(
696        &mut self,
697        updates: Publications,
698        now: jiff::Timestamp,
699        counters: &SnapshotCounters,
700    ) -> bool {
701        let principals = updates.principals();
702        let result = self
703            .publication
704            .as_ref()
705            .expect("a running manager owns publication")
706            .apply(updates, now);
707        if let Err(error) = result {
708            counters
709                .publication_failures
710                .fetch_add(1, Ordering::Relaxed);
711            tracing::warn!(%error, "snapshot publication refused; affected principals remain unresolved and retry");
712            for principal in principals {
713                self.discard(principal);
714            }
715            false
716        } else {
717            true
718        }
719    }
720
721    fn discard(&mut self, principal: Principal) {
722        if let Some(previous) = self.by_principal.remove(&principal) {
723            self.forget(principal, previous);
724        }
725    }
726
727    fn capacity(&self) -> usize {
728        self.publication
729            .as_ref()
730            .expect("a running manager owns publication")
731            .map
732            .generation_capacity()
733            .get()
734    }
735
736    fn needs_refresh(&self, principal: Principal) -> bool {
737        !self.by_principal.contains_key(&principal)
738            || self
739                .publication
740                .as_ref()
741                .expect("a running manager owns publication")
742                .map
743                .needs_refresh(principal)
744    }
745
746    fn prepare_refreshes(
747        &mut self,
748        principals: &[Principal],
749        counters: &SnapshotCounters,
750    ) -> Result<RefreshBatch, PublicationError> {
751        let publication = self
752            .publication
753            .as_ref()
754            .expect("a running manager owns publication");
755        let batch = publication.map.prepare_refreshes(principals)?;
756        if !batch.evicted.is_empty() {
757            counters
758                .history_evictions
759                .fetch_add(batch.evicted.len() as u64, Ordering::Relaxed);
760            if publication.slots.observes() {
761                publication
762                    .slots
763                    .observe_many(batch.evicted.iter().map(|&principal| (principal, None)));
764            }
765            tracing::warn!(
766                evicted = batch.evicted.len(),
767                capacity = publication.map.generation_capacity().get(),
768                "snapshot history reclaimed; evicted principals require authoritative refresh"
769            );
770            for &principal in &batch.evicted {
771                self.discard(principal);
772            }
773        }
774        Ok(batch)
775    }
776
777    fn is_tracked(&self, principal: Principal) -> bool {
778        self.tracked.contains(&principal)
779    }
780
781    /// Start tracking one principal, learned from a push rather than an
782    /// enumeration. Idempotent, and leaves an existing resolution alone.
783    fn track(&mut self, principal: Principal) {
784        self.tracked.insert(principal);
785        if let Some(publication) = &self.publication {
786            publication.slots.track(principal);
787        }
788    }
789
790    /// Adopt a discovered set.
791    ///
792    /// Untracking drops the resolution and its deadline-index entries, but
793    /// **keeps nothing behind** — which is safe only because a principal
794    /// leaves this set by disappearing from the source's catalogue, not by
795    /// being revoked. A revoked principal is still enumerated (its tombstone
796    /// is the record of the revocation), so it stays tracked and negative;
797    /// conflating the two would drop a generation watermark and let a
798    /// replayed older snapshot resurrect it (INVARIANTS.md GL-15).
799    fn retain(&mut self, discovered: HashSet<Principal>) {
800        let removed: Vec<Principal> = self
801            .tracked
802            .difference(&discovered)
803            .copied()
804            .collect::<Vec<_>>();
805        if let Some(publication) = &self.publication {
806            publication.map.remove_many(&removed);
807            if publication.slots.observes() {
808                publication
809                    .slots
810                    .observe_many(removed.iter().map(|&principal| (principal, None)));
811            }
812        }
813        for principal in removed {
814            if let Some(previous) = self.by_principal.remove(&principal) {
815                self.forget(principal, previous);
816            }
817        }
818        if let Some(publication) = &self.publication {
819            publication.slots.retain(&discovered);
820        }
821        self.tracked = discovered;
822    }
823
824    /// The durable watermark, if this principal has ever been resolved.
825    fn watermark_of(&self, principal: Principal) -> Option<Watermark> {
826        self.by_principal
827            .get(&principal)
828            .and_then(|resolution| resolution.watermark())
829    }
830
831    /// Whether an incoming positive at `generation` may replace what this
832    /// principal holds.
833    ///
834    /// Delegates to the admission layer's decision function rather than
835    /// restating the comparison. The manager gates *before* the map is ever
836    /// called, so a second copy of the rule here would decide the outcome on
837    /// its own — which is how GL-53 survived a fix to the map alone.
838    ///
839    /// History reclamation discards this resolution before refetching. A
840    /// separate visible eviction leaves history intact; probe cache membership
841    /// without recording a request-frequency hit so an equal-generation pull
842    /// repairs that miss. The manager owns publication, so its resolution
843    /// determines whether a retained visible entry is positive or negative.
844    fn accepts_positive(
845        &self,
846        principal: Principal,
847        generation: Generation,
848        origin: UpdateOrigin,
849        counters: &SnapshotCounters,
850    ) -> bool {
851        let current = self.by_principal.get(&principal).copied();
852        let (_, accepted) = accept_positive(
853            current.and_then(Resolution::watermark),
854            generation,
855            current.is_some_and(|resolution| matches!(resolution, Resolution::Present { .. }))
856                && self
857                    .publication
858                    .as_ref()
859                    .is_none_or(|publication| publication.map.contains_cached(&principal)),
860        );
861        // The shared admission rule also returns false for a visible equal
862        // generation. Classify that idempotent no-op without changing its
863        // acceptance or retry behavior, and keep a healthy refresh quiet.
864        let unchanged = matches!(current, Some(Resolution::Present { generation: held, .. })
865            if held == generation);
866        if !accepted && !unchanged {
867            self.record_refusal(principal, generation, "positive", origin, counters);
868        }
869        accepted
870    }
871
872    fn revocation(
873        &self,
874        principal: Principal,
875        generation: Generation,
876        origin: UpdateOrigin,
877        counters: &SnapshotCounters,
878    ) -> (Option<Watermark>, bool) {
879        let decision = accept_revoked(self.watermark_of(principal), generation);
880        if !decision.1 {
881            self.record_refusal(principal, generation, "revoked", origin, counters);
882        }
883        decision
884    }
885
886    /// Both ingress paths obtain their decision and its evidence together.
887    fn record_refusal(
888        &self,
889        principal: Principal,
890        generation: Generation,
891        kind: &'static str,
892        origin: UpdateOrigin,
893        counters: &SnapshotCounters,
894    ) {
895        counters.refused_updates.fetch_add(1, Ordering::Relaxed);
896        tracing::warn!(
897            %principal,
898            offered = generation.0,
899            retained = ?self.watermark_of(principal),
900            origin = origin.name(),
901            kind,
902            "snapshot update refused; previous resolution retained"
903        );
904    }
905
906    /// Record a resolution, replacing whatever this principal held.
907    ///
908    /// The single write path: index bookkeeping cannot be skipped because
909    /// there is nowhere else to write.
910    fn insert(&mut self, principal: Principal, resolution: Resolution) {
911        if let Some(previous) = self.by_principal.insert(principal, resolution) {
912            self.forget(principal, previous);
913        }
914        match resolution {
915            Resolution::Present { deadline, .. } => {
916                self.present.insert((deadline, principal));
917            }
918            Resolution::Negative {
919                deadline,
920                next_refetch,
921                ..
922            } => {
923                self.negative.insert((deadline, principal));
924                self.refetch.insert((next_refetch, principal));
925            }
926        }
927    }
928
929    /// Push a negative principal's next attempt out, after a failed refetch.
930    ///
931    /// The one in-place edit the manager makes. It moves `next_refetch` only,
932    /// never `deadline`, so it changes when the retry fires and never whether
933    /// the instance is ready.
934    fn back_off(&mut self, principal: Principal, retry_at: jiff::Timestamp) {
935        let Some(Resolution::Negative { next_refetch, .. }) = self.by_principal.get_mut(&principal)
936        else {
937            return;
938        };
939        let previous = std::mem::replace(next_refetch, retry_at);
940        self.refetch.remove(&(previous, principal));
941        self.refetch.insert((retry_at, principal));
942    }
943
944    /// Drop a superseded resolution's index entries.
945    fn forget(&mut self, principal: Principal, resolution: Resolution) {
946        match resolution {
947            Resolution::Present { deadline, .. } => {
948                self.present.remove(&(deadline, principal));
949            }
950            Resolution::Negative {
951                deadline,
952                next_refetch,
953                ..
954            } => {
955                self.negative.remove(&(deadline, principal));
956                self.refetch.remove(&(next_refetch, principal));
957            }
958        }
959    }
960
961    /// Drop deadline entries that `now` has passed, so the live sets hold
962    /// exactly the still-valid resolutions.
963    ///
964    /// Amortized O(1) per entry: each is drained at most once per insert.
965    fn expire_through(&mut self, now: jiff::Timestamp) {
966        // Pop the expired front, rather than splitting the set. `split_off`
967        // reads naturally but rebuilds the collection on *every* call even
968        // when nothing has expired, which is O(N) per query — and since the
969        // sweep queries once per completed fetch, that reintroduces exactly
970        // the quadratic this change exists to remove. Peeking the front is
971        // O(1) when nothing is due, and each entry is popped at most once per
972        // insert.
973        for set in [&mut self.present, &mut self.negative] {
974            while let Some((deadline, _)) = set.first() {
975                if *deadline > now {
976                    break;
977                }
978                set.pop_first();
979            }
980        }
981    }
982
983    /// Principals with no currently valid resolution.
984    ///
985    /// Readiness is derived from this rather than computed beside it: a
986    /// separate predicate could drift from the gauge an operator reads,
987    /// leaving `ready` false with `unresolved` at zero and nothing to explain.
988    fn unresolved(&mut self, now: jiff::Timestamp) -> usize {
989        self.expire_through(now);
990        self.tracked.len() - self.present.len() - self.negative.len()
991    }
992
993    /// How long until the earliest resolution lapses — when readiness could
994    /// next change on its own.
995    fn next_readiness_check(&mut self, now: jiff::Timestamp) -> std::time::Duration {
996        self.expire_through(now);
997        let earliest = match (self.present.first(), self.negative.first()) {
998            (Some((a, _)), Some((b, _))) => Some((*a).min(*b)),
999            (Some((only, _)), None) | (None, Some((only, _))) => Some(*only),
1000            (None, None) => None,
1001        };
1002        earliest.map_or(IDLE_WAKEUP, |deadline| until(deadline, now))
1003    }
1004
1005    /// How long until the control plane next has work: a live positive
1006    /// lapsing, or a negative becoming due for another attempt.
1007    fn next_control_wakeup(&mut self, now: jiff::Timestamp) -> std::time::Duration {
1008        self.expire_through(now);
1009        let earliest = match (self.present.first(), self.refetch.first()) {
1010            (Some((a, _)), Some((b, _))) => Some((*a).min(*b)),
1011            (Some((only, _)), None) | (None, Some((only, _))) => Some(*only),
1012            (None, None) => None,
1013        };
1014        earliest.map_or(IDLE_WAKEUP, |deadline| until(deadline, now))
1015    }
1016
1017    /// Principals the full refresh should cover: everything tracked whose
1018    /// resolution is not negative.
1019    ///
1020    /// Negatives are excluded because they already have a schedule of their
1021    /// own — [`Self::due_for_refetch`], on the TTL their kind carries. Sweeping
1022    /// them as well fetched every tombstone twice per cycle for nothing (GL-52).
1023    ///
1024    /// Principals with *no* resolution stay in: that is the initial load, and
1025    /// every principal discovery has just added.
1026    ///
1027    /// Ordered, so a sweep visits principals the same way twice running.
1028    fn due_for_sweep(&self) -> Vec<Principal> {
1029        #[allow(
1030            clippy::disallowed_methods,
1031            reason = "sorted below before it is returned, so the hash order never reaches the caller"
1032        )]
1033        let mut principals: Vec<Principal> = self
1034            .tracked
1035            .iter()
1036            .filter(|principal| {
1037                !matches!(
1038                    self.by_principal.get(principal),
1039                    Some(Resolution::Negative { .. })
1040                )
1041            })
1042            .copied()
1043            .collect();
1044        principals.sort_unstable();
1045        principals
1046    }
1047
1048    /// Every tracked principal, ordered — what the lag-recovery path sweeps.
1049    ///
1050    /// Deliberately *not* [`Self::due_for_sweep`]. A dropped push is most
1051    /// likely a reinstatement, Negative → Present, so the principals this
1052    /// path exists to repair are exactly the ones `due_for_sweep` filters
1053    /// out. Lag means local resolutions are untrustworthy; filtering by them
1054    /// would be assuming the answer (GL-52).
1055    fn all_tracked(&self) -> Vec<Principal> {
1056        #[allow(
1057            clippy::disallowed_methods,
1058            reason = "sorted on the next line before it is returned, so the hash order never reaches the caller"
1059        )]
1060        let mut principals: Vec<Principal> = self.tracked.iter().copied().collect();
1061        principals.sort_unstable();
1062        principals
1063    }
1064
1065    /// Negative principals whose next attempt is due, earliest first, at most
1066    /// `limit` of them.
1067    ///
1068    /// The cap is what keeps a due *population* from becoming one unbounded
1069    /// await. Before GL-52 every sweep re-armed each negative's deadline, so the
1070    /// refetch index rarely fired at all; now a catalogue resolved in one
1071    /// initial load shares a deadline and comes due together. The caller
1072    /// awaits this batch inline, so an uncapped set would hold the select
1073    /// loop — and the `tick` arm with it — for the whole population, which is
1074    /// the arm that carries revocation within `refresh_interval`. Chunking
1075    /// returns to the loop between waves without losing any: whatever is
1076    /// still due stays due, and the range is ordered by deadline, so the
1077    /// earliest go first and nothing starves.
1078    ///
1079    fn due_for_refetch(&self, now: jiff::Timestamp, limit: usize) -> Vec<Principal> {
1080        self.refetch
1081            .range(..(next_instant(now), Principal(0)))
1082            .take(limit)
1083            .map(|(_, principal)| *principal)
1084            .collect()
1085    }
1086}
1087
1088/// The smallest instant strictly after `now`, so a `..(bound)` range includes
1089/// everything at or before `now`. Saturates at the representable maximum.
1090fn next_instant(now: jiff::Timestamp) -> jiff::Timestamp {
1091    now.checked_add(SignedDuration::from_nanos(1))
1092        .unwrap_or(jiff::Timestamp::MAX)
1093}
1094
1095/// Time remaining until `deadline`, floored at zero for one already passed.
1096fn until(deadline: jiff::Timestamp, now: jiff::Timestamp) -> std::time::Duration {
1097    if deadline <= now {
1098        return std::time::Duration::ZERO;
1099    }
1100    u64::try_from(deadline.duration_since(now).as_nanos())
1101        .map(std::time::Duration::from_nanos)
1102        .unwrap_or(std::time::Duration::MAX)
1103}
1104
1105/// Whether the instance should be in rotation, and why the answer differs by
1106/// mode (INVARIANTS.md GL-10).
1107///
1108/// Under [`TrackedPrincipals::Fixed`] readiness is unchanged: every tracked
1109/// principal resolved. The set is small and hand-configured, so anything less
1110/// is a real gap.
1111///
1112/// Under [`TrackedPrincipals::All`] that rule inverts into a fault. The set is
1113/// the whole customer base, so it would hold an instance serving 15,999 of
1114/// 16,000 principals out of rotation for the one the source cannot answer for
1115/// — fail-closed correctness masquerading as *un*availability, which is the
1116/// same error GL-10 exists to prevent, pointed the other way. Per-principal
1117/// admissibility does not need readiness to enforce it: the map already denies
1118/// fail-closed for anything unresolved.
1119///
1120/// So under `All`, unready means **every** tracked principal is unresolved —
1121/// the point at which the instance can serve nobody and belongs out of
1122/// rotation. That still catches the case that matters (a source unreachable
1123/// long enough for snapshots to lapse), and it degenerates correctly for an
1124/// empty catalogue: an instance tracking nobody is healthy, not broken.
1125///
1126/// Both readings come from the same pass that sets the `unresolved` gauge, so
1127/// the bit and the number cannot drift.
1128fn ready_now(mode: &TrackedPrincipals, outstanding: usize, tracked: usize) -> bool {
1129    match mode {
1130        TrackedPrincipals::Fixed(_) => outstanding == 0,
1131        TrackedPrincipals::All { .. } => tracked == 0 || outstanding < tracked,
1132    }
1133}
1134
1135fn update_ready(
1136    ready: &watch::Sender<bool>,
1137    resolutions: &mut Resolutions,
1138    clock: &Arc<dyn Clock>,
1139    counters: &SnapshotCounters,
1140    mode: &TrackedPrincipals,
1141) {
1142    let outstanding = resolutions.unresolved(clock.now());
1143    counters.set_unresolved(outstanding as u64);
1144    let now_ready = ready_now(mode, outstanding, resolutions.tracked.len());
1145    // Readiness transitions are the operator-visible half of INVARIANTS GL-10;
1146    // report the edges, not every recomputation.
1147    if ready.borrow().ne(&now_ready) {
1148        tracing::info!(ready = now_ready, "snapshot readiness changed");
1149    }
1150    crate::signal(ready, now_ready, "snapshot-manager readiness");
1151}
1152
1153fn after_std(now: jiff::Timestamp, duration: std::time::Duration) -> jiff::Timestamp {
1154    let nanos = i64::try_from(duration.as_nanos()).unwrap_or(i64::MAX);
1155    now.checked_add(SignedDuration::from_nanos(nanos))
1156        .unwrap_or(jiff::Timestamp::MAX)
1157}
1158
1159/// Which negative TTL a resolution takes.
1160///
1161/// The distinction is the source's answer, not this instance's memory. An
1162/// earlier cut of GL-52 keyed the TTL on the merged generation, reasoning that
1163/// `Some(_)` meant "once served, now withdrawn". It does not: a principal the
1164/// instance has served resolves `Unknown` whenever the source's row is merely
1165/// *absent* — a store rebuilding after restart, a lagging replica, a failover
1166/// — and that inherited the hour-long reinstatement TTL. Combined with the
1167/// sweep no longer covering negatives, a transient absence blackholed a live
1168/// principal until `revoked_ttl`, with readiness still reporting healthy
1169/// because a negative counts as resolved.
1170///
1171/// A tombstone is a statement the source published; an absence is not.
1172#[derive(Clone, Copy, Debug, PartialEq, Eq)]
1173enum NegativeKind {
1174    /// The source published a revocation tombstone. Durable, and undone only
1175    /// by an operator reinstating the account, so it takes `revoked_ttl`.
1176    Revoked,
1177    /// The source has no row for this principal. May be a signup in flight or
1178    /// a source that has not finished coming up, so it takes `unknown_ttl`.
1179    Unknown,
1180}
1181
1182/// When a negative resolution should next be rechecked.
1183///
1184/// Keyed on [`NegativeKind`] — what the *source* answered — never on the
1185/// generation this instance happens to remember (GL-52).
1186fn negative_deadline(
1187    now: jiff::Timestamp,
1188    config: &SnapshotManagerConfig,
1189    kind: NegativeKind,
1190) -> jiff::Timestamp {
1191    let ttl = match kind {
1192        NegativeKind::Revoked => config.revoked_ttl,
1193        NegativeKind::Unknown => config.unknown_ttl,
1194    };
1195    now.checked_add(ttl).unwrap_or(jiff::Timestamp::MAX)
1196}
1197
1198/// Refresh a set of principals with bounded concurrency, then apply every
1199/// positive and negative result in one logical map write. A shutdown signal
1200/// aborts outstanding source futures immediately.
1201#[allow(
1202    clippy::too_many_arguments,
1203    reason = "one refresh pass's inputs, owned by the enclosing task; a parameter struct would rename that task's state"
1204)]
1205async fn refresh_all_cancellable(
1206    source: &Arc<dyn SnapshotSource>,
1207    slots: &Arc<SlotRegistry>,
1208    clock: &Arc<dyn Clock>,
1209    config: &SnapshotManagerConfig,
1210    principals: &[Principal],
1211    resolutions: &mut Resolutions,
1212    shutdown: &mut watch::Receiver<bool>,
1213    ready: &watch::Sender<bool>,
1214    counters: &SnapshotCounters,
1215) -> Option<Vec<Principal>> {
1216    let mut pending = Vec::new();
1217    for principals in principals.chunks(resolutions.capacity()) {
1218        pending.extend(
1219            refresh_chunk_cancellable(
1220                source,
1221                slots,
1222                clock,
1223                config,
1224                principals,
1225                resolutions,
1226                shutdown,
1227                ready,
1228                counters,
1229            )
1230            .await?,
1231        );
1232    }
1233    Some(pending)
1234}
1235
1236// One refresh pass's inputs, threaded rather than bundled: the enclosing task
1237// owns them all and hands each chunk the same set, so a parameter struct would
1238// be a second name for that task's state.
1239#[allow(
1240    clippy::too_many_arguments,
1241    reason = "each chunk takes the same set the enclosing pass owns"
1242)]
1243async fn refresh_chunk_cancellable(
1244    source: &Arc<dyn SnapshotSource>,
1245    slots: &Arc<SlotRegistry>,
1246    clock: &Arc<dyn Clock>,
1247    config: &SnapshotManagerConfig,
1248    principals: &[Principal],
1249    resolutions: &mut Resolutions,
1250    shutdown: &mut watch::Receiver<bool>,
1251    ready: &watch::Sender<bool>,
1252    counters: &SnapshotCounters,
1253) -> Option<Vec<Principal>> {
1254    let mut pending = match resolutions.prepare_refreshes(principals, counters) {
1255        Ok(batch) => batch.reads.into_iter(),
1256        Err(error) => {
1257            counters
1258                .publication_failures
1259                .fetch_add(1, Ordering::Relaxed);
1260            tracing::warn!(%error, "snapshot refresh reservation failed; principals retry with backoff");
1261            return Some(principals.to_vec());
1262        }
1263    };
1264    // Eviction changes readiness before source I/O, even if the source hangs.
1265    update_ready(ready, resolutions, clock, counters, &config.principals);
1266    // The inner result is the source's answer; the outer one says whether it
1267    // arrived at all. Abandoning the fetch — rather than the task holding a
1268    // future that never resolves — is what lets `tasks` empty and the sweep
1269    // return (GL-103).
1270    type Fetched = Result<Result<SnapshotResolution, StoreError>, tokio::time::error::Elapsed>;
1271    let mut tasks = JoinSet::<Refreshed<Fetched>>::new();
1272    for _ in 0..config.max_concurrent_fetches {
1273        let Some(read) = pending.next() else {
1274            break;
1275        };
1276        let source = Arc::clone(source);
1277        let bound = config.fetch_timeout;
1278        let principal = read.principal();
1279        tasks.spawn(async move {
1280            read.fetch(|| tokio::time::timeout(bound, source.snapshot(principal)))
1281                .await
1282        });
1283    }
1284
1285    let mut results = Vec::with_capacity(principals.len());
1286    while !tasks.is_empty() {
1287        // A source future may hang indefinitely. Keep readiness tied to the
1288        // actual resolution deadlines even while this sweep is in flight.
1289        let readiness_check = tokio::time::sleep(resolutions.next_readiness_check(clock.now()));
1290        tokio::pin!(readiness_check);
1291        tokio::select! {
1292            joined = tasks.join_next() => {
1293                match joined {
1294                    Some(Ok(result)) => results.push(result),
1295                    // A fetch task that panicked leaves its principal simply
1296                    // absent from the results, indistinguishable from one
1297                    // that failed cleanly — so the panic itself must be said
1298                    // out loud.
1299                    Some(Err(error)) => {
1300                        tracing::error!(%error, "snapshot fetch task died");
1301                    }
1302                    None => {}
1303                }
1304                if let Some(read) = pending.next() {
1305                    let source = Arc::clone(source);
1306                    let bound = config.fetch_timeout;
1307                    let principal = read.principal();
1308                    tasks.spawn(async move {
1309                        read.fetch(|| tokio::time::timeout(bound, source.snapshot(principal))).await
1310                    });
1311                }
1312            }
1313            changed = shutdown.changed() => {
1314                if changed.is_err() || *shutdown.borrow() {
1315                    tasks.abort_all();
1316                    return None;
1317                }
1318            }
1319            _ = &mut readiness_check => {
1320                update_ready(ready, resolutions, clock, counters, &config.principals);
1321            }
1322        }
1323    }
1324
1325    let mut updates = Vec::with_capacity(results.len());
1326    let mut completed = HashSet::with_capacity(results.len());
1327    for read in results {
1328        let principal = read.principal();
1329        let update = read.filter_map(|result| {
1330            // One attempt per principal per pass, counted whatever the outcome:
1331            // an attempt rate that has gone to zero is itself the signal that the
1332            // refresh loop has stopped.
1333            counters.record_attempt();
1334            // Unwrap the bound before the source's own answer, so every arm below
1335            // reads exactly as it did when a fetch could only succeed or fail.
1336            // Left out of `completed` like a refusal, which is what re-arms
1337            // `next_refetch` and puts this principal behind the backoff rather
1338            // than into a zero-delay refetch loop against a slow source.
1339            let Ok(result) = result else {
1340                counters.record_timeout();
1341                tracing::warn!(
1342                    %principal,
1343                    timeout_ms = config.fetch_timeout.as_millis(),
1344                    "snapshot fetch abandoned at its bound; principal keeps its \
1345                     previous resolution and retries with backoff"
1346                );
1347                return None;
1348            };
1349            match result {
1350                Ok(SnapshotResolution::Present(snapshot)) => {
1351                    if !resolutions.accepts_positive(
1352                        principal,
1353                        snapshot.generation,
1354                        UpdateOrigin::Refresh,
1355                        counters,
1356                    ) {
1357                        // Keep refusals/no-ops out of completed: it rearms a
1358                        // negative's retry backoff, never its validity deadline.
1359                        return None;
1360                    }
1361                    completed.insert(principal);
1362                    let slot = slots.slot(snapshot.account_id);
1363                    resolutions.insert(
1364                        principal,
1365                        Resolution::Present {
1366                            deadline: snapshot.valid_until,
1367                            generation: snapshot.generation,
1368                        },
1369                    );
1370                    Some(PublishableSnapshotUpdate::Present {
1371                        principal,
1372                        snapshot,
1373                        lease: slot,
1374                    })
1375                }
1376                Ok(SnapshotResolution::Revoked {
1377                    generation: incoming,
1378                }) => {
1379                    let now = clock.now();
1380                    let (watermark, accepted) = resolutions.revocation(
1381                        principal,
1382                        incoming,
1383                        UpdateOrigin::Refresh,
1384                        counters,
1385                    );
1386                    if !accepted {
1387                        // Refused answers still drive the bounded retry schedule.
1388                        return None;
1389                    }
1390                    let until = negative_deadline(now, config, NegativeKind::Revoked);
1391                    completed.insert(principal);
1392                    resolutions.insert(
1393                        principal,
1394                        Resolution::Negative {
1395                            deadline: until,
1396                            next_refetch: until,
1397                            watermark,
1398                        },
1399                    );
1400                    Some(PublishableSnapshotUpdate::Revoked {
1401                        principal,
1402                        until,
1403                        generation: incoming,
1404                    })
1405                }
1406                Ok(SnapshotResolution::Unknown) => {
1407                    completed.insert(principal);
1408                    let now = clock.now();
1409                    // The watermark passes through untouched. The source said
1410                    // nothing about any generation, so there is nothing here to
1411                    // raise or re-tag -- and re-tagging it as a revocation is what
1412                    // stranded the principal at its own generation (GL-53).
1413                    let (watermark, _) = accept_unknown(resolutions.watermark_of(principal));
1414                    let until = negative_deadline(now, config, NegativeKind::Unknown);
1415                    resolutions.insert(
1416                        principal,
1417                        Resolution::Negative {
1418                            deadline: until,
1419                            next_refetch: until,
1420                            watermark,
1421                        },
1422                    );
1423                    Some(PublishableSnapshotUpdate::Unknown { principal, until })
1424                }
1425                // The principal keeps whatever resolution it already had and
1426                // stays pending for the next sweep. Without this event a source
1427                // that is down looks exactly like one with nothing to say.
1428                Err(error) => {
1429                    counters.record_failure();
1430                    tracing::warn!(
1431                        %principal,
1432                        %error,
1433                        "snapshot fetch failed; principal keeps its previous resolution"
1434                    );
1435                    None
1436                }
1437            }
1438        });
1439        if let Some(update) = update {
1440            updates.push(update);
1441        }
1442    }
1443    if !updates.is_empty() {
1444        let published = updates.iter().map(Refreshed::principal).collect::<Vec<_>>();
1445        if !resolutions.publish(Publications::Refreshed(updates), clock.now(), counters) {
1446            for principal in published {
1447                completed.remove(&principal);
1448            }
1449        }
1450    }
1451    Some(
1452        principals
1453            .iter()
1454            .copied()
1455            .filter(|principal| !completed.contains(principal))
1456            .collect(),
1457    )
1458}
1459
1460/// Re-enumerate the tracked set.
1461///
1462/// A source that cannot enumerate (`Ok(None)`) leaves the set alone — that is
1463/// the configured-set behaviour, not an empty catalogue. A source that *fails*
1464/// also leaves it alone, but counts as a refresh failure, because an instance
1465/// quietly narrowing to nothing on a transient error would deny every request
1466/// while reporting itself perfectly healthy (GL-48).
1467/// Returns `None` when shutdown was observed while the source was
1468/// enumerating, which the caller must treat as "stop", exactly as it treats
1469/// the same answer from [`refresh_all_cancellable`].
1470///
1471/// Racing it matters because no snapshot-source call carries a wall-clock
1472/// timeout — the manager bounds them by cancellation instead — so this was
1473/// the one loop-body await a hung source could park indefinitely. Same defect
1474/// as issue GL-78 in the lease manager, found in this crate's other background
1475/// loop while fixing that one.
1476async fn discover(
1477    source: &Arc<dyn SnapshotSource>,
1478    resolutions: &mut Resolutions,
1479    counters: &SnapshotCounters,
1480    config: &SnapshotManagerConfig,
1481    shutdown: &mut watch::Receiver<bool>,
1482) -> Option<()> {
1483    // Bounded as well as raced. The shutdown race means a wedged enumeration
1484    // cannot hold shutdown open, but nothing else escaped it: during ordinary
1485    // operation the loop stayed parked here, stopped sweeping, and never
1486    // recovered — the same shape GL-103 fixed for fetches, on the one source
1487    // call it did not cover (GL-59).
1488    //
1489    // Both awaits carry the bound. The second is the spurious-wake path: a
1490    // watch change that is not a shutdown falls through to a fresh call, and
1491    // leaving that one bare would have kept the trap open for exactly the
1492    // wake-up that is not shutting anything down.
1493    let enumerate = || tokio::time::timeout(config.enumeration_timeout, source.principals());
1494    let enumerated = tokio::select! {
1495        enumerated = enumerate() => enumerated,
1496        changed = shutdown.changed() => {
1497            if changed.is_err() || *shutdown.borrow() {
1498                tracing::debug!("shutdown observed during principal enumeration");
1499                return None;
1500            }
1501            enumerate().await
1502        }
1503    };
1504    let Ok(enumerated) = enumerated else {
1505        // Counted as a discovery failure, because the consequence is the same
1506        // one GL-48 names: the tracked set is frozen, everything already known
1507        // keeps being refreshed, and nothing new is ever discovered. The event
1508        // says which of the two it was.
1509        counters.record_discovery_failure();
1510        tracing::warn!(
1511            timeout_ms = config.enumeration_timeout.as_millis(),
1512            "principal enumeration abandoned at its bound; keeping the current set"
1513        );
1514        return Some(());
1515    };
1516    match enumerated {
1517        Ok(Some(discovered)) => resolutions.retain(discovered.into_iter().collect()),
1518        // Cannot enumerate: keep the configured set. Deliberately not the same
1519        // as an empty catalogue, which would mean "forget everyone".
1520        Ok(None) => {}
1521        Err(error) => {
1522            counters.record_discovery_failure();
1523            tracing::warn!(%error, "principal enumeration failed; keeping the current set");
1524        }
1525    }
1526    Some(())
1527}
1528
1529// The task entry point: every argument is a handle the loop must own for its
1530// lifetime, and they are already assembled once by `spawn`.
1531#[allow(
1532    clippy::too_many_arguments,
1533    reason = "task entry point: every argument is a handle the loop owns for its lifetime, assembled once by spawn"
1534)]
1535async fn run(
1536    source: Arc<dyn SnapshotSource>,
1537    map: Arc<dyn SnapshotMap>,
1538    slots: Arc<SlotRegistry>,
1539    clock: Arc<dyn Clock>,
1540    config: SnapshotManagerConfig,
1541    mut shutdown: watch::Receiver<bool>,
1542    readiness: crate::task_health::TaskHealth,
1543    counters: Arc<SnapshotCounters>,
1544) {
1545    let ready = readiness.sender();
1546    let mut updates = source.subscribe();
1547    let mut resolutions = Resolutions::publishing(
1548        config.principals.initial().iter().copied(),
1549        Arc::clone(&map),
1550        Arc::clone(&slots),
1551    );
1552
1553    // Discover before the initial load, so a stateless instance starts from
1554    // the real set rather than its (usually empty) seed. Enumeration failure
1555    // is not fatal: the retry loop below covers it, and the seed keeps the
1556    // instance serving whatever it was told about meanwhile.
1557    if config.principals.discovers()
1558        && discover(&source, &mut resolutions, &counters, &config, &mut shutdown)
1559            .await
1560            .is_none()
1561    {
1562        return;
1563    }
1564
1565    // Initial load: retry until every tracked principal is resolved, then
1566    // report ready. The map denies (fail closed) for anything unresolved in
1567    // the meantime.
1568    #[allow(
1569        clippy::disallowed_methods,
1570        reason = "a work queue drained to empty, not an output: every tracked principal is fetched, and the pass reports the same result whichever order they were fetched in"
1571    )]
1572    let mut pending: Vec<Principal> = resolutions.tracked.iter().copied().collect();
1573    while !pending.is_empty() {
1574        if *shutdown.borrow() {
1575            return;
1576        }
1577        let Some(failed) = refresh_all_cancellable(
1578            &source,
1579            &slots,
1580            &clock,
1581            &config,
1582            &pending,
1583            &mut resolutions,
1584            &mut shutdown,
1585            ready,
1586            &counters,
1587        )
1588        .await
1589        else {
1590            return;
1591        };
1592        pending = failed;
1593        update_ready(
1594            ready,
1595            &mut resolutions,
1596            &clock,
1597            &counters,
1598            &config.principals,
1599        );
1600        if !pending.is_empty() {
1601            tokio::select! {
1602                _ = tokio::time::sleep(config.retry_backoff) => {}
1603                changed = shutdown.changed() => {
1604                    if changed.is_err() || *shutdown.borrow() {
1605                        return;
1606                    }
1607                }
1608            }
1609        }
1610    }
1611    update_ready(
1612        ready,
1613        &mut resolutions,
1614        &clock,
1615        &counters,
1616        &config.principals,
1617    );
1618
1619    let mut tick = tokio::time::interval(config.refresh_interval);
1620    tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
1621    tick.reset(); // the initial load already counts as a refresh
1622    let mut updates_closed = false;
1623    loop {
1624        if *shutdown.borrow() {
1625            return;
1626        }
1627        let control_wakeup = tokio::time::sleep(resolutions.next_control_wakeup(clock.now()));
1628        tokio::pin!(control_wakeup);
1629        tokio::select! {
1630            push = updates.recv(), if !updates_closed => match push {
1631                Ok(push) => {
1632                    // Under discovery every push is ours: a push for a
1633                    // principal we have not enumerated yet *is* the discovery,
1634                    // and dropping it was how a newly provisioned customer
1635                    // stayed invisible until a restart (GL-48).
1636                    if config.principals.discovers() {
1637                        resolutions.track(push.principal);
1638                    }
1639                    if resolutions.is_tracked(push.principal) {
1640                        if resolutions.needs_refresh(push.principal) {
1641                            if refresh_all_cancellable(&source, &slots, &clock, &config, &[push.principal],
1642                                &mut resolutions, &mut shutdown, ready, &counters).await.is_none() { return; }
1643                            update_ready(ready, &mut resolutions, &clock, &counters, &config.principals);
1644                            continue;
1645                        }
1646
1647                        match push.resolution {
1648                            SnapshotResolution::Present(snapshot) => {
1649                                if !resolutions
1650                                    .accepts_positive(push.principal, snapshot.generation, UpdateOrigin::Push, &counters)
1651                                {
1652                                    continue;
1653                                }
1654                                let slot = slots.slot(snapshot.account_id);
1655                                resolutions.insert(
1656                                    push.principal,
1657                                    Resolution::Present {
1658                                        deadline: snapshot.valid_until,
1659                                        generation: snapshot.generation,
1660                                    },
1661                                );
1662                                resolutions.publish(
1663                                    Publications::Push(vec![PublishableSnapshotUpdate::Present {
1664                                        principal: push.principal,
1665                                        snapshot,
1666                                        lease: slot,
1667                                    }]),
1668                                    clock.now(),
1669                                    &counters,
1670                                );
1671                            }
1672                            resolution @ (SnapshotResolution::Revoked { .. }
1673                            | SnapshotResolution::Unknown) => {
1674                                let current = resolutions.watermark_of(push.principal);
1675                                let (watermark, accepted, kind, update) = match resolution {
1676                                    SnapshotResolution::Revoked { generation } => {
1677                                        let (watermark, accepted) =
1678                                            resolutions.revocation(push.principal, generation, UpdateOrigin::Push, &counters);
1679                                        (watermark, accepted, NegativeKind::Revoked, Some(generation))
1680                                    }
1681                                    SnapshotResolution::Unknown => {
1682                                        let (watermark, accepted) = accept_unknown(current);
1683                                        (watermark, accepted, NegativeKind::Unknown, None)
1684                                    }
1685                                    SnapshotResolution::Present(_) => unreachable!(),
1686                                };
1687                                if !accepted {
1688                                    continue;
1689                                }
1690                                let now = clock.now();
1691                                let until = negative_deadline(now, &config, kind);
1692                                resolutions.insert(
1693                                    push.principal,
1694                                    Resolution::Negative {
1695                                        deadline: until,
1696                                        next_refetch: until,
1697                                        watermark,
1698                                    },
1699                                );
1700                                let update = match update {
1701                                    Some(generation) => PublishableSnapshotUpdate::Revoked {
1702                                        principal: push.principal,
1703                                        until,
1704                                        generation,
1705                                    },
1706                                    None => PublishableSnapshotUpdate::Unknown {
1707                                        principal: push.principal,
1708                                        until,
1709                                    },
1710                                };
1711                                resolutions.publish(Publications::Push(vec![update]), now, &counters);
1712                            }
1713                        }
1714                        update_ready(ready, &mut resolutions, &clock, &counters, &config.principals);
1715                    }
1716                }
1717                // Lagged: missed pushes — refetch everything rather than
1718                // guess what was dropped. `all_tracked`, not `due_for_sweep`:
1719                // a dropped reinstatement leaves a stale Negative behind, and
1720                // that is the one thing a filtered sweep would never revisit.
1721                Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => {
1722                    if refresh_all_cancellable(
1723                        &source, &slots, &clock, &config, &resolutions.all_tracked(),
1724                        &mut resolutions, &mut shutdown,
1725                        ready,
1726                        &counters,
1727                    ).await.is_none() {
1728                        return;
1729                    }
1730                    update_ready(ready, &mut resolutions, &clock, &counters, &config.principals);
1731                }
1732                // Push stream gone (e.g. HTTP transport): periodic refresh
1733                // remains the freshness path.
1734                Err(tokio::sync::broadcast::error::RecvError::Closed) => {
1735                    updates_closed = true;
1736                }
1737            },
1738            _ = tick.tick() => {
1739                // Re-enumerate first: the sweep should cover principals added
1740                // since the last one, and stop covering any the source has
1741                // dropped. Over HTTP this is the only propagation path there
1742                // is, because `subscribe` is a closed channel.
1743                if config.principals.discovers()
1744                    && discover(&source, &mut resolutions, &counters, &config, &mut shutdown)
1745                        .await
1746                        .is_none()
1747                {
1748                    return;
1749                }
1750                if refresh_all_cancellable(
1751                    &source, &slots, &clock, &config, &resolutions.due_for_sweep(),
1752                    &mut resolutions, &mut shutdown,
1753                    ready,
1754                    &counters,
1755                ).await.is_none() {
1756                    return;
1757                }
1758                update_ready(ready, &mut resolutions, &clock, &counters, &config.principals);
1759            }
1760            _ = &mut control_wakeup => {
1761                let now = clock.now();
1762                update_ready(ready, &mut resolutions, &clock, &counters, &config.principals);
1763                let due = resolutions.due_for_refetch(now, config.max_concurrent_fetches);
1764                if !due.is_empty() {
1765                    let Some(failed) = refresh_all_cancellable(
1766                        &source, &slots, &clock, &config, &due,
1767                        &mut resolutions, &mut shutdown, ready, &counters,
1768                    ).await else {
1769                        return;
1770                    };
1771                    let retry_at = after_std(clock.now(), config.retry_backoff);
1772                    for principal in failed {
1773                        resolutions.back_off(principal, retry_at);
1774                    }
1775                    update_ready(ready, &mut resolutions, &clock, &counters, &config.principals);
1776                }
1777            }
1778            changed = shutdown.changed() => {
1779                if changed.is_err() || *shutdown.borrow() {
1780                    return;
1781                }
1782            }
1783        }
1784    }
1785}
1786
1787#[cfg(test)]
1788mod tests {
1789    use proptest::prelude::*;
1790
1791    use super::*;
1792
1793    fn t(seconds: i64) -> jiff::Timestamp {
1794        jiff::Timestamp::from_second(seconds).unwrap()
1795    }
1796
1797    /// Publishing and observing is control-plane work, and it must not spend
1798    /// a request-serving thread's affinity to do it (GL-124).
1799    ///
1800    /// `Locality` is handed out on a thread's *first* access from one
1801    /// process-global counter that never recycles, and every number it hands
1802    /// out pushes the live request-serving threads one step closer to sharing
1803    /// a shard. `apply` used to observe its own publication through
1804    /// `SnapshotMap::get`, which resolves `Locality::current()` — so a control
1805    /// plane on a thread that serves no requests charged the request path for
1806    /// its own bookkeeping.
1807    ///
1808    /// Asserted against the call the map receives rather than against the
1809    /// counter. The counter is process-global and every other test in this
1810    /// binary spends from it concurrently, so a before/after delta is a race;
1811    /// which affinity the observation *asks for* is a fact about this code.
1812    #[test]
1813    fn observing_a_publication_does_not_consume_an_affinity() {
1814        use std::sync::Mutex;
1815        use tollgate_admission::{
1816            AdmissionCounters, ArcSwapSnapshotMap, LeaseSlot, MapEntry, PublicationError,
1817        };
1818        use tollgate_core::{
1819            AccountId, AccountSnapshot, AccountStatus, CostTable, CostUnits, Generation,
1820            LocalSharding, PermissionBits, PublishableSnapshot, ResolvedLimits,
1821        };
1822
1823        /// Records how each read asked, so the distinction this test is about
1824        /// is visible: `None` is a `get`, which resolves — and therefore may
1825        /// assign — the calling thread's affinity.
1826        struct RecordingMap {
1827            inner: ArcSwapSnapshotMap,
1828            reads: Mutex<Vec<Option<Locality>>>,
1829        }
1830
1831        impl SnapshotMap for RecordingMap {
1832            fn get(&self, principal: &Principal) -> Option<MapEntry> {
1833                self.reads
1834                    .lock()
1835                    .expect("recording map poisoned")
1836                    .push(None);
1837                self.inner.get(principal)
1838            }
1839            fn get_at(&self, principal: &Principal, locality: Locality) -> Option<MapEntry> {
1840                self.reads
1841                    .lock()
1842                    .expect("recording map poisoned")
1843                    .push(Some(locality));
1844                self.inner.get_at(principal, locality)
1845            }
1846            fn counters(&self) -> &Arc<AdmissionCounters> {
1847                self.inner.counters()
1848            }
1849            fn install(
1850                &self,
1851                principal: Principal,
1852                snapshot: Arc<AccountSnapshot>,
1853                lease: Arc<LeaseSlot>,
1854            ) -> Result<(), PublicationError> {
1855                self.inner.install(principal, snapshot, lease)
1856            }
1857            fn install_revoked(
1858                &self,
1859                principal: Principal,
1860                until: jiff::Timestamp,
1861                generation: Generation,
1862            ) -> Result<(), PublicationError> {
1863                self.inner.install_revoked(principal, until, generation)
1864            }
1865            fn install_unknown(
1866                &self,
1867                principal: Principal,
1868                until: jiff::Timestamp,
1869            ) -> Result<(), PublicationError> {
1870                self.inner.install_unknown(principal, until)
1871            }
1872            fn remove(&self, principal: &Principal) {
1873                self.inner.remove(principal);
1874            }
1875        }
1876
1877        let map = Arc::new(RecordingMap {
1878            inner: ArcSwapSnapshotMap::new(),
1879            reads: Mutex::new(Vec::new()),
1880        });
1881        let (slots, _changes) = SlotRegistry::observed(LocalSharding::SINGLE);
1882        assert!(
1883            slots.observes(),
1884            "an observing registry is what reaches the read under test"
1885        );
1886        let publication = Publication {
1887            map: Arc::clone(&map) as Arc<dyn SnapshotMap>,
1888            slots,
1889        };
1890        let principal = Principal(7);
1891        let snapshot = PublishableSnapshot::try_new(Arc::new(
1892            AccountSnapshot::builder(
1893                AccountId(1),
1894                Generation(1),
1895                AccountStatus::Active,
1896                t(3_600),
1897                PermissionBits::bit(0),
1898                ResolvedLimits::new(16),
1899                Arc::new(CostTable::builder(CostUnits(1), CostUnits(1)).build()),
1900            )
1901            .build(),
1902        ))
1903        .expect("test snapshot limits are valid");
1904
1905        publication
1906            .apply(
1907                Publications::Push(vec![PublishableSnapshotUpdate::Present {
1908                    principal,
1909                    snapshot,
1910                    lease: LeaseSlot::for_account(AccountId(1)),
1911                }]),
1912                t(0),
1913            )
1914            .expect("the fixture publishes a well-formed snapshot");
1915
1916        let reads = map.reads.lock().expect("recording map poisoned").clone();
1917        assert_eq!(
1918            reads,
1919            vec![Some(Locality::OBSERVER)],
1920            "the publication must observe itself at the observer affinity, never \
1921             through a read that resolves the calling thread's"
1922        );
1923        assert!(
1924            map.get_at(&principal, Locality::OBSERVER).is_some(),
1925            "and the publication still landed"
1926        );
1927    }
1928
1929    #[test]
1930    fn history_reclamation_bounds_resolution_and_deadline_indexes() {
1931        use std::num::NonZeroUsize;
1932        use tollgate_admission::ArcSwapSnapshotMap;
1933        use tollgate_core::LocalSharding;
1934        let map = Arc::new(ArcSwapSnapshotMap::with_capacities(
1935            LocalSharding::SINGLE,
1936            7,
1937            NonZeroUsize::new(7).unwrap(),
1938        ));
1939        let mut resolutions =
1940            Resolutions::publishing((0..1000).map(Principal), map.clone(), SlotRegistry::new());
1941        let counters = SnapshotCounters::new();
1942        for principal in (0..1000).map(Principal) {
1943            let read = resolutions
1944                .prepare_refreshes(&[principal], &counters)
1945                .unwrap()
1946                .reads
1947                .pop()
1948                .unwrap();
1949            resolutions.insert(
1950                principal,
1951                Resolution::Negative {
1952                    deadline: t(10),
1953                    next_refetch: t(10),
1954                    watermark: None,
1955                },
1956            );
1957            assert!(resolutions.publish(
1958                Publications::Refreshed(vec![read.read(|| PublishableSnapshotUpdate::Unknown {
1959                    principal,
1960                    until: t(10)
1961                })]),
1962                t(0),
1963                &counters
1964            ));
1965            let retained = ((principal.0 + 1) as usize).min(7);
1966            assert_eq!(resolutions.by_principal.len(), retained);
1967            assert_eq!(resolutions.negative.len(), retained);
1968            assert_eq!(resolutions.refetch.len(), retained);
1969            assert!(resolutions.present.is_empty());
1970            assert_eq!(map.history_stats().unwrap().retained, retained);
1971        }
1972        assert_eq!(resolutions.unresolved(t(0)), 993);
1973        assert_eq!(counters.snapshot().history_evictions, 993);
1974        assert_eq!(counters.snapshot().publication_failures, 0);
1975    }
1976
1977    #[test]
1978    fn a_superseded_publication_cannot_claim_a_resolved_deadline() {
1979        use std::num::NonZeroUsize;
1980        use tollgate_admission::ArcSwapSnapshotMap;
1981        use tollgate_core::LocalSharding;
1982        let map = Arc::new(ArcSwapSnapshotMap::with_capacities(
1983            LocalSharding::SINGLE,
1984            1,
1985            NonZeroUsize::new(1).unwrap(),
1986        ));
1987        let mut resolutions = Resolutions::publishing(
1988            [Principal(1), Principal(2)],
1989            map.clone(),
1990            SlotRegistry::new(),
1991        );
1992        let counters = SnapshotCounters::new();
1993        let stale = resolutions
1994            .prepare_refreshes(&[Principal(1)], &counters)
1995            .unwrap()
1996            .reads
1997            .pop()
1998            .unwrap();
1999        // Another control-plane reservation overtakes this source response.
2000        resolutions
2001            .prepare_refreshes(&[Principal(2)], &counters)
2002            .unwrap();
2003        resolutions.insert(
2004            Principal(1),
2005            Resolution::Negative {
2006                deadline: t(10),
2007                next_refetch: t(10),
2008                watermark: None,
2009            },
2010        );
2011        assert!(!resolutions.publish(
2012            Publications::Refreshed(vec![stale.read(|| PublishableSnapshotUpdate::Unknown {
2013                principal: Principal(1),
2014                until: t(10)
2015            })]),
2016            t(0),
2017            &counters
2018        ));
2019        assert_eq!(resolutions.unresolved(t(0)), 2);
2020        assert!(resolutions.by_principal.is_empty());
2021        assert!(resolutions.negative.is_empty());
2022        assert!(resolutions.refetch.is_empty());
2023        assert_eq!(counters.snapshot().publication_failures, 1);
2024        assert_eq!(counters.snapshot().refresh_failures, 0);
2025        assert!(map.get(&Principal(1)).is_none());
2026    }
2027
2028    #[test]
2029    fn control_wakeup_tracks_future_positive_deadline_once() {
2030        let principal = Principal(1);
2031        let mut resolutions = Resolutions::new([Principal(1)]);
2032        resolutions.insert(
2033            principal,
2034            Resolution::Present {
2035                deadline: t(15),
2036                generation: Generation(1),
2037            },
2038        );
2039        assert_eq!(
2040            resolutions.next_control_wakeup(t(10)),
2041            std::time::Duration::from_secs(5)
2042        );
2043
2044        resolutions.insert(
2045            principal,
2046            Resolution::Present {
2047                deadline: t(10),
2048                generation: Generation(1),
2049            },
2050        );
2051        assert_eq!(
2052            resolutions.next_control_wakeup(t(10)),
2053            std::time::Duration::from_secs(3_600)
2054        );
2055    }
2056
2057    #[test]
2058    fn control_wakeup_tracks_negative_refetch_deadline() {
2059        let principal = Principal(1);
2060        let mut resolutions = Resolutions::new([Principal(1)]);
2061        resolutions.insert(
2062            principal,
2063            Resolution::Negative {
2064                deadline: t(15),
2065                next_refetch: t(15),
2066                watermark: None,
2067            },
2068        );
2069        assert_eq!(
2070            resolutions.next_control_wakeup(t(10)),
2071            std::time::Duration::from_secs(5)
2072        );
2073        assert_eq!(
2074            resolutions.next_control_wakeup(t(15)),
2075            std::time::Duration::ZERO
2076        );
2077    }
2078
2079    // ---- GL-22: the indexed set must answer what the scans answered --------
2080
2081    /// The pre-GL-22 implementations, kept verbatim as the reference the
2082    /// indexed set is checked against. If these and `Resolutions` ever
2083    /// disagree, the refactor changed behaviour — which is the one thing it
2084    /// must not do.
2085    mod naive {
2086        use super::*;
2087
2088        pub(super) fn unresolved(
2089            principals: &[Principal],
2090            resolutions: &HashMap<Principal, Resolution>,
2091            now: jiff::Timestamp,
2092        ) -> usize {
2093            principals
2094                .iter()
2095                .filter(|principal| {
2096                    !resolutions
2097                        .get(principal)
2098                        .is_some_and(|resolution| now < resolution.deadline())
2099                })
2100                .count()
2101        }
2102
2103        pub(super) fn next_readiness_check(
2104            resolutions: &HashMap<Principal, Resolution>,
2105            now: jiff::Timestamp,
2106        ) -> std::time::Duration {
2107            #[allow(
2108                clippy::disallowed_methods,
2109                reason = "reduces deadlines to a minimum; a minimum does not depend on the order it is taken in"
2110            )]
2111            resolutions
2112                .values()
2113                .map(|resolution| resolution.deadline())
2114                .filter(|deadline| *deadline > now)
2115                .map(|deadline| deadline.duration_since(now).as_nanos())
2116                .min()
2117                .and_then(|nanos| u64::try_from(nanos).ok())
2118                .map(std::time::Duration::from_nanos)
2119                .unwrap_or_else(|| std::time::Duration::from_secs(3_600))
2120        }
2121
2122        pub(super) fn next_control_wakeup(
2123            resolutions: &HashMap<Principal, Resolution>,
2124            now: jiff::Timestamp,
2125        ) -> std::time::Duration {
2126            #[allow(
2127                clippy::disallowed_methods,
2128                reason = "reduces deadlines to a minimum; a minimum does not depend on the order it is taken in"
2129            )]
2130            resolutions
2131                .values()
2132                .filter_map(|resolution| match resolution {
2133                    Resolution::Present { deadline, .. } if *deadline > now => Some(*deadline),
2134                    Resolution::Present { .. } => None,
2135                    Resolution::Negative { next_refetch, .. } => Some(*next_refetch),
2136                })
2137                .map(|deadline| {
2138                    if deadline <= now {
2139                        std::time::Duration::ZERO
2140                    } else {
2141                        let nanos = deadline.duration_since(now).as_nanos();
2142                        u64::try_from(nanos)
2143                            .map(std::time::Duration::from_nanos)
2144                            .unwrap_or(std::time::Duration::MAX)
2145                    }
2146                })
2147                .min()
2148                .unwrap_or_else(|| std::time::Duration::from_secs(3_600))
2149        }
2150
2151        pub(super) fn due_for_refetch(
2152            resolutions: &HashMap<Principal, Resolution>,
2153            now: jiff::Timestamp,
2154            limit: usize,
2155        ) -> Vec<Principal> {
2156            #[allow(
2157                clippy::disallowed_methods,
2158                reason = "sorted below before the limit is applied, so which principals a bounded refetch takes is decided by deadline"
2159            )]
2160            let mut due: Vec<_> = resolutions
2161                .iter()
2162                .filter_map(|(principal, resolution)| match resolution {
2163                    Resolution::Negative {
2164                        next_refetch,
2165                        deadline,
2166                        ..
2167                    } if *next_refetch <= now => Some((*next_refetch, *deadline, *principal)),
2168                    _ => None,
2169                })
2170                .collect();
2171            // The index is keyed on (next_refetch, principal), so the naive
2172            // reference has to take the earliest by that same order before
2173            // truncating — otherwise the two disagree on *which* are dropped.
2174            due.sort_unstable_by_key(|(next_refetch, _, principal)| (*next_refetch, *principal));
2175            due.into_iter()
2176                .take(limit)
2177                .map(|(_, _, principal)| principal)
2178                .collect()
2179        }
2180    }
2181
2182    /// One step of the manager's life, as the property test drives it.
2183    #[derive(Debug, Clone, Copy)]
2184    enum Step {
2185        Present { principal: u8, deadline: i64 },
2186        Negative { principal: u8, deadline: i64 },
2187        BackOff { principal: u8, retry_in: i64 },
2188        Advance { seconds: i64 },
2189    }
2190
2191    const TRACKED: usize = 6;
2192
2193    fn step() -> impl Strategy<Value = Step> {
2194        prop_oneof![
2195            (0..TRACKED as u8, 0i64..400).prop_map(|(principal, deadline)| Step::Present {
2196                principal,
2197                deadline
2198            }),
2199            (0..TRACKED as u8, 0i64..400).prop_map(|(principal, deadline)| Step::Negative {
2200                principal,
2201                deadline
2202            }),
2203            (0..TRACKED as u8, 0i64..200).prop_map(|(principal, retry_in)| Step::BackOff {
2204                principal,
2205                retry_in
2206            }),
2207            (0i64..50).prop_map(|seconds| Step::Advance { seconds }),
2208        ]
2209    }
2210
2211    proptest! {
2212        /// Every question the indexed set answers must match the scan it
2213        /// replaced, after any sequence of resolutions, backoffs and clock
2214        /// advances (GL-22). Time only moves forward, as it does in the
2215        /// manager.
2216        #[test]
2217        fn indexed_resolutions_answer_exactly_what_scanning_answered(
2218            steps in proptest::collection::vec(step(), 1..60),
2219        ) {
2220            let principals: Vec<Principal> =
2221                (0..TRACKED as u128).map(Principal).collect();
2222            let mut indexed = Resolutions::new(principals.iter().copied());
2223            let mut reference: HashMap<Principal, Resolution> = HashMap::new();
2224            let mut now = t(0);
2225
2226            for step in steps {
2227                match step {
2228                    Step::Present { principal, deadline } => {
2229                        let principal = Principal(u128::from(principal));
2230                        let resolution = Resolution::Present {
2231                            deadline: t(deadline),
2232                            generation: Generation(1),
2233                        };
2234                        indexed.insert(principal, resolution);
2235                        reference.insert(principal, resolution);
2236                    }
2237                    Step::Negative { principal, deadline } => {
2238                        let principal = Principal(u128::from(principal));
2239                        let resolution = Resolution::Negative {
2240                            deadline: t(deadline),
2241                            next_refetch: t(deadline),
2242                            watermark: None,
2243                        };
2244                        indexed.insert(principal, resolution);
2245                        reference.insert(principal, resolution);
2246                    }
2247                    Step::BackOff { principal, retry_in } => {
2248                        let principal = Principal(u128::from(principal));
2249                        let retry_at = t(now.as_second() + retry_in);
2250                        indexed.back_off(principal, retry_at);
2251                        if let Some(Resolution::Negative { next_refetch, .. }) =
2252                            reference.get_mut(&principal)
2253                        {
2254                            *next_refetch = retry_at;
2255                        }
2256                    }
2257                    Step::Advance { seconds } => {
2258                        now = t(now.as_second() + seconds);
2259                    }
2260                }
2261
2262                prop_assert_eq!(
2263                    indexed.unresolved(now),
2264                    naive::unresolved(&principals, &reference, now),
2265                    "unresolved disagreed at {:?}", now
2266                );
2267                prop_assert_eq!(
2268                    indexed.next_readiness_check(now),
2269                    naive::next_readiness_check(&reference, now),
2270                    "next_readiness_check disagreed at {:?}", now
2271                );
2272                prop_assert_eq!(
2273                    indexed.next_control_wakeup(now),
2274                    naive::next_control_wakeup(&reference, now),
2275                    "next_control_wakeup disagreed at {:?}", now
2276                );
2277                for limit in [1usize, 3, usize::MAX] {
2278                    let mut due = indexed.due_for_refetch(now, limit);
2279                    let mut expected = naive::due_for_refetch(&reference, now, limit);
2280                    due.sort_unstable();
2281                    expected.sort_unstable();
2282                    prop_assert_eq!(
2283                        due,
2284                        expected,
2285                        "due_for_refetch disagreed at {:?} under limit {}", now, limit
2286                    );
2287                }
2288            }
2289        }
2290    }
2291
2292    /// A resolution inserted already expired counts as unresolved at once —
2293    /// it never enters the live index, so nothing has to expire it later.
2294    #[test]
2295    fn a_resolution_inserted_expired_is_never_live() {
2296        let mut resolutions = Resolutions::new([Principal(0)]);
2297        resolutions.insert(
2298            Principal(0),
2299            Resolution::Present {
2300                deadline: t(5),
2301                generation: Generation(1),
2302            },
2303        );
2304        assert_eq!(resolutions.unresolved(t(10)), 1);
2305        // And asking again does not double-count or underflow.
2306        assert_eq!(resolutions.unresolved(t(10)), 1);
2307        assert_eq!(resolutions.unresolved(t(20)), 1);
2308    }
2309
2310    /// Re-resolving a principal replaces its index entry rather than adding
2311    /// one: the live count is per principal, not per insert.
2312    #[test]
2313    fn re_resolving_a_principal_does_not_double_count_it() {
2314        let mut resolutions = Resolutions::new([Principal(0), Principal(1)]);
2315        for deadline in [t(50), t(60), t(60), t(70)] {
2316            resolutions.insert(
2317                Principal(0),
2318                Resolution::Present {
2319                    deadline,
2320                    generation: Generation(1),
2321                },
2322            );
2323        }
2324        assert_eq!(
2325            resolutions.unresolved(t(10)),
2326            1,
2327            "one principal resolved, one still outstanding"
2328        );
2329        assert_eq!(resolutions.next_readiness_check(t(10)), secs(60));
2330    }
2331
2332    /// A negative superseded by a positive must leave neither a stale
2333    /// deadline nor a stale refetch behind.
2334    #[test]
2335    fn a_positive_replacing_a_negative_clears_both_of_its_indexes() {
2336        let mut resolutions = Resolutions::new([Principal(0)]);
2337        resolutions.insert(
2338            Principal(0),
2339            Resolution::Negative {
2340                deadline: t(30),
2341                next_refetch: t(30),
2342                watermark: None,
2343            },
2344        );
2345        resolutions.insert(
2346            Principal(0),
2347            Resolution::Present {
2348                deadline: t(90),
2349                generation: Generation(2),
2350            },
2351        );
2352
2353        assert_eq!(resolutions.unresolved(t(40)), 0, "the positive is live");
2354        assert!(
2355            resolutions.due_for_refetch(t(40), usize::MAX).is_empty(),
2356            "the superseded negative must not still ask to be refetched"
2357        );
2358        assert_eq!(resolutions.next_control_wakeup(t(40)), secs(50));
2359    }
2360
2361    /// A whole due population is handed back in bounded waves, earliest
2362    /// deadline first.
2363    ///
2364    /// Tombstones resolved in one initial load share a deadline and come due
2365    /// together. The caller awaits the batch inline, so an uncapped set would
2366    /// hold the select loop — and with it the `tick` arm that carries
2367    /// revocation within `refresh_interval` — for the entire catalogue. The
2368    /// remainder must stay due rather than be dropped, and the order must be
2369    /// by deadline so a large population cannot starve its own tail (GL-52).
2370    #[test]
2371    fn a_due_population_is_refetched_in_bounded_waves() {
2372        let principals: Vec<Principal> = (0..5).map(Principal).collect();
2373        let mut resolutions = Resolutions::new(principals.iter().copied());
2374        for (offset, principal) in principals.iter().enumerate() {
2375            resolutions.insert(
2376                *principal,
2377                Resolution::Negative {
2378                    deadline: t(10 + offset as i64),
2379                    next_refetch: t(10 + offset as i64),
2380                    watermark: Some(Watermark::Revoked(Generation(9))),
2381                },
2382            );
2383        }
2384
2385        assert_eq!(
2386            resolutions.due_for_refetch(t(100), 2),
2387            vec![Principal(0), Principal(1)],
2388            "one wave, and the earliest deadlines lead it"
2389        );
2390        assert_eq!(
2391            resolutions.due_for_refetch(t(100), usize::MAX).len(),
2392            5,
2393            "capping a wave must not retire the rest: they are still due"
2394        );
2395    }
2396
2397    /// Lag recovery must not filter by local resolutions, because lag is
2398    /// exactly the state in which they cannot be trusted.
2399    ///
2400    /// A dropped push is most likely a reinstatement — Negative → Present —
2401    /// so the principals the recovery path exists to repair are precisely the
2402    /// ones `due_for_sweep` leaves out. The two sets must stay distinct: if
2403    /// `all_tracked` ever starts filtering, a lagged broadcast strands a
2404    /// reinstated principal until its tombstone TTL, with nothing else on the
2405    /// HTTP topology to notice (GL-52).
2406    #[test]
2407    fn lag_recovery_covers_the_negatives_a_sweep_skips() {
2408        let mut resolutions = Resolutions::new([Principal(0), Principal(1)]);
2409        resolutions.insert(
2410            Principal(0),
2411            Resolution::Present {
2412                deadline: t(90),
2413                generation: Generation(2),
2414            },
2415        );
2416        resolutions.insert(
2417            Principal(1),
2418            Resolution::Negative {
2419                deadline: t(3_600),
2420                next_refetch: t(3_600),
2421                watermark: Some(Watermark::Revoked(Generation(9))),
2422            },
2423        );
2424
2425        assert_eq!(
2426            resolutions.due_for_sweep(),
2427            vec![Principal(0)],
2428            "a routine sweep leaves the tombstone to its own schedule"
2429        );
2430        assert_eq!(
2431            resolutions.all_tracked(),
2432            vec![Principal(0), Principal(1)],
2433            "lag recovery refetches everything, tombstones included"
2434        );
2435    }
2436
2437    /// The generation watermark outlives the resolution that carried it, or a
2438    /// replayed older generation could resurrect a revoked principal
2439    /// (INVARIANTS.md GL-15).
2440    #[test]
2441    fn an_expired_resolution_keeps_its_generation() {
2442        let mut resolutions = Resolutions::new([Principal(0)]);
2443        resolutions.insert(
2444            Principal(0),
2445            Resolution::Negative {
2446                deadline: t(30),
2447                next_refetch: t(30),
2448                watermark: Some(Watermark::Revoked(Generation(7))),
2449            },
2450        );
2451        assert_eq!(resolutions.unresolved(t(100)), 1, "expired");
2452        assert_eq!(
2453            resolutions.watermark_of(Principal(0)),
2454            Some(Watermark::Revoked(Generation(7))),
2455            "the watermark must survive the expiry, provenance included"
2456        );
2457    }
2458
2459    fn secs(seconds: u64) -> std::time::Duration {
2460        std::time::Duration::from_secs(seconds)
2461    }
2462
2463    #[test]
2464    fn deadline_helpers_are_exact_in_the_supported_domain() {
2465        let config = SnapshotManagerConfig {
2466            principals: TrackedPrincipals::Fixed(vec![Principal(1)]),
2467            refresh_interval: std::time::Duration::from_secs(60),
2468            unknown_ttl: SignedDuration::from_secs(30),
2469            revoked_ttl: SignedDuration::from_secs(3_600),
2470            retry_backoff: std::time::Duration::from_secs(2),
2471            max_concurrent_fetches: 1,
2472            fetch_timeout: std::time::Duration::from_secs(5),
2473            enumeration_timeout: std::time::Duration::from_secs(30),
2474        };
2475        // An absent row: the short TTL, so a signup in flight — or a source
2476        // still coming up — is picked up soon. A published tombstone: the long
2477        // one, because coming back means an operator reinstated the account
2478        // (GL-52).
2479        assert_eq!(
2480            negative_deadline(t(10), &config, NegativeKind::Unknown),
2481            t(40)
2482        );
2483        assert_eq!(
2484            negative_deadline(t(10), &config, NegativeKind::Revoked),
2485            t(3_610)
2486        );
2487        assert_eq!(after_std(t(10), config.retry_backoff), t(12));
2488    }
2489}