moq_net/model/cache.rs
1//! A shared byte budget for cached groups, repaid by write-time eviction.
2//!
3//! Every group charges its cached bytes into a [`Pool`] through a crate-internal
4//! `Charge`, billed to its track's `Track` account. The pool itself never evicts: it is
5//! a handful of atomic counters. While the pool is over capacity, each track accrues
6//! eviction debt as it writes (`accrue`),
7//! sized proportionally to what it wrote, and pays that debt by aborting its own oldest
8//! groups with [`Error::Evicted`](crate::Error::Evicted). Reclamation is therefore
9//! distributed across every writing track and converges on the capacity without any
10//! global eviction task.
11//!
12//! Cross-track ordering comes from one statistic: the mean last-access time of the
13//! evictable population (every cached group except each live track's protected latest).
14//! A group accessed more recently than that mean is never evicted, so freshly read
15//! or fetched content in one track can't die while another track holds staler
16//! content, and a track
17//! whose oldest group is staler than the mean accrues debt at double rate. Evicting
18//! old entries and inserting new ones both advance the mean, so the eviction
19//! frontier moves with cache turnover on its own.
20//!
21//! The pool also owns the wall-clock LRU window ([`Pool::expiry`]): a group that
22//! nobody has read or written for that long is reclaimed, no matter what retention
23//! its track advertises. Only a live track's latest group is exempt: once a track
24//! ends, a stale consumer can't pin any of it. Track retention
25//! ([`max_age`](crate::track::Info::max_age)) is measured in media timestamps, so a
26//! congestion stall can't age content out; the pool's expiry is the orthogonal
27//! wall-clock bound that keeps unwatched content from pinning RAM.
28//!
29//! Expiry is driven by [`Pool::gc`], also called by each origin driver.
30//! Reads and writes clear the expiration timestamp without reading a clock.
31//! The next cleanup pass dates that activity at its supplied instant. Delayed
32//! cleanup extends retention; standalone pools must call `gc` too.
33//! Byte-pressure eviction still runs inline on writes.
34//!
35//! A bare pool is inert by default ([`Pool::unbounded`]): publishers and subscribers
36//! that never set a capacity or expiry pay only a couple of atomic counters, and
37//! register nothing. A standalone [`origin`](crate::origin::Config) enables
38//! [`DEFAULT_EXPIRY`], while a relay creates one configured pool and shares it across
39//! every origin so the whole process caches into a single policy.
40
41use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
42use std::sync::{Arc, Mutex, OnceLock, Weak};
43use std::time::Duration;
44
45use super::group;
46use super::track::{self, TrackState};
47
48/// Fixed bookkeeping charged per cached group on top of its frame payload bytes.
49///
50/// A group that holds one small frame is almost entirely bookkeeping: the kio channel
51/// carrying its state, the containers the track indexes it by, and the frame slots
52/// themselves dwarf a chat-sized payload. Billing payload alone lets such a track cache
53/// millions of groups while the pool believes it is inside budget, so the process is
54/// killed before anything is evicted.
55///
56/// Derived from `size_of` rather than pasted from a measured process, so it follows the
57/// structs instead of rotting: each half lives beside the types it sizes, in
58/// [`group::CACHE_OVERHEAD`] and [`track::CACHE_OVERHEAD`]. It excludes what the
59/// allocator rounds up and what a group with many frames grows into, both of which only
60/// matter for shapes payload already dominates.
61///
62/// Also bounds the live group count (`used / ENTRY_OVERHEAD`), which keeps the
63/// access-time sum below u64 (see [`TICK_MS`]).
64pub(crate) const ENTRY_OVERHEAD: u64 = group::CACHE_OVERHEAD + track::CACHE_OVERHEAD;
65
66/// Sub-tick boosts applied to the last-access stamp, breaking ties within one
67/// coarse tick: a frame write outranks merely-inserted content, and a read (a
68/// delivered or fetched group, a frame read, a backfill's birth) outranks both.
69const WRITE_BOOST: u64 = 1;
70const READ_BOOST: u64 = 2;
71const ACCESS_SHIFT: u32 = 2;
72
73/// Milliseconds per tick of the coarse clock behind access timestamps.
74///
75/// Coarse ticks keep the count-weighted timestamp sum far from u64 overflow: the
76/// sum is bounded by `elapsed_ticks * live_groups`, plus two low tie-breaking
77/// bits. Live groups are bounded by `used / ENTRY_OVERHEAD`, and twenty years of
78/// ticks (6.3e9) times a 64 GiB target's worst case of ~70M groups is ~1.8e18 after
79/// that encoding, a tenth of `u64::MAX`. A byte-weighted mean would overflow u64
80/// even at whole-second ticks, which is why the mean is count-weighted.
81const TICK_MS: u64 = 100;
82
83/// Default idle window for standalone origins and relays.
84pub const DEFAULT_EXPIRY: Duration = Duration::from_secs(30);
85
86/// The initial policy for a [`Pool`].
87///
88/// The default is inert: no byte target and no idle expiry. Use
89/// [`Self::with_capacity`] and [`Self::with_expiry`] before creating the pool.
90#[derive(Clone, Debug, Default)]
91pub struct Config {
92 capacity: Option<u64>,
93 expiry: Option<Duration>,
94}
95
96impl Config {
97 /// Set the initial byte target. `None` leaves it unbounded.
98 pub fn with_capacity(mut self, capacity: impl Into<Option<u64>>) -> Self {
99 self.capacity = capacity.into();
100 self
101 }
102
103 /// Set the wall-clock LRU window. `None` disables idle reclamation.
104 ///
105 /// A cached group (other than a live track's latest) that nobody reads or writes
106 /// for this long is reclaimed, surfacing to any remaining reader as
107 /// [`Error::Old`](crate::Error::Old). This is independent of track retention:
108 /// [`max_age`](crate::track::Info::max_age) uses media timestamps, while this
109 /// window keeps idle content from pinning memory. Reclaiming without a write behind
110 /// it needs [`Pool::gc`] called periodically (origins do this automatically).
111 /// The value is fixed when the
112 /// pool is created. Values are rounded up to the pool's 100 ms clock tick, with
113 /// 100 ms as the minimum effective window.
114 pub fn with_expiry(mut self, expiry: impl Into<Option<Duration>>) -> Self {
115 self.expiry = expiry.into();
116 self
117 }
118}
119
120/// A shared cache policy and byte budget; cloning shares both.
121///
122/// The pool tracks how many payload bytes are cached across every registered group,
123/// plus the mean last-access time of the evictable ones. It never evicts on its own:
124/// tracks accrue eviction debt as they write and evict their own oldest groups to
125/// pay it. Idle expiration runs during [`Self::gc`]. The capacity is
126/// therefore a target usage converges toward, not a hard limit: carried debt, capped
127/// payments, and the always-protected live edge all let usage transiently exceed it.
128#[derive(Clone)]
129pub struct Pool {
130 inner: Arc<Inner>,
131}
132
133impl Default for Pool {
134 fn default() -> Self {
135 Self::unbounded()
136 }
137}
138
139struct Inner {
140 // Total bytes currently charged, including per-entry overhead.
141 used: AtomicU64,
142 // u64::MAX means unbounded.
143 capacity: AtomicU64,
144 // Wall-clock LRU window in milliseconds; u64::MAX means never expire by idleness.
145 expiry: u64,
146 // Reference point for the coarse tick clock.
147 clock: Mutex<Option<Clock>>,
148 tick: AtomicU64,
149 // Sum and count of last-access ticks across the evictable population, giving a
150 // count-weighted mean. Tracks add a group when it becomes evictable (demoted
151 // from the live edge, or inserted behind it) and remove it when it leaves.
152 access_sum: AtomicU64,
153 access_count: AtomicU64,
154 // Live track accounts, so [`Pool::sweep`] can expire idle groups in a track that
155 // has stopped writing. Empty and never touched when expiry is disabled: the byte
156 // budget needs no registry, since a track that never writes never grows the pool.
157 // Weak, because a track owns its account and the account must not outlive it.
158 tracks: kio::Lock<slab::Slab<Weak<Track>>>,
159}
160
161struct Clock {
162 epoch: crate::time::Instant,
163 now: crate::time::Instant,
164 sweep: Option<crate::time::Instant>,
165}
166
167impl Pool {
168 /// Create a pool from an initial policy.
169 ///
170 /// The budget counts frame payload bytes plus a fixed cost per cached group, which
171 /// is most of what a group carrying one small frame occupies. It is not process
172 /// RSS, and it is a convergence target rather than a hard limit; leave headroom
173 /// when sizing it from real memory. The expiry is fixed, while the capacity can
174 /// later be changed with [`Self::resize`].
175 pub fn new(config: Config) -> Self {
176 let expiry = config.expiry.map_or(u64::MAX, |expiry| {
177 let ms = u64::try_from(expiry.as_millis()).unwrap_or(u64::MAX);
178 if ms == u64::MAX {
179 return u64::MAX;
180 }
181 ms.max(1).div_ceil(TICK_MS).saturating_mul(TICK_MS)
182 });
183 let pool = Self {
184 inner: Arc::new(Inner {
185 used: AtomicU64::new(0),
186 capacity: AtomicU64::new(config.capacity.unwrap_or(u64::MAX)),
187 expiry,
188 clock: Mutex::new(None),
189 tick: AtomicU64::new(0),
190 access_sum: AtomicU64::new(0),
191 access_count: AtomicU64::new(0),
192 tracks: kio::Lock::new(slab::Slab::new()),
193 }),
194 };
195 #[cfg(test)]
196 crate::model::clock::register(&pool);
197 pool
198 }
199
200 /// Create a pool that never evicts. This is the [`Default`].
201 pub fn unbounded() -> Self {
202 Self::new(Config::default())
203 }
204
205 /// The configured byte target, or `None` when unbounded.
206 pub fn capacity(&self) -> Option<u64> {
207 match self.inner.capacity.load(Ordering::Relaxed) {
208 u64::MAX => None,
209 capacity => Some(capacity),
210 }
211 }
212
213 /// Bytes currently cached across every registered group.
214 pub fn used(&self) -> u64 {
215 self.inner.used.load(Ordering::Relaxed)
216 }
217
218 /// Change the capacity. `None` makes the pool unbounded.
219 ///
220 /// Takes effect as tracks write: a shrink leaves the pool over budget, which every
221 /// subsequent write pays down proportionally. Nothing is reclaimed synchronously.
222 pub fn resize(&self, capacity: impl Into<Option<u64>>) {
223 let capacity = capacity.into().unwrap_or(u64::MAX);
224 self.inner.capacity.store(capacity, Ordering::Relaxed);
225 }
226
227 /// The wall-clock LRU window, or `None` when idle content is never reclaimed.
228 pub fn expiry(&self) -> Option<Duration> {
229 match self.inner.expiry {
230 u64::MAX => None,
231 ms => Some(Duration::from_millis(ms)),
232 }
233 }
234
235 /// The LRU window in coarse ticks; effectively infinite when disabled.
236 pub(crate) fn expiry_ticks(&self) -> u64 {
237 match self.inner.expiry {
238 u64::MAX => u64::MAX,
239 ms => ms / TICK_MS,
240 }
241 }
242
243 /// Sample recency periodically while either cache policy is enabled.
244 ///
245 /// Expiry is approximate: activity is dated on the following cleanup pass,
246 /// and passes run at half the idle window.
247 pub(crate) fn sweep_interval(&self) -> Option<Duration> {
248 self.expiry()
249 .or_else(|| self.capacity().map(|_| DEFAULT_EXPIRY))
250 .map(|window| window / 2)
251 }
252
253 /// Expire idle groups across the registered tracks.
254 pub(crate) fn sweep(&self) {
255 // Upgrade outside each track's lock: dropping its last account unregisters it.
256 let tracks: Vec<_> = self
257 .inner
258 .tracks
259 .lock()
260 .iter()
261 .filter_map(|(_, track)| track.upgrade())
262 .collect();
263 for track in tracks {
264 track.sweep();
265 }
266 }
267
268 /// Collect idle cache entries and return the next cleanup time.
269 ///
270 /// Call after polling and at the returned deadline, including when idle.
271 /// Calls before that deadline only advance the pool's sampled clock. A due
272 /// pass visits every cached group, dating accesses since the last pass and
273 /// reclaiming idle groups except each live track's latest. Delayed calls extend
274 /// retention. Shared pools use the latest supplied instant.
275 ///
276 /// `None` means both cache policies are disabled. After enabling a capacity
277 /// with [`Self::resize`], call this again to resume periodic clock sampling.
278 pub fn gc(&self, now: crate::time::Instant) -> Option<crate::time::Instant> {
279 self.advance(now, true)
280 }
281
282 fn advance(&self, now: crate::time::Instant, sweep: bool) -> Option<crate::time::Instant> {
283 // Shared origins must not run overlapping collection passes.
284 let mut clock = self.inner.clock.lock().unwrap();
285 let clock = clock.get_or_insert(Clock {
286 epoch: now,
287 now,
288 sweep: None,
289 });
290 let now = now.max(clock.now);
291 let tick = u64::try_from(now.duration_since(clock.epoch).as_millis() / u128::from(TICK_MS))
292 .expect("cache clock overflow");
293 self.inner.tick.store(tick, Ordering::Relaxed);
294 clock.now = now;
295 if self.sweep_interval().is_none() {
296 clock.sweep = None;
297 } else if sweep && clock.sweep.is_none_or(|at| at <= now) {
298 self.sweep();
299 clock.sweep = self.sweep_interval().and_then(|interval| now.checked_add(interval));
300 }
301 clock.sweep
302 }
303
304 #[cfg(test)]
305 pub(crate) fn advance_test(&self, now: crate::time::Instant) {
306 self.advance(now, false);
307 let tracks: Vec<_> = self
308 .inner
309 .tracks
310 .lock()
311 .iter()
312 .filter_map(|(_, track)| track.upgrade())
313 .collect();
314 for track in tracks {
315 if let Some(state) = track.state.upgrade() {
316 state.read().date_cache_accesses(self.now());
317 }
318 }
319 }
320
321 /// Enter a track account into the sweep registry, returning its key. `None` when
322 /// idle reclamation is off, which is what keeps a bare pool free of bookkeeping.
323 fn register(&self, track: &Arc<Track>) -> Option<usize> {
324 self.expiry()?;
325 Some(self.inner.tracks.lock().insert(Arc::downgrade(track)))
326 }
327
328 /// Drop a track account from the sweep registry.
329 fn unregister(&self, key: usize) {
330 self.inner.tracks.lock().remove(key);
331 }
332
333 /// Returns true if both handles share the same underlying pool.
334 #[cfg(test)]
335 pub(crate) fn same_pool(&self, other: &Self) -> bool {
336 Arc::ptr_eq(&self.inner, &other.inner)
337 }
338
339 /// A handle that reaches this budget without keeping it alive.
340 pub fn downgrade(&self) -> PoolWeak {
341 PoolWeak {
342 inner: Arc::downgrade(&self.inner),
343 }
344 }
345
346 /// Charge `n` more cached bytes.
347 pub(crate) fn add(&self, n: u64) {
348 self.inner.used.fetch_add(n, Ordering::Relaxed);
349 }
350
351 /// Release `n` cached bytes.
352 pub(crate) fn sub(&self, n: u64) {
353 self.inner.used.fetch_sub(n, Ordering::Relaxed);
354 }
355
356 /// Coarse ticks since the first cleanup call.
357 pub(crate) fn now(&self) -> u64 {
358 self.inner.tick.load(Ordering::Relaxed)
359 }
360
361 /// Encode the current clock tick and an access-priority tie breaker.
362 fn stamp(&self, boost: u64) -> u64 {
363 self.now().saturating_mul(1 << ACCESS_SHIFT).saturating_add(boost)
364 }
365
366 /// Mean last-access tick across the evictable population, or `None` when it is
367 /// empty. The sum and count are read separately, so the mean is approximate
368 /// under concurrent updates; eviction only needs a rough frontier.
369 pub(crate) fn average(&self) -> Option<u64> {
370 let count = self.inner.access_count.load(Ordering::Relaxed);
371 if count == 0 {
372 return None;
373 }
374 Some(self.inner.access_sum.load(Ordering::Relaxed) / count)
375 }
376
377 /// A group with last-access tick `ts` joined the evictable population.
378 pub(crate) fn access_insert(&self, ts: u64) {
379 self.inner.access_sum.fetch_add(ts, Ordering::Relaxed);
380 self.inner.access_count.fetch_add(1, Ordering::Relaxed);
381 }
382
383 /// A group with last-access tick `ts` left the evictable population.
384 pub(crate) fn access_remove(&self, ts: u64) {
385 self.inner.access_sum.fetch_sub(ts, Ordering::Relaxed);
386 self.inner.access_count.fetch_sub(1, Ordering::Relaxed);
387 }
388
389 /// An evictable group's last-access tick moved from `old` to `new` (a FETCH hit).
390 pub(crate) fn access_refresh(&self, old: u64, new: u64) {
391 // A single wrapping add keeps the sum exact even under racing refreshes.
392 self.inner
393 .access_sum
394 .fetch_add(new.wrapping_sub(old), Ordering::Relaxed);
395 }
396
397 /// The eviction debt a track takes on by writing `written` bytes, or `None` while
398 /// the pool is under capacity (the caller should forget any outstanding debt).
399 ///
400 /// The debt is `written * used / capacity`, so paying it evicts slightly more
401 /// than was written and the overshoot decays toward the capacity. Tracks double
402 /// it when their oldest content is staler than [`Self::average`]. Saturates: a
403 /// tiny capacity must not wrap a huge debt into a small one.
404 pub(crate) fn accrue(&self, written: u64) -> Option<u64> {
405 let used = self.inner.used.load(Ordering::Relaxed);
406 let capacity = self.inner.capacity.load(Ordering::Relaxed);
407 if used <= capacity {
408 return None;
409 }
410 let debt = written as u128 * used as u128 / capacity.max(1) as u128;
411 Some(u64::try_from(debt).unwrap_or(u64::MAX))
412 }
413}
414
415impl std::fmt::Debug for Pool {
416 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
417 f.debug_struct("Pool")
418 .field("used", &self.used())
419 .field("capacity", &self.capacity())
420 .field("expiry", &self.expiry())
421 .finish()
422 }
423}
424
425/// A handle to a [`Pool`] that does not keep the budget alive.
426///
427/// [`upgrade`](Self::upgrade) stops returning a [`Pool`] once every strong handle has
428/// dropped, which is how a background resizer learns the budget it manages is gone and
429/// nothing can cache into it any more.
430#[derive(Clone)]
431pub struct PoolWeak {
432 inner: std::sync::Weak<Inner>,
433}
434
435impl PoolWeak {
436 /// Recover a [`Pool`], or `None` once every strong handle has dropped.
437 pub fn upgrade(&self) -> Option<Pool> {
438 self.inner.upgrade().map(|inner| Pool { inner })
439 }
440}
441
442impl std::fmt::Debug for PoolWeak {
443 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
444 match self.upgrade() {
445 Some(pool) => pool.fmt(f),
446 None => f.debug_struct("PoolWeak").finish_non_exhaustive(),
447 }
448 }
449}
450
451/// Gross bytes a track writes before a frame write settles its eviction debt itself,
452/// so a track appending frames to open groups (never inserting another group) still
453/// pays. Coarse: the cost is one track-state lock per threshold crossing.
454const WRITE_CHARGE_THRESHOLD: u64 = 256 * 1024;
455
456/// Maximum cadence for write-driven expiry scans, in the pool's coarse ticks.
457///
458/// Byte debt settles only after enough data accumulates, but expiry is a time
459/// policy and must also run for low-bitrate tracks. Limiting that extra track lock
460/// to once per second keeps the write hot path cheap while a bounded scan drains
461/// stale backlogs steadily.
462const EXPIRY_SCAN_TICKS: u64 = 1000 / TICK_MS;
463
464/// One track's account against the [`Pool`], shared with every group it creates.
465///
466/// Groups charge their bytes here (through a [`Charge`]) rather than straight into the
467/// pool, so the track can drain what its own groups wrote into eviction debt and pay it
468/// off by evicting them. The link back to the track is a [`kio::Weak`] because the track
469/// owns its cached groups and each of those owns this account: anything stronger would
470/// make a track's cache immortal.
471///
472/// The default account is detached: an unbounded pool and no track, so every operation
473/// is a no-op.
474#[derive(Default)]
475pub(crate) struct Track {
476 pool: Pool,
477
478 // Gross bytes charged by this track's groups (payload plus overhead), never
479 // decremented here: the track swaps it out as it accrues debt.
480 written: AtomicU64,
481
482 // Earliest coarse tick when a frame write may run another expiry scan.
483 next_expiry: AtomicU64,
484
485 // Rotating position of the expiry scan over the track's eviction order.
486 expiry_cursor: AtomicUsize,
487
488 // The track that pays this account off, holding the groups being charged.
489 state: kio::Weak<TrackState>,
490
491 // This account's slot in the pool's sweep registry, absent when the pool has no
492 // expiry window (nothing is registered) or for the detached default account.
493 sweep: OnceLock<usize>,
494}
495
496impl Track {
497 /// Open an account against `pool` for the track behind `state`.
498 pub(crate) fn new(pool: Pool, state: kio::Weak<TrackState>) -> Arc<Self> {
499 let track = Arc::new(Self {
500 pool,
501 written: AtomicU64::new(0),
502 next_expiry: AtomicU64::new(0),
503 expiry_cursor: AtomicUsize::new(0),
504 state,
505 sweep: OnceLock::new(),
506 });
507 if let Some(key) = track.pool.register(&track) {
508 let _ = track.sweep.set(key);
509 }
510 track
511 }
512
513 /// The pool this track caches into.
514 pub(crate) fn pool(&self) -> &Pool {
515 &self.pool
516 }
517
518 /// Charge a new group's fixed overhead, returning its [`Charge`].
519 pub(crate) fn charge(self: &Arc<Self>) -> Charge {
520 self.pool.add(ENTRY_OVERHEAD);
521 self.written.fetch_add(ENTRY_OVERHEAD, Ordering::Relaxed);
522 let access = Arc::new(Access::new(self.pool.stamp(0)));
523 Charge {
524 track: Some(self.clone()),
525 bytes: ENTRY_OVERHEAD,
526 access,
527 counted: false,
528 }
529 }
530
531 /// Take everything written since the last call, to be turned into eviction debt.
532 pub(crate) fn take_written(&self) -> u64 {
533 self.written.swap(0, Ordering::Relaxed)
534 }
535
536 /// Settle eviction debt and expire idle groups from a frame write.
537 ///
538 /// Called with no group lock held (locks are ordered track then group). Cheap
539 /// until the byte or time gate crosses: relaxed atomics plus a coarse clock read
540 /// when expiry is enabled. This is what makes a track that only appends frames to
541 /// open groups, never inserting another group, still pay its debt and age its
542 /// idle content out.
543 ///
544 /// `now` is the coarse tick a cache access on this path just sampled (see
545 /// [`Charge::add`]), reused so a frame write reads the clock once instead of
546 /// twice. `None` leaves the gate to sample it, and only if it needs it.
547 pub(crate) fn settle(&self, now: Option<u64>) {
548 self.settle_inner(now, false);
549 }
550
551 /// Settle from [`Pool::sweep`], dating activity and expiring every idle candidate.
552 pub(crate) fn sweep(&self) {
553 self.settle_inner(None, true);
554 }
555
556 fn settle_inner(&self, now: Option<u64>, full: bool) {
557 let settle_debt = self.written.load(Ordering::Relaxed) >= WRITE_CHARGE_THRESHOLD;
558 let scan_expiry = if full {
559 self.pool.expiry().is_some()
560 } else {
561 self.expiry_due(now)
562 };
563 if !settle_debt && !scan_expiry {
564 return;
565 }
566 // Counts as a producer while it lives, which is why `track::Producer` gates
567 // its teardown on its own clone count rather than the state's.
568 let Some(state) = self.state.upgrade() else {
569 // An ended track's channel is closed, but a stale consumer can still hold its
570 // groups: the sweep expires them in place, or they would never go.
571 if full && scan_expiry {
572 self.state.read(|state| state.expire_closed(state.expiry_scan_drain()));
573 }
574 return;
575 };
576 let expiry = if scan_expiry {
577 let state = state.read();
578 let scan = if full {
579 state.expiry_scan_drain()
580 } else {
581 state.expiry_scan()
582 };
583 state.expiry_mutation_due(scan).then_some(scan)
584 } else {
585 None
586 };
587 if !settle_debt && expiry.is_none() {
588 return;
589 }
590 if let Ok(mut state) = state.write() {
591 if settle_debt {
592 state.charge_debt();
593 }
594 if let Some(scan) = expiry {
595 state.evict_expired_scan(scan);
596 }
597 }
598 }
599
600 /// Claim the next rotating window in the track's eviction order.
601 pub(crate) fn next_expiry_scan(&self, width: usize) -> usize {
602 self.expiry_cursor.fetch_add(width, Ordering::Relaxed)
603 }
604
605 /// Claim the next write-driven expiry scan when its time gate is due.
606 fn expiry_due(&self, now: Option<u64>) -> bool {
607 let expiry = self.pool.expiry_ticks();
608 if expiry == u64::MAX {
609 return false;
610 }
611
612 // Sampled here, below the gate: a pool with no expiry window returns above
613 // without ever reading the clock, whether or not a caller had a tick.
614 let now = now.unwrap_or_else(|| self.pool.now());
615 let interval = expiry.clamp(1, EXPIRY_SCAN_TICKS);
616 let deadline = now.saturating_add(interval);
617 let next = self.next_expiry.load(Ordering::Relaxed);
618 if now < next {
619 return false;
620 }
621
622 self.next_expiry
623 .compare_exchange(next, deadline, Ordering::Relaxed, Ordering::Relaxed)
624 .is_ok()
625 }
626}
627
628impl Drop for Track {
629 fn drop(&mut self) {
630 if let Some(key) = self.sweep.get() {
631 self.pool.unregister(*key);
632 }
633 }
634}
635
636/// The RAII byte accounting for one cached group, owned by the group's state.
637///
638/// `add`/`sub` mirror the group's cached payload bytes into the pool with plain
639/// atomics, and every charged byte (overhead included) is also accumulated into the
640/// track's account, which the track drains into eviction debt on its next write. The
641/// charge also owns the group's sample in the pool's access mean, so the sample lives
642/// exactly as long as the cached bytes do: aborting or dropping the group removes both,
643/// no matter who does it or when. The default charge is detached: it belongs to no
644/// account and every operation is a no-op.
645#[derive(Default)]
646pub(crate) struct Charge {
647 track: Option<Arc<Track>>,
648 // Bytes currently charged, including ENTRY_OVERHEAD, released on drop.
649 bytes: u64,
650 // The group's last-access stamp, shared with the handles that read it during a
651 // track scan. Only ever written here, under the group's state lock.
652 access: Arc<Access>,
653 // Whether `access` is currently a sample in the pool's access mean, i.e. the
654 // group is in the evictable population. A plain bool: it is only read while
655 // updating the mean, which the owning state's lock already serializes.
656 counted: bool,
657}
658
659/// The tick of one cached group's last cache access: its creation, every write, and
660/// every read (group delivery, frame reads, FETCH hits, a fetched backfill's birth).
661/// Eviction protection and age expiry both key off it.
662///
663/// Shared, so a track scan can read it without entering the group's state. The
664/// eviction and expiry walks run under the track lock and weigh every candidate
665/// against [`Pool::average`], so reaching this through the group lock would nest one
666/// lock inside the other once per candidate.
667///
668/// Atomic for a second reason on the write side: a kio write guard's release notifies
669/// every parked consumer, and a mere cache access must not wake anyone, so [`Charge`]
670/// stamps this through a shared guard.
671#[derive(Default)]
672pub(crate) struct Access {
673 stamp: AtomicU64,
674 expires: AtomicU64,
675}
676
677impl Access {
678 fn new(stamp: u64) -> Self {
679 Self {
680 stamp: AtomicU64::new(stamp),
681 expires: AtomicU64::new(u64::MAX),
682 }
683 }
684
685 /// The stamp, tie-breaking bits included.
686 pub(crate) fn get(&self) -> u64 {
687 self.stamp.load(Ordering::Relaxed)
688 }
689
690 /// Clear the expiration timestamp until a cleanup pass observes this access.
691 pub(crate) fn touch(&self) {
692 self.expires.store(u64::MAX, Ordering::Relaxed);
693 }
694
695 /// The last access tick, assigning undated activity only during cleanup.
696 pub(crate) fn tick(&self, now: Option<u64>) -> Option<u64> {
697 let tick = match now {
698 Some(now) => match self
699 .expires
700 .compare_exchange(u64::MAX, now, Ordering::Relaxed, Ordering::Relaxed)
701 {
702 Ok(_) => now,
703 Err(tick) => tick,
704 },
705 None => self.expires.load(Ordering::Relaxed),
706 };
707 (tick != u64::MAX).then_some(tick)
708 }
709
710 /// Advance to `target` if it is newer, returning the previous stamp.
711 fn bump(&self, target: u64) -> u64 {
712 // `fetch_max` keeps the stamp monotone, and its prior value makes the
713 // paired mean update exact even for back-to-back accesses.
714 self.stamp.fetch_max(target, Ordering::Relaxed)
715 }
716}
717
718impl Charge {
719 /// Charge `n` more payload bytes, counting them as written.
720 ///
721 /// A write is also an access: it restarts the retention clock and keeps an
722 /// actively-growing group (a straggler or backfill still being filled) from
723 /// being evicted or expired mid-write, even within the same coarse tick as
724 /// content that was merely inserted.
725 ///
726 /// Returns the coarse tick it stamped, which the caller hands to
727 /// [`Track::settle`] so the write path reads the clock once rather than twice.
728 /// `None` when the charge is detached and stamped nothing.
729 pub(crate) fn add(&mut self, n: u64) -> Option<u64> {
730 if let Some(track) = &self.track {
731 track.pool.add(n);
732 track.written.fetch_add(n, Ordering::Relaxed);
733 self.bytes += n;
734 }
735 self.touch(WRITE_BOOST)
736 }
737
738 /// The group's full cached footprint: payload bytes plus overhead.
739 pub(crate) fn size(&self) -> u64 {
740 self.bytes
741 }
742
743 /// The shared handle to this group's last-access stamp, so the group can read it
744 /// without taking the state lock this charge lives behind.
745 pub(crate) fn access(&self) -> Arc<Access> {
746 self.access.clone()
747 }
748
749 /// Tick of the group's last cache access.
750 pub(crate) fn accessed(&self) -> u64 {
751 self.access.get()
752 }
753
754 /// Enter the group into the evictable population (demoted from the live edge,
755 /// or inserted behind it), sampling its access time into the pool's mean.
756 /// Idempotent.
757 pub(crate) fn demote(&mut self) {
758 if let Some(track) = &self.track
759 && !self.counted
760 {
761 track.pool.access_insert(self.accessed());
762 self.counted = true;
763 }
764 }
765
766 /// Record a cache read: a delivered or fetched group, a frame read, or a
767 /// fetched backfill's birth. `&self` so the read paths can stamp through a
768 /// shared guard without waking parked consumers.
769 pub(crate) fn refresh(&self) {
770 self.touch(READ_BOOST);
771 }
772
773 /// Record a write that charges no new bytes (a chunk written into an
774 /// already-charged in-flight frame): restarts the retention clock like any
775 /// other write. `&mut self` deliberately: reaching it through a kio write
776 /// guard marks the guard modified, so its release wakes parked readers. Returns
777 /// the stamped tick like [`Self::add`].
778 pub(crate) fn record_write(&mut self) -> Option<u64> {
779 self.touch(WRITE_BOOST)
780 }
781
782 /// Advance the last-access stamp to the current clock tick with `boost` priority.
783 ///
784 /// The boost breaks ties within one coarse tick: written content outranks
785 /// merely-inserted content, and explicitly read content outranks both, so a
786 /// same-tick access still reads as strictly newer than the population mean of
787 /// weaker accesses. Idempotent within a tick (monotone, never regressing), so
788 /// repeated accesses remain idempotent without advancing the expiry clock.
789 /// Returns the tick it read, or `None` when the charge is detached.
790 fn touch(&self, boost: u64) -> Option<u64> {
791 let track = self.track.as_ref()?;
792 // Cleanup assigns the next supplied timestamp to this access.
793 self.access.touch();
794 let target = track.pool.stamp(boost);
795 let prev = self.access.bump(target);
796 if target > prev && self.counted {
797 track.pool.access_refresh(prev, target);
798 }
799 Some(target >> ACCESS_SHIFT)
800 }
801
802 /// Release everything this charge holds: bytes, overhead, and the access
803 /// sample. Idempotent; used when the group aborts and clears its frames.
804 pub(crate) fn clear(&mut self) {
805 if let Some(track) = &self.track {
806 track.pool.sub(self.bytes);
807 self.bytes = 0;
808 if self.counted {
809 track.pool.access_remove(self.accessed());
810 self.counted = false;
811 }
812 }
813 }
814}
815
816impl Drop for Charge {
817 fn drop(&mut self) {
818 self.clear();
819 }
820}
821
822#[cfg(test)]
823mod test {
824 use super::*;
825
826 fn charge(pool: &Pool) -> Charge {
827 // No track behind the account: nothing here settles debt, it just accounts.
828 Track::new(pool.clone(), kio::Weak::new()).charge()
829 }
830
831 fn bounded(capacity: u64) -> Pool {
832 let config = Config::default().with_capacity(capacity).with_expiry(DEFAULT_EXPIRY);
833 Pool::new(config)
834 }
835
836 #[test]
837 fn unbounded_never_accrues() {
838 let pool = Pool::unbounded();
839 let mut charge = charge(&pool);
840 charge.add(1 << 40);
841 assert_eq!(pool.accrue(1 << 30), None);
842 assert_eq!(pool.used(), (1 << 40) + ENTRY_OVERHEAD);
843 drop(charge);
844 assert_eq!(pool.used(), 0);
845 }
846
847 #[test]
848 fn config_applies_capacity_and_expiry() {
849 let pool = bounded(1000);
850 assert_eq!(pool.capacity(), Some(1000));
851 assert_eq!(pool.expiry(), Some(DEFAULT_EXPIRY));
852 }
853
854 #[test]
855 fn weak_follows_the_last_strong_handle() {
856 let pool = bounded(1000);
857 let clone = pool.clone();
858 let weak = pool.downgrade();
859
860 drop(pool);
861 let upgraded = weak.upgrade().expect("a strong handle remains");
862 assert!(upgraded.same_pool(&clone));
863
864 drop(upgraded);
865 drop(clone);
866 assert!(weak.upgrade().is_none());
867 }
868
869 #[test]
870 fn accrue_none_under_capacity() {
871 let pool = bounded(ENTRY_OVERHEAD + 1000);
872 let mut charge = charge(&pool);
873 charge.add(500);
874 assert_eq!(pool.accrue(100), None);
875 }
876
877 #[test]
878 fn accrue_proportional_over_capacity() {
879 let pool = bounded(1000);
880 let mut charge = charge(&pool);
881 charge.add(2000 - ENTRY_OVERHEAD); // used = 2000, twice the capacity
882
883 // Debt exceeds what was written by the overshoot ratio, so the pool drains.
884 assert_eq!(pool.accrue(100), Some(200));
885 // Zero written accrues zero: an idle track takes on no debt.
886 assert_eq!(pool.accrue(0), Some(0));
887 }
888
889 #[test]
890 fn average_tracks_evictable_population() {
891 let pool = bounded(1000);
892 assert_eq!(pool.average(), None);
893
894 pool.access_insert(10);
895 pool.access_insert(20);
896 assert_eq!(pool.average(), Some(15));
897
898 // A refresh moves one member's contribution, exactly.
899 pool.access_refresh(10, 40);
900 assert_eq!(pool.average(), Some(30));
901
902 pool.access_remove(40);
903 assert_eq!(pool.average(), Some(20));
904 pool.access_remove(20);
905 assert_eq!(pool.average(), None);
906 }
907
908 #[test]
909 fn charge_raii() {
910 let pool = bounded(1000);
911 let mut charge = charge(&pool);
912 assert_eq!(pool.used(), ENTRY_OVERHEAD);
913
914 charge.add(100);
915 assert_eq!(pool.used(), ENTRY_OVERHEAD + 100);
916
917 charge.clear();
918 assert_eq!(pool.used(), 0);
919 // Idempotent: a second clear (and the eventual drop) releases nothing more.
920 charge.clear();
921 drop(charge);
922 assert_eq!(pool.used(), 0);
923 }
924
925 #[test]
926 fn detached_charge_is_noop() {
927 let mut charge = Charge::default();
928 charge.add(123);
929 charge.clear();
930 }
931
932 #[test]
933 fn accrue_saturates() {
934 // A huge overshoot against a tiny capacity must saturate, not wrap.
935 let pool = bounded(1);
936 let mut c = charge(&pool);
937 c.add(1 << 40);
938 assert_eq!(pool.accrue(1 << 40), Some(u64::MAX));
939 }
940
941 #[test]
942 fn charge_counts_gross_writes() {
943 let track = Track::new(bounded(1000), kio::Weak::new());
944 let mut c = track.charge();
945 c.add(100);
946 assert_eq!(track.take_written(), ENTRY_OVERHEAD + 100);
947 assert_eq!(track.take_written(), 0, "taking it drains the counter");
948 }
949
950 #[test]
951 fn charge_owns_access_sample() {
952 let pool = bounded(1000);
953 let mut c = charge(&pool);
954 assert_eq!(pool.average(), None, "not evictable until demoted");
955
956 c.demote();
957 c.demote(); // idempotent
958 assert!(pool.average().is_some());
959
960 // Clearing (an abort, from anyone) removes the sample with the bytes.
961 c.clear();
962 assert_eq!(pool.average(), None, "aborted groups leave no ghost sample");
963 drop(c);
964 assert_eq!(pool.average(), None);
965 }
966
967 #[test]
968 fn refresh_updates_a_counted_sample() {
969 let pool = bounded(1000);
970 let mut c = charge(&pool);
971 c.demote();
972 c.refresh();
973 // The sample in the pool mean moved with the stamp, so releasing the charge
974 // removes exactly what was inserted and leaves no residue.
975 assert_eq!(pool.average(), Some(c.accessed()));
976 c.clear();
977 assert_eq!(pool.average(), None);
978 }
979
980 #[test]
981 fn refresh_protects_within_a_tick() {
982 let pool = bounded(1000);
983 let mut c = charge(&pool);
984 c.demote();
985 let average = pool.average().unwrap();
986 // A refresh in the same coarse tick still lifts the group above the mean.
987 c.refresh();
988 assert!(c.accessed() > average);
989 assert_eq!(c.access().tick(None), None, "undated access is protected until cleanup");
990 assert_eq!(
991 c.access().tick(Some(pool.now())),
992 Some(pool.now()),
993 "cleanup dates the access"
994 );
995 // Repeated same-tick refreshes are idempotent, not runaway.
996 let stamped = c.accessed();
997 c.refresh();
998 assert_eq!(c.accessed(), stamped);
999 }
1000
1001 #[test]
1002 fn expiry_config() {
1003 // A bare pool preserves the unbounded contract in both dimensions.
1004 let pool = Pool::unbounded();
1005 assert_eq!(pool.expiry(), None);
1006
1007 let pool = Pool::new(Config::default().with_expiry(Duration::from_secs(1)));
1008 assert_eq!(pool.expiry(), Some(Duration::from_secs(1)));
1009 assert_eq!(pool.expiry_ticks(), 10);
1010
1011 let pool = Pool::new(Config::default().with_expiry(Duration::from_millis(1)));
1012 assert_eq!(pool.expiry(), Some(Duration::from_millis(TICK_MS)));
1013 assert_eq!(pool.expiry_ticks(), 1);
1014
1015 // Disabled: never reclaimed by idleness, but readers still re-stamp on a
1016 // bounded cadence for byte-eviction protection.
1017 let pool = Pool::new(Config::default());
1018 assert_eq!(pool.expiry(), None);
1019 assert_eq!(pool.expiry_ticks(), u64::MAX);
1020 }
1021
1022 #[test]
1023 fn expiry_gate_reuses_a_supplied_tick() {
1024 let pool = Pool::new(Config::default().with_expiry(Duration::from_secs(1)));
1025 let track = Track::new(pool, kio::Weak::new());
1026 // A supplied tick drives the gate on its own: claimed immediately, closed
1027 // until the interval elapses, claimable again on the tick it reopens.
1028 assert!(track.expiry_due(Some(0)));
1029 assert!(!track.expiry_due(Some(9)));
1030 assert!(track.expiry_due(Some(10)));
1031 // Without one the gate samples the pool clock, frozen here at tick 0, so it
1032 // stays closed rather than inheriting the caller's tick 10.
1033 assert!(!track.expiry_due(None));
1034 }
1035
1036 #[test]
1037 fn sweep_interval_is_half_the_window() {
1038 assert_eq!(Pool::unbounded().sweep_interval(), None);
1039 let pool = Pool::new(Config::default().with_expiry(Duration::from_secs(4)));
1040 assert_eq!(pool.sweep_interval(), Some(Duration::from_secs(2)));
1041 }
1042
1043 #[test]
1044 fn the_sweep_registry_follows_account_lifetime() {
1045 let pool = bounded(1000);
1046 assert!(pool.inner.tracks.lock().is_empty());
1047
1048 let track = Track::new(pool.clone(), kio::Weak::new());
1049 assert_eq!(pool.inner.tracks.lock().len(), 1);
1050
1051 // A registered account whose track is already gone is swept harmlessly.
1052 pool.sweep();
1053
1054 drop(track);
1055 assert!(pool.inner.tracks.lock().is_empty(), "a dropped account leaves no entry");
1056 }
1057
1058 #[test]
1059 fn an_inert_pool_registers_nothing() {
1060 // No window means no sweep, so a bare pool pays no registry cost.
1061 let pool = Pool::unbounded();
1062 let track = Track::new(pool.clone(), kio::Weak::new());
1063 assert!(pool.inner.tracks.lock().is_empty());
1064 pool.sweep();
1065 drop(track);
1066 }
1067
1068 #[test]
1069 fn expiry_gate_stays_closed_without_a_window() {
1070 // No window means no time gate, and no clock read to reach it.
1071 let track = Track::new(Pool::unbounded(), kio::Weak::new());
1072 assert!(!track.expiry_due(None));
1073 assert!(!track.expiry_due(Some(u64::MAX)));
1074 }
1075
1076 #[test]
1077 fn standalone_origin_enables_default_expiry() {
1078 assert_eq!(crate::origin::Config::default().pool.expiry(), Some(DEFAULT_EXPIRY));
1079 }
1080
1081 #[test]
1082 fn collecting_before_the_deadline_does_not_postpone_it() {
1083 let pool = Pool::new(Config::default().with_expiry(Duration::from_secs(2)));
1084 let now = crate::model::clock::now();
1085 let deadline = pool.gc(now);
1086 assert_eq!(pool.gc(now + Duration::from_millis(500)), deadline);
1087 }
1088
1089 #[test]
1090 fn bounded_pools_sample_recency_without_expiration() {
1091 let pool = Pool::unbounded();
1092 let now = crate::model::clock::now();
1093 assert_eq!(pool.gc(now), None);
1094 pool.resize(1024);
1095 assert_eq!(pool.gc(now), Some(now + DEFAULT_EXPIRY / 2));
1096 pool.gc(now + DEFAULT_EXPIRY);
1097 assert!(pool.now() > 0);
1098 pool.resize(None);
1099 assert_eq!(pool.gc(now + DEFAULT_EXPIRY), None);
1100 }
1101
1102 #[test]
1103 fn resize() {
1104 let pool = Pool::unbounded();
1105 assert_eq!(pool.capacity(), None);
1106
1107 let mut charge = charge(&pool);
1108 charge.add(1000);
1109
1110 // Shrinking doesn't reclaim anything synchronously; writers accrue debt instead.
1111 pool.resize(100);
1112 assert_eq!(pool.capacity(), Some(100));
1113 assert!(pool.used() > 100);
1114 assert!(pool.accrue(50).unwrap() > 50);
1115
1116 pool.resize(None);
1117 assert_eq!(pool.capacity(), None);
1118 assert_eq!(pool.accrue(50), None);
1119 }
1120}