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}