use std::{
collections::BTreeMap,
sync::Arc,
time::{Duration, Instant},
};
use arc_swap::ArcSwap;
use dashmap::DashMap;
use futures::future::join_all;
use roaring::RoaringBitmap;
use uuid::Uuid;
use crate::{
runtime_bridge::bridge_sync_to_async,
supertable::wal::{SealRecord, WalStore},
};
pub const DEFAULT_SEAL_TTL: Duration = Duration::from_secs(1);
#[derive(Debug, Default)]
pub struct TombstoneSeqView {
pub manifest_id: u64,
pub seqs: BTreeMap<Uuid, u64>,
}
#[derive(Debug, thiserror::Error)]
pub enum SidecarCacheError {
#[error("tombstone sidecar refresh failed for {superfile_id}: {message}")]
RefreshFailed { superfile_id: Uuid, message: String },
}
#[derive(Debug)]
pub struct SidecarCache {
inner: DashMap<Uuid, CachedSidecar>,
seal_ttl: Duration,
seq_view: ArcSwap<TombstoneSeqView>,
empty_bitmap: Arc<RoaringBitmap>,
wal_store: WalStore,
}
#[derive(Debug, Clone)]
struct CachedSidecar {
#[allow(dead_code)]
etag: Option<String>,
bitmap: Arc<RoaringBitmap>,
seal: Option<SealRecord>,
seq: u64,
fetched_at: Instant,
}
#[derive(Clone, Copy, PartialEq, Eq)]
enum Freshness {
Bitmap,
Seal,
}
impl SidecarCache {
pub fn new(
wal_store: WalStore,
seal_ttl: Duration,
initial_view: Arc<TombstoneSeqView>,
) -> Self {
Self {
inner: DashMap::new(),
seal_ttl,
seq_view: ArcSwap::new(initial_view),
empty_bitmap: Arc::new(RoaringBitmap::new()),
wal_store,
}
}
pub fn reconcile(&self, view: Arc<TombstoneSeqView>) {
loop {
let current = self.seq_view.load();
if view.manifest_id <= current.manifest_id {
return;
}
let prev = self.seq_view.compare_and_swap(&*current, Arc::clone(&view));
if Arc::ptr_eq(&prev, ¤t) {
return;
}
}
}
pub fn view_manifest_id(&self) -> u64 {
self.seq_view.load().manifest_id
}
pub async fn prefetch(&self, superfile_ids: &[Uuid], now: Instant) {
let view = self.seq_view.load();
let stale: Vec<(Uuid, u64)> = superfile_ids
.iter()
.filter_map(|id| {
let expected = view.seqs.get(id).copied()?;
match self.inner.get(id) {
Some(entry) if entry.seq == expected => None,
_ => Some((*id, expected)),
}
})
.collect();
if stale.is_empty() {
return;
}
let fetches = stale.into_iter().map(|(id, expected)| {
let wal_store = self.wal_store.clone();
async move { (id, expected, wal_store.get_tombstones(id).await) }
});
let results = join_all(fetches).await;
for (id, expected, result) in results {
let (bitmap, seal, etag) = match result {
Ok(Some((sidecar, etag))) => (Arc::new(sidecar.bitmap), sidecar.seal, Some(etag)),
Ok(None) => (Arc::clone(&self.empty_bitmap), None, None),
Err(_) => continue,
};
self.inner.insert(
id,
CachedSidecar {
etag,
bitmap,
seal,
seq: expected,
fetched_at: now,
},
);
}
}
fn fetch_sidecar(
&self,
superfile_id: Uuid,
now: Instant,
freshness: Freshness,
) -> Result<(Arc<RoaringBitmap>, Option<SealRecord>), SidecarCacheError> {
let view = self.seq_view.load();
let Some(expected) = view.seqs.get(&superfile_id).copied() else {
return Ok((Arc::clone(&self.empty_bitmap), None));
};
if let Some(entry) = self.inner.get(&superfile_id) {
let seq_fresh = entry.seq == expected;
let fresh = match freshness {
Freshness::Bitmap => seq_fresh,
Freshness::Seal => {
seq_fresh && now.duration_since(entry.fetched_at) < self.seal_ttl
}
};
if fresh {
return Ok((Arc::clone(&entry.bitmap), entry.seal.clone()));
}
}
self.refresh_and_return_sidecar(superfile_id, expected)
}
pub fn bitmap_for(
&self,
superfile_id: Uuid,
now: Instant,
) -> Result<Arc<RoaringBitmap>, SidecarCacheError> {
self.fetch_sidecar(superfile_id, now, Freshness::Bitmap)
.map(|(bitmap, _)| bitmap)
}
pub fn seal_for(
&self,
superfile_id: Uuid,
now: Instant,
) -> Result<Option<SealRecord>, SidecarCacheError> {
self.fetch_sidecar(superfile_id, now, Freshness::Seal)
.map(|(_, seal)| seal)
}
pub fn sidecar_for(
&self,
superfile_id: Uuid,
now: Instant,
) -> Result<(Arc<RoaringBitmap>, Option<SealRecord>), SidecarCacheError> {
self.fetch_sidecar(superfile_id, now, Freshness::Seal)
}
#[cfg(test)]
pub fn clear(&self) {
self.inner.clear();
}
pub fn len(&self) -> usize {
self.inner.len()
}
pub fn is_empty(&self) -> bool {
self.inner.is_empty()
}
fn refresh_and_return_sidecar(
&self,
superfile_id: Uuid,
expected_seq: u64,
) -> Result<(Arc<RoaringBitmap>, Option<SealRecord>), SidecarCacheError> {
let wal_store = self.wal_store.clone();
let result =
bridge_sync_to_async(async move { wal_store.get_tombstones(superfile_id).await });
let (bitmap, seal, etag) = match result {
Ok(Some((sidecar, etag))) => (Arc::new(sidecar.bitmap), sidecar.seal, Some(etag)),
Ok(None) => (Arc::clone(&self.empty_bitmap), None, None),
Err(e) => {
return Err(SidecarCacheError::RefreshFailed {
superfile_id,
message: format!("{e}"),
});
}
};
let entry = CachedSidecar {
etag,
bitmap: Arc::clone(&bitmap),
seal: seal.clone(),
seq: expected_seq,
fetched_at: Instant::now(),
};
self.inner.insert(superfile_id, entry);
Ok((bitmap, seal))
}
}
#[cfg(test)]
mod tests {
use std::iter::once;
use chrono::Utc;
use tempfile::TempDir;
use super::*;
use crate::{
storage::{LocalFsStorageProvider, StorageProvider},
supertable::wal::tombstones_codec::TombstonesSidecar,
};
fn fixture() -> (TempDir, WalStore, SidecarCache) {
let dir = TempDir::new().expect("tempdir");
let storage: Arc<dyn StorageProvider> =
Arc::new(LocalFsStorageProvider::new(dir.path()).expect("provider"));
let ws = WalStore::new(Arc::clone(&storage));
let cache = SidecarCache::new(
ws.clone(),
DEFAULT_SEAL_TTL,
Arc::new(TombstoneSeqView::default()),
);
(dir, ws, cache)
}
fn view(manifest_id: u64, entries: &[(Uuid, u64)]) -> Arc<TombstoneSeqView> {
Arc::new(TombstoneSeqView {
manifest_id,
seqs: entries.iter().copied().collect(),
})
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn unmapped_superfile_is_empty_with_no_entry_and_no_get() {
let (_dir, ws, cache) = fixture();
let sf_id = Uuid::from_u128(0xAB);
let mut bitmap = RoaringBitmap::new();
bitmap.insert(11);
ws.put_tombstones(sf_id, None, &TombstonesSidecar { seal: None, bitmap })
.await
.expect("put");
let cached = cache.bitmap_for(sf_id, Instant::now()).expect("lookup");
assert!(cached.is_empty());
assert_eq!(cache.len(), 0, "no entry is materialized for unmapped ids");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn lookup_reflects_persisted_sidecar() {
let (_dir, ws, cache) = fixture();
let sf_id = Uuid::from_u128(0xCAFE);
let mut bitmap = RoaringBitmap::new();
bitmap.insert(1);
bitmap.insert(3);
bitmap.insert(5);
let sidecar = TombstonesSidecar { seal: None, bitmap };
ws.put_tombstones(sf_id, None, &sidecar).await.expect("put");
cache.reconcile(view(1, &[(sf_id, 1)]));
let cached = cache.bitmap_for(sf_id, Instant::now()).expect("lookup");
let collected: Vec<u32> = cached.iter().collect();
assert_eq!(collected, vec![1u32, 3, 5]);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn matching_seq_serves_cached_view_without_refresh() {
let (_dir, ws, cache) = fixture();
let sf_id = Uuid::from_u128(0xDEAD);
let mut bitmap = RoaringBitmap::new();
bitmap.insert(1);
ws.put_tombstones(sf_id, None, &TombstonesSidecar { seal: None, bitmap })
.await
.expect("put");
cache.reconcile(view(1, &[(sf_id, 1)]));
let now = Instant::now();
let first = cache.bitmap_for(sf_id, now).expect("warm");
assert_eq!(first.iter().collect::<Vec<_>>(), vec![1u32]);
let (_, etag) = ws
.get_tombstones(sf_id)
.await
.expect("read")
.expect("present");
let mut bumped = RoaringBitmap::new();
bumped.insert(1);
bumped.insert(2);
ws.put_tombstones(
sf_id,
Some(&etag),
&TombstonesSidecar {
seal: None,
bitmap: bumped,
},
)
.await
.expect("overwrite");
let cached = cache.bitmap_for(sf_id, now).expect("cached read");
assert_eq!(
cached.iter().collect::<Vec<_>>(),
vec![1u32],
"cache must hold the seq-fresh view"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn seq_bump_forces_next_lookup_to_refresh() {
let (_dir, ws, cache) = fixture();
let sf_id = Uuid::from_u128(0xBEEF);
let mut bitmap = RoaringBitmap::new();
bitmap.insert(7);
ws.put_tombstones(sf_id, None, &TombstonesSidecar { seal: None, bitmap })
.await
.expect("put");
cache.reconcile(view(1, &[(sf_id, 1)]));
let now = Instant::now();
let first = cache.bitmap_for(sf_id, now).expect("warm");
assert_eq!(first.iter().collect::<Vec<_>>(), vec![7u32]);
let (_, etag) = ws
.get_tombstones(sf_id)
.await
.expect("read")
.expect("present");
let mut grown = RoaringBitmap::new();
grown.insert(7);
grown.insert(9);
ws.put_tombstones(
sf_id,
Some(&etag),
&TombstonesSidecar {
seal: None,
bitmap: grown,
},
)
.await
.expect("grow");
cache.reconcile(view(2, &[(sf_id, 2)]));
let cached = cache.bitmap_for(sf_id, now).expect("re-read");
assert_eq!(cached.iter().collect::<Vec<_>>(), vec![7u32, 9]);
}
#[test]
fn reconcile_is_forward_only() {
let (_dir, _ws, cache) = fixture();
let sf_id = Uuid::from_u128(0x77);
cache.reconcile(view(5, &[(sf_id, 5)]));
assert_eq!(cache.view_manifest_id(), 5);
cache.reconcile(view(3, &[]));
assert_eq!(cache.view_manifest_id(), 5);
cache.reconcile(view(6, &[(sf_id, 6)]));
assert_eq!(cache.view_manifest_id(), 6);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn prefetch_populates_mapped_ids_in_one_batch() {
let (_dir, ws, cache) = fixture();
let present = Uuid::from_u128(0x01);
let mut bitmap = RoaringBitmap::new();
bitmap.insert(9);
ws.put_tombstones(present, None, &TombstonesSidecar { seal: None, bitmap })
.await
.expect("put");
cache.reconcile(view(1, &[(present, 1)]));
let ids: Vec<Uuid> = once(present)
.chain((2..32u128).map(Uuid::from_u128))
.collect();
let now = Instant::now();
cache.prefetch(&ids, now).await;
assert_eq!(cache.len(), 1);
assert_eq!(
cache
.bitmap_for(present, now)
.expect("present")
.iter()
.collect::<Vec<_>>(),
vec![9u32]
);
assert!(cache.bitmap_for(ids[1], now).expect("absent").is_empty());
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn seal_for_returns_cached_seal() {
let (_dir, ws, cache) = fixture();
let sf_id = Uuid::from_u128(0xFFFF);
let mut bitmap = RoaringBitmap::new();
bitmap.insert(42);
ws.put_tombstones(sf_id, None, &TombstonesSidecar { seal: None, bitmap })
.await
.expect("put");
cache.reconcile(view(1, &[(sf_id, 1)]));
let now = Instant::now();
let seal = cache.seal_for(sf_id, now).expect("lookup");
assert!(seal.is_none(), "initially unsealed");
let seal_2 = cache.seal_for(sf_id, now).expect("cached");
assert!(seal_2.is_none());
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn seal_lookup_past_ttl_observes_new_seal_without_seq_bump() {
let dir = TempDir::new().expect("tempdir");
let storage: Arc<dyn StorageProvider> =
Arc::new(LocalFsStorageProvider::new(dir.path()).expect("provider"));
let ws = WalStore::new(Arc::clone(&storage));
let cache = SidecarCache::new(
ws.clone(),
Duration::ZERO,
Arc::new(TombstoneSeqView::default()),
);
let sf_id = Uuid::from_u128(0x5EA1);
let mut bitmap = RoaringBitmap::new();
bitmap.insert(3);
ws.put_tombstones(
sf_id,
None,
&TombstonesSidecar {
seal: None,
bitmap: bitmap.clone(),
},
)
.await
.expect("put");
cache.reconcile(view(1, &[(sf_id, 1)]));
let now = Instant::now();
assert!(cache.seal_for(sf_id, now).expect("unsealed").is_none());
assert_eq!(
cache
.bitmap_for(sf_id, now)
.expect("bitmap")
.iter()
.collect::<Vec<_>>(),
vec![3u32]
);
let (sidecar, etag) = ws
.get_tombstones(sf_id)
.await
.expect("read")
.expect("present");
let sealed = TombstonesSidecar {
seal: Some(SealRecord {
compaction_id: Uuid::from_u128(0xC0),
sealed_at: Utc::now(),
}),
bitmap: sidecar.bitmap,
};
ws.put_tombstones(sf_id, Some(&etag), &sealed)
.await
.expect("seal");
let seal = cache.seal_for(sf_id, now).expect("sealed");
assert!(seal.is_some(), "seal observed once the TTL window closed");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn sidecar_for_returns_both_bitmap_and_seal() {
let (_dir, ws, cache) = fixture();
let sf_id = Uuid::from_u128(0xABCD);
let mut bitmap = RoaringBitmap::new();
bitmap.insert(1);
bitmap.insert(2);
ws.put_tombstones(sf_id, None, &TombstonesSidecar { seal: None, bitmap })
.await
.expect("put");
cache.reconcile(view(1, &[(sf_id, 1)]));
let now = Instant::now();
let (cached_bitmap, seal) = cache.sidecar_for(sf_id, now).expect("lookup");
let collected: Vec<u32> = cached_bitmap.iter().collect();
assert_eq!(collected, vec![1u32, 2]);
assert!(seal.is_none());
}
#[test]
fn cache_is_empty_on_construction() {
let (_dir, _ws, cache) = fixture();
assert!(cache.is_empty());
assert_eq!(cache.len(), 0);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn clear_empties_the_cache() {
let (_dir, ws, cache) = fixture();
let sf_id = Uuid::from_u128(0x1234);
let mut bitmap = RoaringBitmap::new();
bitmap.insert(1);
ws.put_tombstones(sf_id, None, &TombstonesSidecar { seal: None, bitmap })
.await
.expect("put");
cache.reconcile(view(1, &[(sf_id, 1)]));
let _ = cache.bitmap_for(sf_id, Instant::now()).expect("lookup");
assert_eq!(cache.len(), 1);
cache.clear();
assert!(cache.is_empty(), "clear drops all entries");
assert_eq!(cache.len(), 0);
}
}