1use super::runs::MetadataRunManifest;
5use crate::metadata::MetadataState;
6use crate::recency::Recency;
7use loonfs_api::wire::manifest::NamespaceManifestEnvelope;
8use loonfs_api::wire::sst_blocks::{DecodedDataBlock, SegmentFilter, SegmentIndexEntry};
9use loonfs_api::{ChangeSeq, ManifestId, NamespaceId};
10use serde::{Deserialize, Serialize};
11use std::collections::{HashMap, VecDeque};
12use std::sync::atomic::{AtomicUsize, Ordering};
13use std::sync::{Arc, Mutex};
14use tokio::sync::OnceCell;
15
16pub const DEFAULT_METADATA_TABLE_CACHE_DECODED_BYTES: usize = 256 * 1024 * 1024;
21
22#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
23pub struct MetadataTableCacheConfig {
24 pub max_decoded_bytes: usize,
25}
26
27impl Default for MetadataTableCacheConfig {
28 fn default() -> Self {
29 Self {
30 max_decoded_bytes: DEFAULT_METADATA_TABLE_CACHE_DECODED_BYTES,
31 }
32 }
33}
34
35#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
36pub struct MetadataTableCacheStats {
37 pub hits: usize,
38 pub misses: usize,
39 pub inserts: usize,
40 pub evictions: usize,
41 pub filter_skips: usize,
44 pub filter_false_positives: usize,
48}
49
50#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
51pub(super) enum MetadataTableBlockKind {
52 Index,
53 Filter,
54 Data,
55 Manifest,
56}
57
58#[derive(Debug, Clone, PartialEq, Eq, Hash)]
59pub(super) struct MetadataTableCacheKey {
60 pub(super) identity: String,
64 pub(super) block_kind: MetadataTableBlockKind,
65 pub(super) block_offset: u64,
66}
67
68#[derive(Debug, Clone)]
74pub(super) enum DecodedMetadataTableBlock {
75 Index {
76 entries: Arc<Vec<SegmentIndexEntry>>,
77 decoded_byte_len: usize,
78 },
79 Filter {
80 filter: Arc<SegmentFilter>,
81 decoded_byte_len: usize,
82 },
83 Data {
84 block: Arc<DecodedDataBlock>,
85 decoded_byte_len: usize,
86 },
87 Manifest {
88 manifest: Arc<NamespaceManifestEnvelope>,
89 scan_runs: Arc<Vec<MetadataRunManifest>>,
93 decoded_byte_len: usize,
94 },
95}
96
97impl DecodedMetadataTableBlock {
98 pub(super) fn decoded_byte_len(&self) -> usize {
99 match self {
100 Self::Index {
101 decoded_byte_len, ..
102 }
103 | Self::Filter {
104 decoded_byte_len, ..
105 }
106 | Self::Data {
107 decoded_byte_len, ..
108 }
109 | Self::Manifest {
110 decoded_byte_len, ..
111 } => *decoded_byte_len,
112 }
113 }
114}
115
116#[derive(Debug)]
117pub struct MetadataTableCache {
118 config: MetadataTableCacheConfig,
119 inner: Mutex<MetadataTableCacheInner>,
120 stats: MetadataTableCacheStatsInner,
121 in_flight: Mutex<HashMap<MetadataTableCacheKey, Arc<OnceCell<DecodedMetadataTableBlock>>>>,
124}
125
126#[derive(Debug, Default)]
127struct MetadataTableCacheInner {
128 entries: HashMap<MetadataTableCacheKey, CacheSlot>,
129 order: Recency<MetadataTableCacheKey>,
130 decoded_byte_len: usize,
131}
132
133#[derive(Debug)]
134struct CacheSlot {
135 block: DecodedMetadataTableBlock,
136 last_touch: u64,
139}
140
141#[derive(Debug, Default)]
142struct MetadataTableCacheStatsInner {
143 hits: AtomicUsize,
144 misses: AtomicUsize,
145 inserts: AtomicUsize,
146 evictions: AtomicUsize,
147 filter_skips: AtomicUsize,
148 filter_false_positives: AtomicUsize,
149}
150
151impl MetadataTableCache {
152 pub fn new(config: MetadataTableCacheConfig) -> Self {
153 Self {
154 config,
155 inner: Mutex::new(MetadataTableCacheInner::default()),
156 stats: MetadataTableCacheStatsInner::default(),
157 in_flight: Mutex::new(HashMap::new()),
158 }
159 }
160
161 pub(super) async fn get_or_load<E, F, Fut>(
163 &self,
164 cache_key: &MetadataTableCacheKey,
165 fetch: F,
166 ) -> Result<DecodedMetadataTableBlock, E>
167 where
168 F: FnOnce() -> Fut,
169 Fut: std::future::Future<Output = Result<DecodedMetadataTableBlock, E>>,
170 {
171 let cell = {
172 let mut in_flight = self
173 .in_flight
174 .lock()
175 .expect("metadata table cache in-flight lock should not be poisoned");
176 Arc::clone(
177 in_flight
178 .entry(cache_key.clone())
179 .or_insert_with(|| Arc::new(OnceCell::new())),
180 )
181 };
182 let result = cell
183 .get_or_try_init(|| async {
184 if let Some(block) = self.get(cache_key) {
185 return Ok(block);
186 }
187 let block = fetch().await?;
188 self.insert(cache_key.clone(), block.clone());
189 Ok(block)
190 })
191 .await
192 .cloned();
193 let mut in_flight = self
194 .in_flight
195 .lock()
196 .expect("metadata table cache in-flight lock should not be poisoned");
197 if in_flight
198 .get(cache_key)
199 .is_some_and(|current| Arc::ptr_eq(current, &cell))
200 {
201 in_flight.remove(cache_key);
202 }
203 result
204 }
205
206 pub fn stats(&self) -> MetadataTableCacheStats {
207 MetadataTableCacheStats {
208 hits: self.stats.hits.load(Ordering::SeqCst),
209 misses: self.stats.misses.load(Ordering::SeqCst),
210 inserts: self.stats.inserts.load(Ordering::SeqCst),
211 evictions: self.stats.evictions.load(Ordering::SeqCst),
212 filter_skips: self.stats.filter_skips.load(Ordering::SeqCst),
213 filter_false_positives: self.stats.filter_false_positives.load(Ordering::SeqCst),
214 }
215 }
216
217 pub(super) fn record_filter_skip(&self) {
218 self.stats.filter_skips.fetch_add(1, Ordering::SeqCst);
219 }
220
221 pub(super) fn record_filter_false_positive(&self) {
222 self.stats
223 .filter_false_positives
224 .fetch_add(1, Ordering::SeqCst);
225 }
226
227 pub(super) fn get(&self, key: &MetadataTableCacheKey) -> Option<DecodedMetadataTableBlock> {
228 if self.config.max_decoded_bytes == 0 {
229 return None;
230 }
231 let mut inner = self
232 .inner
233 .lock()
234 .expect("metadata table cache lock should not be poisoned");
235 let Some(block) = inner.entries.get(key).map(|slot| slot.block.clone()) else {
236 self.stats.misses.fetch_add(1, Ordering::SeqCst);
237 return None;
238 };
239 inner.touch(key);
240 self.stats.hits.fetch_add(1, Ordering::SeqCst);
241 Some(block)
242 }
243
244 pub(super) fn insert(&self, key: MetadataTableCacheKey, block: DecodedMetadataTableBlock) {
245 if self.config.max_decoded_bytes == 0 {
246 return;
247 }
248 let mut inner = self
249 .inner
250 .lock()
251 .expect("metadata table cache lock should not be poisoned");
252 let decoded_byte_len = block.decoded_byte_len();
253 if let Some(previous) = inner.entries.insert(
254 key.clone(),
255 CacheSlot {
256 block,
257 last_touch: 0,
258 },
259 ) {
260 inner.decoded_byte_len = inner
261 .decoded_byte_len
262 .saturating_sub(previous.block.decoded_byte_len());
263 }
264 inner.decoded_byte_len = inner.decoded_byte_len.saturating_add(decoded_byte_len);
265 inner.touch(&key);
266 self.stats.inserts.fetch_add(1, Ordering::SeqCst);
267 let MetadataTableCacheInner {
268 entries,
269 order,
270 decoded_byte_len,
271 } = &mut *inner;
272 while *decoded_byte_len > self.config.max_decoded_bytes {
273 let Some(candidate) = order.pop_oldest(|key, stamp| slot_is_live(entries, key, stamp))
274 else {
275 break;
276 };
277 if let Some(slot) = entries.remove(&candidate) {
278 *decoded_byte_len = decoded_byte_len.saturating_sub(slot.block.decoded_byte_len());
279 self.stats.evictions.fetch_add(1, Ordering::SeqCst);
280 }
281 }
282 }
283}
284
285impl MetadataTableCacheInner {
286 fn touch(&mut self, key: &MetadataTableCacheKey) {
287 let stamp = self.order.touch(key);
288 if let Some(slot) = self.entries.get_mut(key) {
289 slot.last_touch = stamp;
290 }
291 let entries = &self.entries;
292 self.order.compact(entries.len(), |key, stamp| {
293 slot_is_live(entries, key, stamp)
294 });
295 }
296}
297
298fn slot_is_live(
301 entries: &HashMap<MetadataTableCacheKey, CacheSlot>,
302 key: &MetadataTableCacheKey,
303 stamp: u64,
304) -> bool {
305 entries
306 .get(key)
307 .is_some_and(|slot| slot.last_touch == stamp)
308}
309
310pub const DEFAULT_WAL_TAIL_PROJECTION_ROWS: usize = 1_000_000;
313pub const DEFAULT_WAL_TAIL_PROJECTION_DECODED_BYTES: usize = 256 * 1024 * 1024;
314
315#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
318pub struct WalTailProjectionCacheConfig {
319 pub max_entries: usize,
320 pub max_rows: usize,
321 pub max_decoded_bytes: usize,
322}
323
324#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
325pub struct WalTailProjectionCacheStats {
326 pub hits: usize,
327 pub misses: usize,
328 pub inserts: usize,
329 pub evictions: usize,
330 pub evicted_rows: usize,
331 pub evicted_decoded_bytes: usize,
332 pub uncacheable_count: usize,
333 pub uncacheable_rows: usize,
334 pub uncacheable_decoded_bytes: usize,
335 pub cached_rows: usize,
336 pub cached_decoded_bytes: usize,
337}
338
339#[derive(Debug, Clone, PartialEq, Eq, Hash)]
340pub struct WalTailProjectionCacheKey {
341 pub namespace_id: NamespaceId,
342 pub manifest_id: ManifestId,
343 pub manifest_head_seq: ChangeSeq,
344 pub head_seq: ChangeSeq,
345 pub head_etag: String,
346}
347
348#[derive(Debug, Clone)]
349struct CachedWalTailProjection {
350 rows: Arc<MetadataState>,
351 row_count: usize,
352 decoded_bytes: usize,
353}
354
355impl CachedWalTailProjection {
356 fn new(rows: Arc<MetadataState>) -> Self {
357 Self {
358 row_count: rows.row_count(),
359 decoded_bytes: rows.decoded_bytes(),
360 rows,
361 }
362 }
363
364 fn rows(&self) -> Arc<MetadataState> {
365 Arc::clone(&self.rows)
366 }
367
368 fn weight(&self) -> (usize, usize) {
369 (self.row_count, self.decoded_bytes)
370 }
371}
372
373#[derive(Debug)]
374pub struct WalTailProjectionCache {
375 config: WalTailProjectionCacheConfig,
376 inner: Mutex<WalTailProjectionCacheInner>,
377 stats: WalTailProjectionCacheStatsInner,
378}
379
380#[derive(Debug, Default)]
381struct WalTailProjectionCacheInner {
382 entries: HashMap<WalTailProjectionCacheKey, CachedWalTailProjection>,
383 order: VecDeque<WalTailProjectionCacheKey>,
384 cached_rows: usize,
385 cached_decoded_bytes: usize,
386}
387
388#[derive(Debug, Default)]
389struct WalTailProjectionCacheStatsInner {
390 hits: AtomicUsize,
391 misses: AtomicUsize,
392 inserts: AtomicUsize,
393 evictions: AtomicUsize,
394 evicted_rows: AtomicUsize,
395 evicted_decoded_bytes: AtomicUsize,
396 uncacheable_count: AtomicUsize,
397 uncacheable_rows: AtomicUsize,
398 uncacheable_decoded_bytes: AtomicUsize,
399}
400
401impl WalTailProjectionCache {
402 pub fn new(config: WalTailProjectionCacheConfig) -> Self {
403 Self {
404 config,
405 inner: Mutex::new(WalTailProjectionCacheInner::default()),
406 stats: WalTailProjectionCacheStatsInner::default(),
407 }
408 }
409
410 pub fn stats(&self) -> WalTailProjectionCacheStats {
411 let inner = self
412 .inner
413 .lock()
414 .expect("wal tail projection cache lock should not be poisoned");
415 WalTailProjectionCacheStats {
416 hits: self.stats.hits.load(Ordering::SeqCst),
417 misses: self.stats.misses.load(Ordering::SeqCst),
418 inserts: self.stats.inserts.load(Ordering::SeqCst),
419 evictions: self.stats.evictions.load(Ordering::SeqCst),
420 evicted_rows: self.stats.evicted_rows.load(Ordering::SeqCst),
421 evicted_decoded_bytes: self.stats.evicted_decoded_bytes.load(Ordering::SeqCst),
422 uncacheable_count: self.stats.uncacheable_count.load(Ordering::SeqCst),
423 uncacheable_rows: self.stats.uncacheable_rows.load(Ordering::SeqCst),
424 uncacheable_decoded_bytes: self.stats.uncacheable_decoded_bytes.load(Ordering::SeqCst),
425 cached_rows: inner.cached_rows,
426 cached_decoded_bytes: inner.cached_decoded_bytes,
427 }
428 }
429
430 pub fn get(&self, key: &WalTailProjectionCacheKey) -> Option<Arc<MetadataState>> {
431 if self.config.max_entries == 0 {
432 return None;
433 }
434 let mut inner = self
435 .inner
436 .lock()
437 .expect("wal tail projection cache lock should not be poisoned");
438 let Some(rows) = inner.entries.get(key).map(CachedWalTailProjection::rows) else {
439 self.stats.misses.fetch_add(1, Ordering::SeqCst);
440 return None;
441 };
442 inner.touch(key);
443 self.stats.hits.fetch_add(1, Ordering::SeqCst);
444 Some(rows)
445 }
446
447 pub fn insert(&self, key: WalTailProjectionCacheKey, rows: Arc<MetadataState>) {
448 if self.config.max_entries == 0 {
449 return;
450 }
451 let cached = CachedWalTailProjection::new(rows);
452 let (row_count, decoded_bytes) = cached.weight();
453 if row_count > self.config.max_rows || decoded_bytes > self.config.max_decoded_bytes {
454 self.stats.uncacheable_count.fetch_add(1, Ordering::SeqCst);
455 self.stats
456 .uncacheable_rows
457 .fetch_add(row_count, Ordering::SeqCst);
458 self.stats
459 .uncacheable_decoded_bytes
460 .fetch_add(decoded_bytes, Ordering::SeqCst);
461 return;
462 }
463
464 let mut inner = self
465 .inner
466 .lock()
467 .expect("wal tail projection cache lock should not be poisoned");
468 if let Some(previous) = inner.entries.insert(key.clone(), cached) {
469 let (rows, bytes) = previous.weight();
470 inner.cached_rows = inner.cached_rows.saturating_sub(rows);
471 inner.cached_decoded_bytes = inner.cached_decoded_bytes.saturating_sub(bytes);
472 }
473 inner.cached_rows = inner.cached_rows.saturating_add(row_count);
474 inner.cached_decoded_bytes = inner.cached_decoded_bytes.saturating_add(decoded_bytes);
475 inner.touch(&key);
476 self.stats.inserts.fetch_add(1, Ordering::SeqCst);
477
478 while inner.entries.len() > self.config.max_entries
479 || inner.cached_rows > self.config.max_rows
480 || inner.cached_decoded_bytes > self.config.max_decoded_bytes
481 {
482 let Some(evicted) = inner.order.pop_front() else {
483 break;
484 };
485 if let Some(previous) = inner.entries.remove(&evicted) {
486 let (rows, bytes) = previous.weight();
487 inner.cached_rows = inner.cached_rows.saturating_sub(rows);
488 inner.cached_decoded_bytes = inner.cached_decoded_bytes.saturating_sub(bytes);
489 self.stats.evictions.fetch_add(1, Ordering::SeqCst);
490 self.stats.evicted_rows.fetch_add(rows, Ordering::SeqCst);
491 self.stats
492 .evicted_decoded_bytes
493 .fetch_add(bytes, Ordering::SeqCst);
494 }
495 }
496 }
497
498 pub fn invalidate_namespace(&self, namespace_id: &NamespaceId) {
499 let mut inner = self
500 .inner
501 .lock()
502 .expect("wal tail projection cache lock should not be poisoned");
503 let keys = inner
504 .entries
505 .keys()
506 .filter(|key| &key.namespace_id == namespace_id)
507 .cloned()
508 .collect::<Vec<_>>();
509 for key in keys {
510 if let Some(previous) = inner.entries.remove(&key) {
511 let (rows, bytes) = previous.weight();
512 inner.cached_rows = inner.cached_rows.saturating_sub(rows);
513 inner.cached_decoded_bytes = inner.cached_decoded_bytes.saturating_sub(bytes);
514 self.stats.evictions.fetch_add(1, Ordering::SeqCst);
515 self.stats.evicted_rows.fetch_add(rows, Ordering::SeqCst);
516 self.stats
517 .evicted_decoded_bytes
518 .fetch_add(bytes, Ordering::SeqCst);
519 }
520 }
521 inner.order.retain(|key| &key.namespace_id != namespace_id);
522 }
523}
524
525impl WalTailProjectionCacheInner {
526 fn touch(&mut self, key: &WalTailProjectionCacheKey) {
527 self.order.retain(|candidate| candidate != key);
528 self.order.push_back(key.clone());
529 }
530}
531
532#[cfg(test)]
533mod tests {
534 use super::{
535 DecodedMetadataTableBlock, MetadataTableBlockKind, MetadataTableCache,
536 MetadataTableCacheConfig, MetadataTableCacheKey,
537 };
538 use loonfs_api::wire::sst_blocks::DecodedDataBlock;
539 use std::sync::Arc;
540
541 fn block(decoded_byte_len: usize) -> DecodedMetadataTableBlock {
542 DecodedMetadataTableBlock::Data {
543 block: Arc::new(DecodedDataBlock {
544 row_keys: Vec::new(),
545 rows: Vec::new(),
546 }),
547 decoded_byte_len,
548 }
549 }
550
551 fn key(digest: &str) -> MetadataTableCacheKey {
552 MetadataTableCacheKey {
553 identity: digest.to_owned(),
554 block_kind: MetadataTableBlockKind::Data,
555 block_offset: 0,
556 }
557 }
558
559 #[test]
560 fn default_config_budgets_bytes() {
561 let config = MetadataTableCacheConfig::default();
562 assert_eq!(
563 config.max_decoded_bytes,
564 super::DEFAULT_METADATA_TABLE_CACHE_DECODED_BYTES
565 );
566 }
567
568 #[test]
569 fn byte_budget_evicts_the_oldest_block() {
570 let cache = MetadataTableCache::new(MetadataTableCacheConfig {
571 max_decoded_bytes: 1000,
572 });
573 cache.insert(key("a"), block(600));
574 cache.insert(key("b"), block(600));
575 assert!(
576 cache.get(&key("a")).is_none(),
577 "oldest block should evict once the byte budget is exceeded"
578 );
579 assert!(cache.get(&key("b")).is_some());
580 assert_eq!(cache.stats().evictions, 1);
581 }
582
583 #[test]
584 fn replacing_a_block_reaccounts_its_decoded_bytes() {
585 let cache = MetadataTableCache::new(MetadataTableCacheConfig {
586 max_decoded_bytes: 1000,
587 });
588 cache.insert(key("a"), block(600));
589 cache.insert(key("a"), block(100));
590 cache.insert(key("b"), block(600));
592 assert!(cache.get(&key("a")).is_some());
593 assert!(cache.get(&key("b")).is_some());
594 assert_eq!(cache.stats().evictions, 0);
595 }
596
597 #[test]
598 fn recency_queue_stays_bounded_under_repeated_hits() {
599 let cache = MetadataTableCache::new(MetadataTableCacheConfig {
600 max_decoded_bytes: 10_000,
601 });
602 for index in 0..8 {
603 cache.insert(key(&format!("k{index}")), block(10));
604 }
605 for _ in 0..10_000 {
606 cache.get(&key("k0"));
607 cache.get(&key("k3"));
608 }
609 let inner = cache.inner.lock().expect("cache lock");
610 assert_eq!(inner.entries.len(), 8);
611 assert!(
612 inner.order.positions() <= (inner.entries.len() * 2).max(16),
613 "hits must not grow the recency queue unboundedly, queue = {}",
614 inner.order.positions()
615 );
616 }
617
618 #[test]
619 fn cache_hits_share_the_decoded_row_allocation() {
620 let cache = MetadataTableCache::new(MetadataTableCacheConfig::default());
621 let inserted = block(64);
622 let rows = match &inserted {
623 DecodedMetadataTableBlock::Data { block: rows, .. } => Arc::clone(rows),
624 DecodedMetadataTableBlock::Index { .. }
625 | DecodedMetadataTableBlock::Filter { .. }
626 | DecodedMetadataTableBlock::Manifest { .. } => {
627 unreachable!("fixture builds a data block")
628 }
629 };
630 cache.insert(key("a"), inserted);
631 let hit = cache.get(&key("a")).expect("inserted block should hit");
632 let shares_allocation = match &hit {
633 DecodedMetadataTableBlock::Data {
634 block: hit_rows, ..
635 } => Arc::ptr_eq(hit_rows, &rows),
636 DecodedMetadataTableBlock::Index { .. }
637 | DecodedMetadataTableBlock::Filter { .. }
638 | DecodedMetadataTableBlock::Manifest { .. } => false,
639 };
640 assert!(
641 shares_allocation,
642 "a cache hit should share the decoded rows, not clone them"
643 );
644 }
645
646 #[tokio::test]
647 async fn get_or_load_retries_after_a_failed_load() {
648 let cache = MetadataTableCache::new(MetadataTableCacheConfig::default());
649 let failed: Result<_, String> = cache
650 .get_or_load(&key("a"), || async { Err("transport".to_owned()) })
651 .await;
652 assert!(failed.is_err());
653 let recovered: Result<_, String> = cache
654 .get_or_load(&key("a"), || async { Ok(block(1)) })
655 .await;
656 assert!(
657 recovered.is_ok(),
658 "a failed fetch should leave nothing behind for the next caller"
659 );
660 }
661
662 #[tokio::test]
663 async fn get_or_load_counts_one_miss_and_populates_for_later_hits() {
664 let cache = MetadataTableCache::new(MetadataTableCacheConfig::default());
665 let fetched: Result<_, String> = cache
666 .get_or_load(&key("a"), || async { Ok(block(1)) })
667 .await;
668 assert!(fetched.is_ok());
669 let cached: Result<_, String> = cache
670 .get_or_load(&key("a"), || async {
671 Err("a populated key must not re-fetch".to_owned())
672 })
673 .await;
674 assert!(cached.is_ok(), "the cached block should answer the access");
675
676 let stats = cache.stats();
677 assert_eq!(stats.misses, 1);
678 assert_eq!(stats.inserts, 1);
679 assert_eq!(stats.hits, 1);
680 }
681}