use std::{
collections::HashSet,
fmt, fs, io,
io::SeekFrom,
os::unix::fs::FileExt,
path::{Path, PathBuf},
sync::{
Arc, OnceLock, Weak,
atomic::{AtomicBool, AtomicU64, Ordering},
},
thread,
time::{Duration, Instant},
};
use async_trait::async_trait;
use bytes::Bytes;
use dashmap::{DashMap, mapref::entry::Entry};
use futures::{
future::try_join_all,
stream::{FuturesUnordered, StreamExt},
};
use memmap2::{Mmap, UncheckedAdvice};
use thiserror::Error;
use tokio::{
io::{AsyncSeekExt, AsyncWriteExt},
sync::{Notify, OnceCell, Semaphore, oneshot},
task::{JoinHandle, spawn_blocking},
};
use super::{
block_source::BlockCachedSource,
config::{ColdFetchMode, DiskCacheConfig, EvictionCandidate},
};
use crate::{
config::global as global_config,
runtime_metrics::io::scope_background,
storage::{StorageError, StorageProvider},
superfile::{
BytesLazyByteSource, LazyByteSource, LazyByteSourceError, PrefetchedSource,
format::{footer, kv},
reader::{OpenOptions, SuperfileReader},
},
supertable::{
StorageRangeSource,
manifest::{SubsectionOffsets, SuperfileUri},
},
};
const PARQUET_TAIL_SPEC_BYTES: u64 = 64 * 1024;
const VECTOR_OPEN_HEADER_FALLBACK_BYTES: u64 = 32;
const FTS_OPEN_HEADER_FALLBACK_BYTES: u64 = 48;
const MMAP_PROMOTION_POLL_INTERVAL: Duration = Duration::from_millis(10);
const STORE_UPGRADE_RETRY_INTERVAL: Duration = Duration::from_millis(10);
const BLOCKS_FILE_SUFFIX: &str = ".blocks";
static FOREGROUND_QUERIES: AtomicU64 = AtomicU64::new(0);
static FOREGROUND_NOTIFY: OnceLock<Notify> = OnceLock::new();
fn foreground_notify() -> &'static Notify {
FOREGROUND_NOTIFY.get_or_init(Notify::new)
}
pub struct ForegroundQueryGuard(());
impl ForegroundQueryGuard {
pub fn enter() -> Self {
FOREGROUND_QUERIES.fetch_add(1, Ordering::AcqRel);
foreground_notify().notify_waiters();
ForegroundQueryGuard(())
}
}
impl Drop for ForegroundQueryGuard {
fn drop(&mut self) {
FOREGROUND_QUERIES.fetch_sub(1, Ordering::AcqRel);
foreground_notify().notify_waiters();
}
}
fn reader_blocks_background_fill(reader: &Weak<SuperfileReader>) -> bool {
reader.strong_count() > 1
}
#[derive(Debug, Error)]
pub enum DiskCacheError {
#[error("storage error during cold fetch")]
Storage(#[from] StorageError),
#[error("local filesystem error: {0}")]
Io(#[from] std::io::Error),
#[error("superfile reader failed to open mmap'd bytes: {0}")]
SuperfileOpen(String),
#[error("superfile reader failed to open bytes")]
SuperfileOpenRead(#[from] crate::superfile::ReadError),
#[error("disk cache budget exceeded with no eligible victims")]
BudgetExceeded,
#[error("config: {0}")]
Config(String),
}
struct CachedEntry {
reader: Arc<SuperfileReader>,
mmap: Option<Arc<Mmap>>,
size_bytes: Arc<AtomicU64>,
accounting: EntryAccounting,
block_token: Option<Arc<()>>,
block_source: Option<Arc<BlockCachedSource>>,
fill_spawned: AtomicBool,
last_access_us: AtomicU64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum EntryAccounting {
Eager,
#[cfg(test)]
SourceOwned,
}
type Coordinator = Arc<OnceCell<Result<Arc<CachedEntry>, DiskCacheError>>>;
#[derive(Debug, Clone, Default)]
pub struct CacheStats {
pub n_entries: u64,
pub current_bytes: u64,
pub budget_bytes: u64,
pub n_cold_fetches: u64,
pub n_evictions: u64,
pub n_madvise_calls: u64,
}
pub struct DiskCacheStore {
storage: Arc<dyn StorageProvider>,
config: DiskCacheConfig,
started_at: Instant,
cached: DashMap<SuperfileUri, Arc<CachedEntry>>,
coordinators: DashMap<SuperfileUri, Coordinator>,
current_bytes: AtomicU64,
budget_bytes: AtomicU64,
budget_auto_sized: AtomicBool,
budget_warned: AtomicBool,
n_cold_fetches: AtomicU64,
n_evictions: AtomicU64,
n_madvise_calls: AtomicU64,
n_promotion_waiters: AtomicU64,
pinned_fn: std::sync::Mutex<Arc<dyn Fn() -> HashSet<SuperfileUri> + Send + Sync>>,
prefetch_semaphore: Arc<Semaphore>,
}
impl fmt::Debug for DiskCacheStore {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("DiskCacheStore")
.field("cache_root", &self.config.cache_root)
.field("budget_bytes", &self.disk_budget_bytes())
.field("current_bytes", &self.current_bytes.load(Ordering::Acquire))
.field("n_entries", &self.cached.len())
.field(
"n_cold_fetches",
&self.n_cold_fetches.load(Ordering::Acquire),
)
.finish()
}
}
impl DiskCacheStore {
pub fn new(
storage: Arc<dyn StorageProvider>,
config: DiskCacheConfig,
pinned_fn: Arc<dyn Fn() -> HashSet<SuperfileUri> + Send + Sync>,
) -> Result<Arc<Self>, DiskCacheError> {
if config.cold_fetch_mode == ColdFetchMode::RangeOnly {
return Err(DiskCacheError::Config(
"range_only does not currently use a disk cache; \
omit cache_dir or choose a different cold_fetch_mode"
.into(),
));
}
fs::create_dir_all(&config.cache_root)?;
let threshold_secs = config.mmap_cold_threshold_secs;
let interval_secs = config.mmap_sweep_interval_secs.max(1);
let configured_budget = config.disk_budget_bytes;
let prefetch_semaphore = Arc::new(Semaphore::new(config.prefetch_concurrency.max(1)));
let store = Arc::new(Self {
storage,
config,
started_at: Instant::now(),
cached: DashMap::new(),
coordinators: DashMap::new(),
current_bytes: AtomicU64::new(0),
budget_bytes: AtomicU64::new(configured_budget),
budget_auto_sized: AtomicBool::new(false),
budget_warned: AtomicBool::new(false),
n_cold_fetches: AtomicU64::new(0),
n_evictions: AtomicU64::new(0),
n_madvise_calls: AtomicU64::new(0),
n_promotion_waiters: AtomicU64::new(0),
pinned_fn: std::sync::Mutex::new(pinned_fn),
prefetch_semaphore,
});
store.restore_from_cache_root();
if threshold_secs > 0 {
let weak = Arc::downgrade(&store);
let _ = thread::Builder::new()
.name("infino-disk-cache-sweep".into())
.spawn(move || {
loop {
thread::sleep(Duration::from_secs(interval_secs));
match weak.upgrade() {
None => break,
Some(strong) => {
strong.sweep_once();
}
}
}
});
}
Ok(store)
}
pub fn sweep_once(&self) -> u64 {
let threshold_us = self
.config
.mmap_cold_threshold_secs
.saturating_mul(1_000_000);
let now_us = self.now_us();
let snapshot: Vec<(SuperfileUri, Arc<Mmap>, u64)> = self
.cached
.iter()
.filter_map(|e| {
let mmap = e.value().mmap.clone()?;
let last = e.value().last_access_us.load(Ordering::Acquire);
Some((*e.key(), mmap, last))
})
.collect();
let mut n_advised = 0u64;
for (_uri, mmap, last_access) in snapshot {
let idle = now_us.saturating_sub(last_access);
if idle >= threshold_us {
let _ = unsafe { mmap.unchecked_advise(UncheckedAdvice::DontNeed) };
n_advised += 1;
}
}
if n_advised > 0 {
self.n_madvise_calls.fetch_add(n_advised, Ordering::AcqRel);
}
n_advised
}
pub fn new_unpinned(
storage: Arc<dyn StorageProvider>,
config: DiskCacheConfig,
) -> Result<Arc<Self>, DiskCacheError> {
Self::new(storage, config, Arc::new(HashSet::new))
}
fn resolve_storage(
&self,
storage: Option<&Arc<dyn StorageProvider>>,
) -> Arc<dyn StorageProvider> {
storage
.map(Arc::clone)
.unwrap_or_else(|| Arc::clone(&self.storage))
}
pub async fn reader(
self: &Arc<Self>,
uri: &SuperfileUri,
) -> Result<Arc<SuperfileReader>, DiskCacheError> {
self.reader_with_hints(uri, None, None, true).await
}
pub async fn reader_with_hints(
self: &Arc<Self>,
uri: &SuperfileUri,
offsets: Option<&SubsectionOffsets>,
storage: Option<&Arc<dyn StorageProvider>>,
allow_background_fill: bool,
) -> Result<Arc<SuperfileReader>, DiskCacheError> {
match self.config.cold_fetch_mode {
ColdFetchMode::HybridWithPrefetch => self.reader_hybrid(uri, storage).await,
ColdFetchMode::RangeOnly => Err(DiskCacheError::SuperfileOpen(
"ColdFetchMode::RangeOnly bypasses the disk cache; \
construct StorageRangeSource + open_lazy directly"
.into(),
)),
ColdFetchMode::LazyForegroundWithBackgroundFill => {
self.reader_lazy_with_bg_fill_hinted(uri, offsets, storage, allow_background_fill)
.await
}
}
}
pub async fn open_range_only(
self: &Arc<Self>,
uri: &SuperfileUri,
offsets: Option<&SubsectionOffsets>,
storage: Option<&Arc<dyn StorageProvider>>,
) -> Result<Arc<SuperfileReader>, DiskCacheError> {
let fetch_storage = self.resolve_storage(storage);
let storage_uri = Self::storage_path(uri);
let range_src: Arc<dyn LazyByteSource> = match offsets {
Some(o) if o.total_size > 0 => Arc::new(StorageRangeSource::with_known_size(
fetch_storage,
storage_uri,
o.total_size,
)),
_ => Arc::new(StorageRangeSource::with_unknown_size(
fetch_storage,
storage_uri,
)),
};
let reader =
SuperfileReader::open_lazy_with(range_src, OpenOptions { verify_crc: false }).await?;
Ok(Arc::new(reader))
}
pub async fn reader_synchronous(
self: &Arc<Self>,
uri: &SuperfileUri,
) -> Result<Arc<SuperfileReader>, DiskCacheError> {
let storage = Arc::clone(&self.storage);
self.reader_synchronous_with_storage(uri, storage).await
}
pub async fn reader_synchronous_with_storage(
self: &Arc<Self>,
uri: &SuperfileUri,
fetch_storage: Arc<dyn StorageProvider>,
) -> Result<Arc<SuperfileReader>, DiskCacheError> {
if let Some(entry) = self.cached.get(uri) {
if entry.mmap.is_some() {
entry.last_access_us.store(self.now_us(), Ordering::Release);
return Ok(Arc::clone(&entry.reader));
}
drop(entry);
if let Some((_, removed)) = self.cached.remove(uri) {
self.release_entry_accounting(&removed);
}
self.coordinators.remove(uri);
let replacement = self.cold_fetch(uri, Arc::clone(&fetch_storage)).await?;
return Ok(Arc::clone(&replacement.reader));
}
let cell = self
.coordinators
.entry(*uri)
.or_insert_with(|| Arc::new(OnceCell::new()))
.clone();
let result = cell
.get_or_init(|| async { self.cold_fetch(uri, Arc::clone(&fetch_storage)).await })
.await;
match result {
Ok(entry) => {
self.coordinators.remove(uri);
Ok(Arc::clone(&entry.reader))
}
Err(_e) => {
self.coordinators.remove(uri);
Err(self
.cold_fetch(uri, Arc::clone(&fetch_storage))
.await
.err()
.unwrap_or(DiskCacheError::SuperfileOpen("cold fetch error".into())))
}
}
}
async fn reader_hybrid(
self: &Arc<Self>,
uri: &SuperfileUri,
storage: Option<&Arc<dyn StorageProvider>>,
) -> Result<Arc<SuperfileReader>, DiskCacheError> {
if let Some(entry) = self.cached.get(uri) {
entry.last_access_us.store(self.now_us(), Ordering::Release);
return Ok(Arc::clone(&entry.reader));
}
let cell = self
.coordinators
.entry(*uri)
.or_insert_with(|| Arc::new(OnceCell::new()))
.clone();
let result = cell
.get_or_init(|| async {
let fetch_storage = self.resolve_storage(storage);
self.cold_fetch_hybrid(uri, fetch_storage).await
})
.await;
match result {
Ok(entry) => Ok(Arc::clone(&entry.reader)),
Err(DiskCacheError::BudgetExceeded) => {
self.coordinators.remove(uri);
Err(DiskCacheError::BudgetExceeded)
}
Err(_) => {
self.coordinators.remove(uri);
let fetch_storage = self.resolve_storage(storage);
self.cold_fetch_hybrid(uri, fetch_storage)
.await
.map(|entry| Arc::clone(&entry.reader))
}
}
}
pub fn is_cached(&self, uri: &SuperfileUri) -> bool {
self.cached.contains_key(uri)
}
pub fn is_mmap_promoted(&self, uri: &SuperfileUri) -> bool {
self.cached
.get(uri)
.map(|e| e.mmap.is_some())
.unwrap_or(false)
}
pub async fn wait_until_mmap_promoted(
self: &Arc<Self>,
uri: &SuperfileUri,
timeout: Duration,
) -> Result<(), DiskCacheError> {
let _guard = PromotionWaitGuard::new(&self.n_promotion_waiters);
let start = Instant::now();
while start.elapsed() < timeout {
if self.is_mmap_promoted(uri) {
return Ok(());
}
tokio::time::sleep(MMAP_PROMOTION_POLL_INTERVAL).await;
}
Err(DiskCacheError::SuperfileOpen(format!(
"superfile {uri:?} not mmap-promoted within {timeout:?}"
)))
}
pub async fn wait_until_fills_settled(
self: &Arc<Self>,
timeout: Duration,
) -> Result<(), DiskCacheError> {
let _guard = PromotionWaitGuard::new(&self.n_promotion_waiters);
let start = Instant::now();
loop {
let pending = self.cached.iter().any(|entry| {
entry.value().fill_spawned.load(Ordering::Acquire) && entry.value().mmap.is_none()
});
if !pending {
return Ok(());
}
if start.elapsed() >= timeout {
return Err(DiskCacheError::SuperfileOpen(format!(
"background fills not settled within {timeout:?}"
)));
}
tokio::time::sleep(MMAP_PROMOTION_POLL_INTERVAL).await;
}
}
pub fn stats(&self) -> CacheStats {
CacheStats {
n_entries: self.cached.len() as u64,
current_bytes: self.current_bytes.load(Ordering::Acquire),
budget_bytes: self.disk_budget_bytes(),
n_cold_fetches: self.n_cold_fetches.load(Ordering::Acquire),
n_evictions: self.n_evictions.load(Ordering::Acquire),
n_madvise_calls: self.n_madvise_calls.load(Ordering::Acquire),
}
}
pub fn disk_budget_bytes(&self) -> u64 {
self.budget_bytes.load(Ordering::Acquire)
}
pub fn mark_budget_auto_sized(&self) {
self.budget_auto_sized.store(true, Ordering::Release);
}
pub fn reconcile_budget_floor(&self, floor_bytes: u64, footprint_bytes: u64) {
if self.budget_auto_sized.load(Ordering::Acquire) {
let mut current = self.budget_bytes.load(Ordering::Acquire);
while floor_bytes > current {
match self.budget_bytes.compare_exchange_weak(
current,
floor_bytes,
Ordering::AcqRel,
Ordering::Acquire,
) {
Ok(_) => break,
Err(next) => current = next,
}
}
return;
}
let budget = self.disk_budget_bytes();
if footprint_bytes > budget && !self.budget_warned.swap(true, Ordering::AcqRel) {
tracing::warn!(
"disk cache budget ({budget} B) is below the table's on-storage footprint \
({footprint_bytes} B, hidden vector index included): steady-state queries \
will evict and re-fetch. Raise ConnectOptions::with_cache_budget_bytes (or \
storage.disk_budget_bytes), or omit the budget to let the engine size it."
);
}
}
pub fn set_pinned_fn(&self, pinned_fn: Arc<dyn Fn() -> HashSet<SuperfileUri> + Send + Sync>) {
let mut g = self.pinned_fn.lock().expect("pinned_fn mutex poisoned");
*g = pinned_fn;
}
pub fn current_mmap_size_bytes(&self) -> u64 {
self.cached
.iter()
.filter_map(|e| e.value().mmap.as_ref().map(|m| m.len() as u64))
.sum()
}
pub fn sweep_for_budget(&self, budget_bytes: u64) -> u64 {
let mut total = self.current_mmap_size_bytes();
if total <= budget_bytes {
return 0;
}
let mut candidates: Vec<(SuperfileUri, Arc<Mmap>, u64, u64)> = self
.cached
.iter()
.filter_map(|e| {
let mmap = e.value().mmap.clone()?;
Some((
*e.key(),
mmap,
e.value().last_access_us.load(Ordering::Acquire),
e.value().size_bytes.load(Ordering::Acquire),
))
})
.collect();
candidates.sort_by_key(|(_, _, last, _)| *last);
let mut n_advised = 0u64;
for (_uri, mmap, _last, size) in candidates {
if total <= budget_bytes {
break;
}
let _ = unsafe { mmap.unchecked_advise(UncheckedAdvice::DontNeed) };
self.n_madvise_calls.fetch_add(1, Ordering::AcqRel);
total = total.saturating_sub(size);
n_advised += 1;
}
n_advised
}
pub fn current_pinned_uris(&self) -> HashSet<SuperfileUri> {
let f = {
let g = self.pinned_fn.lock().expect("pinned_fn mutex poisoned");
Arc::clone(&g)
};
f()
}
pub async fn insert_warm(
self: &Arc<Self>,
uri: &SuperfileUri,
bytes: Bytes,
) -> Result<(), DiskCacheError> {
if self.cached.contains_key(uri) {
return Ok(());
}
let size = bytes.len() as u64;
self.reserve_manual(size).await?;
let result: Result<Arc<CachedEntry>, DiskCacheError> = async {
let tmp = self.tmp_path(uri);
let final_path = self.cache_path(uri);
{
let mut file = tokio::fs::File::create(&tmp).await?;
file.write_all(&bytes).await?;
file.flush().await?;
}
tokio::fs::rename(&tmp, &final_path).await?;
self.open_cached_entry(&final_path, size, false)
}
.await;
let entry = match result {
Ok(e) => e,
Err(e) => {
self.current_bytes.fetch_sub(size, Ordering::Release);
return Err(e);
}
};
match self.cached.entry(*uri) {
Entry::Vacant(v) => {
v.insert(entry);
}
Entry::Occupied(_) => {
self.current_bytes.fetch_sub(size, Ordering::Release);
let _ = fs::remove_file(self.cache_path(uri));
}
}
Ok(())
}
fn now_us(&self) -> u64 {
self.started_at.elapsed().as_micros() as u64
}
fn open_cached_entry(
&self,
path: &Path,
size: u64,
verify_crc: bool,
) -> Result<Arc<CachedEntry>, DiskCacheError> {
let mmap = open_readonly_mmap(path).map_err(DiskCacheError::Io)?;
let mmap_arc = Arc::new(mmap);
let reader_bytes = Bytes::from_owner(ArcMmapOwner(Arc::clone(&mmap_arc)));
let reader = SuperfileReader::open_with(reader_bytes, OpenOptions { verify_crc })?;
Ok(Arc::new(CachedEntry {
reader: Arc::new(reader),
mmap: Some(mmap_arc),
size_bytes: Arc::new(AtomicU64::new(size)),
accounting: EntryAccounting::Eager,
block_token: None,
block_source: None,
fill_spawned: AtomicBool::new(false),
last_access_us: AtomicU64::new(self.now_us()),
}))
}
fn restore_from_cache_root(self: &Arc<Self>) {
let dir = match fs::read_dir(&self.config.cache_root) {
Ok(d) => d,
Err(_) => return,
};
for entry in dir.flatten() {
let path = entry.path();
let Some(name) = path.file_name().and_then(|n| n.to_str()) else {
continue;
};
if name.ends_with(BLOCKS_FILE_SUFFIX) {
let _ = fs::remove_file(&path);
continue;
}
let Some(uri) = SuperfileUri::from_cache_filename(name) else {
continue; };
let size = match entry.metadata() {
Ok(m) if m.len() > 0 => m.len(),
_ => continue,
};
match self.open_cached_entry(&path, size, self.config.verify_crc_on_open) {
Ok(cached_entry) => {
if self.cached.insert(uri, cached_entry).is_none() {
self.current_bytes.fetch_add(size, Ordering::Release);
}
}
Err(_) => {
let _ = fs::remove_file(&path);
}
}
}
}
fn cache_path(&self, uri: &SuperfileUri) -> PathBuf {
self.config.cache_root.join(uri.cache_filename())
}
fn blocks_path(&self, uri: &SuperfileUri) -> PathBuf {
self.config
.cache_root
.join(format!("{}{BLOCKS_FILE_SUFFIX}", uri.cache_filename()))
}
fn tmp_path(&self, uri: &SuperfileUri) -> PathBuf {
self.config.cache_root.join(uri.cache_tmp_filename())
}
fn storage_path(uri: &SuperfileUri) -> String {
uri.storage_path()
}
async fn cold_fetch_hybrid(
self: &Arc<Self>,
uri: &SuperfileUri,
fetch_storage: Arc<dyn StorageProvider>,
) -> Result<Arc<CachedEntry>, DiskCacheError> {
let storage_uri = Self::storage_path(uri);
let head = fetch_storage.head(&storage_uri).await?;
let size = head.size;
self.reserve_manual(size).await?;
let reserved_bytes = size;
let tmp = self.tmp_path(uri);
let final_path = self.cache_path(uri);
let n_streams = self.config.cold_fetch_streams.max(1) as u64;
let chunk_size = self
.config
.cold_fetch_chunk_bytes
.max(size.div_ceil(n_streams));
let n_chunks = if size == 0 {
0
} else {
size.div_ceil(chunk_size)
};
let file = tokio::fs::File::create(&tmp).await?;
file.set_len(size).await?;
let file = Arc::new(tokio::sync::Mutex::new(file));
let chunks: Arc<tokio::sync::Mutex<Vec<Option<(u64, Bytes)>>>> =
Arc::new(tokio::sync::Mutex::new(vec![None; n_chunks as usize]));
let mut fetch_handles = Vec::with_capacity(n_chunks as usize);
let mut write_handles = Vec::with_capacity(n_chunks as usize);
for i in 0..n_chunks {
let start = i * chunk_size;
let end = (start + chunk_size).min(size);
let storage = Arc::clone(&fetch_storage);
let file = Arc::clone(&file);
let chunks = Arc::clone(&chunks);
let uri_s = storage_uri.clone();
let (write_tx, write_rx) = oneshot::channel::<JoinHandle<Result<(), DiskCacheError>>>();
write_handles.push(write_rx);
fetch_handles.push(tokio::spawn(async move {
let bytes = storage.get_range(&uri_s, start..end).await?;
{
let mut guard = chunks.lock().await;
guard[i as usize] = Some((start, bytes.clone()));
}
let pwrite_handle = tokio::spawn(async move {
let mut guard = file.lock().await;
guard.seek(SeekFrom::Start(start)).await?;
guard.write_all(&bytes).await?;
Ok::<(), DiskCacheError>(())
});
let _ = write_tx.send(pwrite_handle);
Ok::<(), DiskCacheError>(())
}));
}
for h in fetch_handles {
h.await
.map_err(|e| DiskCacheError::SuperfileOpen(format!("fetch join: {e}")))??;
}
let buffer = {
let chunks_guard = chunks.lock().await;
let mut buf = vec![0u8; size as usize];
for (start, bytes) in chunks_guard.iter().flatten() {
let s = *start as usize;
let e = s + bytes.len();
buf[s..e].copy_from_slice(bytes);
}
buf
};
let foreground_bytes = Bytes::from(buffer);
let foreground_reader = SuperfileReader::open_with(
foreground_bytes,
OpenOptions {
verify_crc: self.config.verify_crc_on_open,
},
)?;
let foreground_reader = Arc::new(foreground_reader);
let entry = Arc::new(CachedEntry {
reader: Arc::clone(&foreground_reader),
mmap: None, size_bytes: Arc::new(AtomicU64::new(size)),
accounting: EntryAccounting::Eager,
block_token: None,
block_source: None,
fill_spawned: AtomicBool::new(false),
last_access_us: AtomicU64::new(self.now_us()),
});
self.n_cold_fetches.fetch_add(1, Ordering::AcqRel);
self.cached.insert(*uri, Arc::clone(&entry));
let store = Arc::clone(self);
let uri_owned = *uri;
let tmp_owned = tmp.clone();
let final_owned = final_path.clone();
let file_owned = Arc::clone(&file);
tokio::spawn(async move {
let _ = finalize_to_mmap(
store,
uri_owned,
tmp_owned,
final_owned,
file_owned,
write_handles,
size,
reserved_bytes,
)
.await;
});
Ok(entry)
}
async fn reader_lazy_with_bg_fill_hinted(
self: &Arc<Self>,
uri: &SuperfileUri,
offsets: Option<&SubsectionOffsets>,
storage: Option<&Arc<dyn StorageProvider>>,
allow_background_fill: bool,
) -> Result<Arc<SuperfileReader>, DiskCacheError> {
if let Some(entry) = self.cached.get(uri) {
entry.last_access_us.store(self.now_us(), Ordering::Release);
if allow_background_fill {
self.maybe_spawn_background_fill(uri, &entry, storage);
}
return Ok(Arc::clone(&entry.reader));
}
let cell = self
.coordinators
.entry(*uri)
.or_insert_with(|| Arc::new(OnceCell::new()))
.clone();
let result = cell
.get_or_init(|| async {
let fetch_storage = self.resolve_storage(storage);
self.cold_fetch_lazy(uri, offsets, fetch_storage).await
})
.await;
let fetch_storage = self.resolve_storage(storage);
match result {
Ok(entry) => {
if allow_background_fill {
self.maybe_spawn_background_fill(uri, entry, storage);
}
Ok(Arc::clone(&entry.reader))
}
Err(_e) => {
self.coordinators.remove(uri);
match self.cold_fetch_lazy(uri, offsets, fetch_storage).await {
Ok(entry) => {
if allow_background_fill {
self.maybe_spawn_background_fill(uri, &entry, storage);
}
Ok(Arc::clone(&entry.reader))
}
Err(e) => Err(e),
}
}
}
}
fn maybe_spawn_background_fill(
self: &Arc<Self>,
uri: &SuperfileUri,
entry: &CachedEntry,
storage: Option<&Arc<dyn StorageProvider>>,
) {
if skip_background_fill() || entry.mmap.is_some() {
return;
}
if entry
.fill_spawned
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
.is_err()
{
return;
}
let size = entry.size_bytes.load(Ordering::Acquire);
let skip_vec = vector_blob_range(&entry.reader);
let store = Arc::downgrade(self);
let reader = Arc::downgrade(&entry.reader);
let uri_owned = *uri;
let storage_uri_owned = Self::storage_path(uri);
let fetch_storage = self.resolve_storage(storage);
tokio::spawn(async move {
let _ = lazy_background_fill(
store,
reader,
uri_owned,
storage_uri_owned,
size,
size,
fetch_storage,
skip_vec,
)
.await;
});
}
async fn cold_fetch_lazy(
self: &Arc<Self>,
uri: &SuperfileUri,
offsets: Option<&SubsectionOffsets>,
fetch_storage: Arc<dyn StorageProvider>,
) -> Result<Arc<CachedEntry>, DiskCacheError> {
let storage_uri = Self::storage_path(uri);
let block_source_arc: Arc<BlockCachedSource>;
let (lazy_reader, size) = if let Some(offsets) = offsets {
let total_size = offsets.total_size;
let parquet_tail_len = PARQUET_TAIL_SPEC_BYTES.min(total_size);
let parquet_tail_start = total_size.saturating_sub(parquet_tail_len);
let vec_ranges = if !offsets.vec_open_ranges.is_empty() {
offsets.vec_open_ranges.clone()
} else {
match offsets.vec {
Some((off, len)) if len > 0 => {
vec![(off, VECTOR_OPEN_HEADER_FALLBACK_BYTES.min(len))]
}
_ => Vec::new(),
}
};
let fts_ranges = if !offsets.fts_open_ranges.is_empty() {
offsets.fts_open_ranges.clone()
} else {
match offsets.fts {
Some((off, len)) if len > 0 => {
vec![(off, FTS_OPEN_HEADER_FALLBACK_BYTES.min(len))]
}
_ => Vec::new(),
}
};
let inner: Arc<dyn LazyByteSource> = Arc::new(StorageRangeSource::with_known_size(
Arc::clone(&fetch_storage),
storage_uri.clone(),
total_size,
));
let block_source = BlockCachedSource::new_pre_reserved(
inner,
Arc::downgrade(self),
*uri,
self.blocks_path(uri),
offsets.fts,
);
block_source_arc = Arc::clone(&block_source);
let mut overlay = PrefetchedSource::new(block_source);
if !offsets.open_blob.is_empty() {
for (off, bytes) in &offsets.open_blob {
overlay.install(*off, Bytes::copy_from_slice(bytes));
}
} else {
let storage_for_parquet = Arc::clone(&fetch_storage);
let storage_for_vec = Arc::clone(&fetch_storage);
let storage_for_fts = Arc::clone(&fetch_storage);
let parquet_uri = storage_uri.clone();
let vec_uri = storage_uri.clone();
let fts_uri = storage_uri.clone();
let parquet_fut = async move {
let end = total_size;
let start = parquet_tail_start;
if end == start {
return Ok::<_, StorageError>(Bytes::new());
}
storage_for_parquet
.get_range(&parquet_uri, start..end)
.await
};
let vec_fut =
async move { fetch_hint_ranges(storage_for_vec, vec_uri, vec_ranges).await };
let fts_fut =
async move { fetch_hint_ranges(storage_for_fts, fts_uri, fts_ranges).await };
let (parquet_bytes, vec_pre, fts_pre) =
futures::try_join!(parquet_fut, vec_fut, fts_fut)?;
if !parquet_bytes.is_empty() {
overlay.install(parquet_tail_start, parquet_bytes);
}
for (off, bytes) in vec_pre {
overlay.install(off, bytes);
}
for (off, bytes) in fts_pre {
overlay.install(off, bytes);
}
}
let source: Arc<dyn LazyByteSource> = Arc::new(overlay);
let lazy_reader = SuperfileReader::open_lazy_with(
Arc::clone(&source),
OpenOptions { verify_crc: false },
)
.await?;
(lazy_reader, total_size)
} else {
let range_src: Arc<dyn LazyByteSource> =
Arc::new(StorageRangeSource::with_unknown_size(
Arc::clone(&fetch_storage),
storage_uri.clone(),
));
let block_source = BlockCachedSource::new_pre_reserved(
range_src,
Arc::downgrade(self),
*uri,
self.blocks_path(uri),
None,
);
block_source_arc = Arc::clone(&block_source);
let source: Arc<dyn LazyByteSource> = block_source;
let lazy_reader = SuperfileReader::open_lazy_with(
Arc::clone(&source),
OpenOptions { verify_crc: false },
)
.await?;
let size = source.size();
(lazy_reader, size)
};
self.reserve_manual(size).await?;
let lazy_reader = Arc::new(lazy_reader);
let block_token = block_source_arc.entry_token();
let entry = Arc::new(CachedEntry {
reader: Arc::clone(&lazy_reader),
mmap: None,
size_bytes: Arc::new(AtomicU64::new(size)),
accounting: EntryAccounting::Eager,
block_token: Some(block_token),
block_source: Some(block_source_arc),
fill_spawned: AtomicBool::new(false),
last_access_us: AtomicU64::new(self.now_us()),
});
self.n_cold_fetches.fetch_add(1, Ordering::AcqRel);
self.cached.insert(*uri, Arc::clone(&entry));
Ok(entry)
}
async fn cold_fetch(
&self,
uri: &SuperfileUri,
fetch_storage: Arc<dyn StorageProvider>,
) -> Result<Arc<CachedEntry>, DiskCacheError> {
let storage_uri = Self::storage_path(uri);
let head = fetch_storage.head(&storage_uri).await?;
let size = head.size;
let reservation = self.reserve(size).await?;
let tmp = self.tmp_path(uri);
let final_path = self.cache_path(uri);
self.cold_fetch_to_disk(&fetch_storage, &storage_uri, &tmp, size)
.await?;
tokio::fs::rename(&tmp, &final_path).await?;
let mmap = open_readonly_mmap(&final_path).map_err(DiskCacheError::Io)?;
let mmap_arc = Arc::new(mmap);
let bytes = Bytes::from_owner(ArcMmapOwner(Arc::clone(&mmap_arc)));
let reader = SuperfileReader::open_with(
bytes,
OpenOptions {
verify_crc: self.config.verify_crc_on_open,
},
)?;
let entry = Arc::new(CachedEntry {
reader: Arc::new(reader),
mmap: Some(mmap_arc),
size_bytes: Arc::new(AtomicU64::new(size)),
accounting: EntryAccounting::Eager,
block_token: None,
block_source: None,
fill_spawned: AtomicBool::new(false),
last_access_us: AtomicU64::new(self.now_us()),
});
self.cached.insert(*uri, Arc::clone(&entry));
self.n_cold_fetches.fetch_add(1, Ordering::AcqRel);
reservation.commit();
Ok(entry)
}
async fn reserve_manual(&self, bytes: u64) -> Result<(), DiskCacheError> {
loop {
let budget = self.disk_budget_bytes();
let cur = self.current_bytes.load(Ordering::Acquire);
if cur + bytes <= budget {
if self
.current_bytes
.compare_exchange_weak(cur, cur + bytes, Ordering::AcqRel, Ordering::Acquire)
.is_ok()
{
return Ok(());
}
continue;
}
let needed = (cur + bytes).saturating_sub(budget);
self.evict_at_least(needed).await?;
}
}
pub(super) async fn reserve_block_bytes(&self, bytes: u64) -> Result<(), DiskCacheError> {
self.reserve_manual(bytes).await
}
pub(super) fn release_block_bytes(&self, bytes: u64) {
self.current_bytes.fetch_sub(bytes, Ordering::Release);
}
pub(super) fn lazy_block_entry_is_current(&self, uri: &SuperfileUri, token: &Arc<()>) -> bool {
self.cached
.get(uri)
.and_then(|entry| {
entry
.block_token
.as_ref()
.map(|current| Arc::ptr_eq(current, token))
})
.unwrap_or(false)
}
fn release_entry_accounting(&self, entry: &CachedEntry) {
if entry.accounting == EntryAccounting::Eager {
self.current_bytes
.fetch_sub(entry.size_bytes.load(Ordering::Acquire), Ordering::Release);
}
}
#[cfg(test)]
pub(super) fn install_block_entry_for_test(
&self,
uri: SuperfileUri,
filled: Arc<AtomicU64>,
block_token: Arc<()>,
) {
let reader =
SuperfileReader::open(tests::tiny_superfile_bytes()).expect("tiny superfile opens");
self.cached.insert(
uri,
Arc::new(CachedEntry {
reader: Arc::new(reader),
mmap: None,
size_bytes: filled,
accounting: EntryAccounting::SourceOwned,
block_token: Some(block_token),
block_source: None,
fill_spawned: AtomicBool::new(false),
last_access_us: AtomicU64::new(self.now_us()),
}),
);
}
#[cfg(test)]
pub(super) fn remove_block_entry_for_test(&self, uri: &SuperfileUri) {
let _ = self.cached.remove(uri);
}
async fn reserve(&self, bytes: u64) -> Result<Reservation<'_>, DiskCacheError> {
loop {
let budget = self.disk_budget_bytes();
let cur = self.current_bytes.load(Ordering::Acquire);
if cur + bytes <= budget {
if self
.current_bytes
.compare_exchange_weak(cur, cur + bytes, Ordering::AcqRel, Ordering::Acquire)
.is_ok()
{
return Ok(Reservation {
store: self,
bytes,
committed: false,
});
}
continue;
}
let needed = (cur + bytes).saturating_sub(budget);
self.evict_at_least(needed).await?;
}
}
async fn evict_at_least(&self, bytes_needed: u64) -> Result<(), DiskCacheError> {
let pinned_fn = {
let g = self.pinned_fn.lock().expect("pinned_fn mutex poisoned");
Arc::clone(&g)
};
let pinned = pinned_fn();
let candidates: Vec<EvictionCandidate> = self
.cached
.iter()
.map(|e| EvictionCandidate {
uri: *e.key(),
size_bytes: e.value().size_bytes.load(Ordering::Acquire),
last_access_us: e.value().last_access_us.load(Ordering::Acquire),
})
.collect();
let victims = self
.config
.eviction
.select_for_eviction(&candidates, &pinned, bytes_needed);
if victims.is_empty() {
return Err(DiskCacheError::BudgetExceeded);
}
for uri in victims {
if let Some((_, entry)) = self.cached.remove(&uri) {
let path = self.cache_path(&uri);
let _ = fs::remove_file(&path);
self.release_entry_accounting(&entry);
self.n_evictions.fetch_add(1, Ordering::AcqRel);
}
}
Ok(())
}
async fn cold_fetch_to_disk(
&self,
fetch_storage: &Arc<dyn StorageProvider>,
storage_uri: &str,
dest_path: &Path,
size: u64,
) -> Result<(), DiskCacheError> {
let n_streams = self.config.cold_fetch_streams.max(1);
let chunk_size = self.config.cold_fetch_chunk_bytes.max(1);
let file = {
let f = fs::OpenOptions::new()
.write(true)
.create(true)
.truncate(true)
.open(dest_path)?;
f.set_len(size)?;
Arc::new(f)
};
let n_chunks = if size == 0 {
0
} else {
size.div_ceil(chunk_size)
};
let stream_sem = Arc::new(tokio::sync::Semaphore::new(n_streams));
let mut joins = Vec::with_capacity(n_chunks as usize);
for i in 0..n_chunks {
let start = i * chunk_size;
let end = (start + chunk_size).min(size);
let storage = Arc::clone(fetch_storage);
let file = Arc::clone(&file);
let uri = storage_uri.to_string();
let stream_sem = Arc::clone(&stream_sem);
joins.push(tokio::spawn(async move {
let _permit = stream_sem.acquire_owned().await.map_err(|e| {
DiskCacheError::SuperfileOpen(format!("stream semaphore closed: {e}"))
})?;
let bytes = storage.get_range(&uri, start..end).await?;
spawn_blocking(move || file.write_all_at(&bytes, start))
.await
.map_err(|e| DiskCacheError::SuperfileOpen(format!("write join: {e}")))??;
Ok::<(), DiskCacheError>(())
}));
}
for h in joins {
h.await
.map_err(|e| DiskCacheError::SuperfileOpen(format!("join error: {e}")))??;
}
spawn_blocking(move || file.sync_all())
.await
.map_err(|e| DiskCacheError::SuperfileOpen(format!("fsync join: {e}")))??;
Ok(())
}
}
struct Reservation<'a> {
store: &'a DiskCacheStore,
bytes: u64,
committed: bool,
}
impl<'a> Reservation<'a> {
fn commit(mut self) {
self.committed = true;
}
}
impl<'a> Drop for Reservation<'a> {
fn drop(&mut self) {
if !self.committed {
self.store
.current_bytes
.fetch_sub(self.bytes, Ordering::Release);
}
}
}
struct PromotionWaitGuard<'a>(&'a AtomicU64);
impl<'a> PromotionWaitGuard<'a> {
fn new(counter: &'a AtomicU64) -> Self {
counter.fetch_add(1, Ordering::AcqRel);
Self(counter)
}
}
impl Drop for PromotionWaitGuard<'_> {
fn drop(&mut self) {
self.0.fetch_sub(1, Ordering::AcqRel);
}
}
async fn finalize_to_mmap(
store: Arc<DiskCacheStore>,
uri: SuperfileUri,
tmp_path: PathBuf,
final_path: PathBuf,
file: Arc<tokio::sync::Mutex<tokio::fs::File>>,
pwrite_handles: Vec<oneshot::Receiver<JoinHandle<Result<(), DiskCacheError>>>>,
size: u64,
reserved_bytes: u64,
) -> Result<(), DiskCacheError> {
let res: Result<(), DiskCacheError> = async {
for recv in pwrite_handles {
let handle = recv
.await
.map_err(|e| DiskCacheError::SuperfileOpen(format!("pwrite handle: {e}")))?;
handle
.await
.map_err(|e| DiskCacheError::SuperfileOpen(format!("pwrite join: {e}")))??;
}
{
let mut guard = file.lock().await;
guard.flush().await?;
guard.sync_all().await?;
}
drop(file);
tokio::fs::rename(&tmp_path, &final_path).await?;
let mmap = open_readonly_mmap(&final_path)?;
let mmap_arc = Arc::new(mmap);
let bytes = Bytes::from_owner(ArcMmapOwner(Arc::clone(&mmap_arc)));
let reader = SuperfileReader::open_with(
bytes,
OpenOptions {
verify_crc: store.config.verify_crc_on_open,
},
)?;
match store.cached.entry(uri) {
Entry::Occupied(mut occ) => {
*occ.get_mut() = Arc::new(CachedEntry {
reader: Arc::new(reader),
mmap: Some(mmap_arc),
size_bytes: Arc::new(AtomicU64::new(size)),
accounting: EntryAccounting::Eager,
block_token: None,
block_source: None,
fill_spawned: AtomicBool::new(false),
last_access_us: AtomicU64::new(store.started_at.elapsed().as_micros() as u64),
});
}
Entry::Vacant(_) => {
let _ = fs::remove_file(&final_path);
}
}
store.coordinators.remove(&uri);
Ok::<(), DiskCacheError>(())
}
.await;
if res.is_err() {
if let Some((_, entry)) = store.cached.remove(&uri) {
store.release_entry_accounting(&entry);
}
store.coordinators.remove(&uri);
}
let _ = reserved_bytes;
res
}
async fn fetch_hint_ranges(
storage: Arc<dyn StorageProvider>,
storage_uri: String,
ranges: Vec<(u64, u64)>,
) -> Result<Vec<(u64, Bytes)>, StorageError> {
try_join_all(
ranges
.into_iter()
.filter(|&(_, len)| len > 0)
.map(|(off, len)| {
let storage = Arc::clone(&storage);
let storage_uri = storage_uri.clone();
async move {
let bytes = storage.get_range(&storage_uri, off..off + len).await?;
Ok::<_, StorageError>((off, bytes))
}
}),
)
.await
}
fn background_store_abandoned(store: &Arc<DiskCacheStore>) -> bool {
Arc::strong_count(store) == 1
}
async fn wait_for_lazy_foreground_release(
store: &Weak<DiskCacheStore>,
reader: &Weak<SuperfileReader>,
) -> Option<Arc<DiskCacheStore>> {
loop {
if store.strong_count() == 0 || reader.strong_count() == 0 {
return None;
}
if let Some(strong) = store.upgrade()
&& strong.n_promotion_waiters.load(Ordering::Acquire) > 0
{
return Some(strong);
}
if reader.strong_count() <= 1 {
tokio::time::sleep(STORE_UPGRADE_RETRY_INTERVAL).await;
if reader.strong_count() <= 1 {
return store.upgrade();
}
continue;
}
tokio::time::sleep(STORE_UPGRADE_RETRY_INTERVAL).await;
}
}
async fn wait_for_reader_quiescence(
store: &Arc<DiskCacheStore>,
reader: &Weak<SuperfileReader>,
) -> bool {
loop {
while reader_blocks_background_fill(reader) {
if background_store_abandoned(store) {
return false;
}
tokio::time::sleep(STORE_UPGRADE_RETRY_INTERVAL).await;
}
if reader.strong_count() == 0 {
return false;
}
tokio::time::sleep(STORE_UPGRADE_RETRY_INTERVAL).await;
if reader.strong_count() == 0 {
return false;
}
if !reader_blocks_background_fill(reader) {
return !background_store_abandoned(store);
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum BackgroundFillOutcome {
Complete,
Paused,
Abandoned,
}
async fn cold_fetch_to_disk_cancelable(
store: &Arc<DiskCacheStore>,
reader: &Weak<SuperfileReader>,
fetch_storage: &Arc<dyn StorageProvider>,
storage_uri: &str,
dest_path: &Path,
size: u64,
filled: &mut Vec<bool>,
skip_vec: Option<(u64, u64)>,
) -> Result<BackgroundFillOutcome, DiskCacheError> {
let n_streams = store.config.cold_fetch_streams.max(1);
let chunk_size = store.config.cold_fetch_chunk_bytes.max(1);
let n_chunks = if size == 0 {
0
} else {
size.div_ceil(chunk_size)
};
let first_attempt = filled.len() != n_chunks as usize;
if first_attempt {
filled.clear();
filled.resize(n_chunks as usize, false);
}
let file = {
let mut opts = fs::OpenOptions::new();
opts.write(true).create(true);
if first_attempt {
opts.truncate(true);
}
let file = opts.open(dest_path)?;
if first_attempt {
file.set_len(size)?;
}
Arc::new(file)
};
let mut next_chunk = 0u64;
let mut in_flight = FuturesUnordered::new();
loop {
while next_chunk < n_chunks && in_flight.len() < n_streams {
if filled[next_chunk as usize] {
next_chunk += 1;
continue;
}
if background_store_abandoned(store) {
return Ok(BackgroundFillOutcome::Abandoned);
}
if reader.strong_count() == 0 {
return Ok(BackgroundFillOutcome::Abandoned);
}
if reader_blocks_background_fill(reader) {
return Ok(BackgroundFillOutcome::Paused);
}
let chunk_idx = next_chunk;
let start = chunk_idx * chunk_size;
let end = (start + chunk_size).min(size);
let fetch_ranges = chunk_fetch_ranges(start, end, skip_vec);
if fetch_ranges.is_empty() {
filled[chunk_idx as usize] = true;
next_chunk += 1;
continue;
}
let storage = Arc::clone(fetch_storage);
let file = Arc::clone(&file);
let uri = storage_uri.to_string();
in_flight.push(async move {
for (range_start, range_end) in fetch_ranges {
let len = range_end - range_start;
let bytes =
scope_background(storage.get_range(&uri, range_start..range_end)).await?;
let file = Arc::clone(&file);
spawn_blocking(move || file.write_all_at(&bytes, range_start))
.await
.map_err(|error| {
DiskCacheError::SuperfileOpen(format!("write join: {error}"))
})??;
let _ = len;
}
Ok::<u64, DiskCacheError>(chunk_idx)
});
next_chunk += 1;
}
let foreground = foreground_notify().notified();
tokio::pin!(foreground);
let _ = foreground.as_mut().enable();
if reader.strong_count() == 0 {
return Ok(BackgroundFillOutcome::Abandoned);
}
if reader_blocks_background_fill(reader) {
return Ok(BackgroundFillOutcome::Paused);
}
tokio::select! {
biased;
_ = &mut foreground => {
if reader.strong_count() == 0 {
return Ok(BackgroundFillOutcome::Abandoned);
}
if reader_blocks_background_fill(reader) {
return Ok(BackgroundFillOutcome::Paused);
}
}
result = in_flight.next() => match result {
Some(result) => filled[result? as usize] = true,
None => break,
}
}
if background_store_abandoned(store) {
return Ok(BackgroundFillOutcome::Abandoned);
}
}
if background_store_abandoned(store) {
return Ok(BackgroundFillOutcome::Abandoned);
}
if reader.strong_count() == 0 {
return Ok(BackgroundFillOutcome::Abandoned);
}
if reader_blocks_background_fill(reader) {
return Ok(BackgroundFillOutcome::Paused);
}
spawn_blocking(move || file.sync_all())
.await
.map_err(|error| DiskCacheError::SuperfileOpen(format!("fsync join: {error}")))??;
Ok(BackgroundFillOutcome::Complete)
}
fn rollback_lazy_background_fill(store: &Arc<DiskCacheStore>, uri: &SuperfileUri, tmp: &Path) {
if let Some((_, entry)) = store.cached.remove(uri) {
store.release_entry_accounting(&entry);
}
store.coordinators.remove(uri);
let _ = fs::remove_file(tmp);
}
pub(crate) fn skip_background_fill() -> bool {
global_config().diagnostics.disable_background_fill
}
async fn lazy_background_fill(
store: Weak<DiskCacheStore>,
reader: Weak<SuperfileReader>,
uri: SuperfileUri,
storage_uri: String,
size: u64,
reserved_bytes: u64,
fetch_storage: Arc<dyn StorageProvider>,
skip_vec: Option<(u64, u64)>,
) -> Result<(), DiskCacheError> {
let Some(store) = wait_for_lazy_foreground_release(&store, &reader).await else {
return Ok(());
};
let tmp = store.tmp_path(&uri);
let final_path = store.cache_path(&uri);
if background_store_abandoned(&store) {
rollback_lazy_background_fill(&store, &uri, &tmp);
let _ = reserved_bytes;
return Ok(());
}
let _prefetch_permit = match Arc::clone(&store.prefetch_semaphore).acquire_owned().await {
Ok(permit) => permit,
Err(error) => {
rollback_lazy_background_fill(&store, &uri, &tmp);
return Err(DiskCacheError::SuperfileOpen(format!(
"prefetch semaphore closed: {error}"
)));
}
};
let mut filled: Vec<bool> = Vec::new();
loop {
if !wait_for_reader_quiescence(&store, &reader).await {
rollback_lazy_background_fill(&store, &uri, &tmp);
return Ok(());
}
match cold_fetch_to_disk_cancelable(
&store,
&reader,
&fetch_storage,
&storage_uri,
&tmp,
size,
&mut filled,
skip_vec,
)
.await?
{
BackgroundFillOutcome::Complete => break,
BackgroundFillOutcome::Paused => {}
BackgroundFillOutcome::Abandoned => {
rollback_lazy_background_fill(&store, &uri, &tmp);
return Ok(());
}
}
}
let result: Result<(), DiskCacheError> = async {
if background_store_abandoned(&store) {
return Ok(());
}
tokio::fs::rename(&tmp, &final_path).await?;
let mmap = open_readonly_mmap(&final_path)?;
let mmap_arc = Arc::new(mmap);
let bytes = Bytes::from_owner(ArcMmapOwner(Arc::clone(&mmap_arc)));
let prior_block = store
.cached
.get(&uri)
.and_then(|entry| entry.block_source.clone());
let (promoted_reader, block_token, block_source) = match (skip_vec, prior_block) {
(Some((vec_off, vec_len)), Some(block_source)) => {
let block_token = block_source.entry_token();
let local: Arc<dyn LazyByteSource> =
Arc::new(BytesLazyByteSource::new(bytes.clone()));
let source: Arc<dyn LazyByteSource> = Arc::new(HoleFallbackSource {
local,
hole_start: vec_off,
hole_len: vec_len,
fallback: Arc::clone(&block_source),
});
let mut reader =
SuperfileReader::open_lazy_with(source, OpenOptions { verify_crc: false })
.await?;
reader.install_resident_parquet(bytes)?;
(reader, Some(block_token), Some(block_source))
}
(Some((vec_off, vec_len)), None) => {
let remote: Arc<dyn LazyByteSource> =
Arc::new(StorageRangeSource::with_known_size(
Arc::clone(&fetch_storage),
storage_uri.clone(),
size,
));
let block_source = BlockCachedSource::new_pre_reserved(
remote,
Arc::downgrade(&store),
uri,
store.blocks_path(&uri),
None,
);
let block_token = block_source.entry_token();
let local: Arc<dyn LazyByteSource> =
Arc::new(BytesLazyByteSource::new(bytes.clone()));
let source: Arc<dyn LazyByteSource> = Arc::new(HoleFallbackSource {
local,
hole_start: vec_off,
hole_len: vec_len,
fallback: Arc::clone(&block_source),
});
let mut reader =
SuperfileReader::open_lazy_with(source, OpenOptions { verify_crc: false })
.await?;
reader.install_resident_parquet(bytes)?;
(reader, Some(block_token), Some(block_source))
}
(None, _) => {
let reader = SuperfileReader::open_with(
bytes,
OpenOptions {
verify_crc: store.config.verify_crc_on_open,
},
)?;
(reader, None, None)
}
};
match store.cached.entry(uri) {
Entry::Occupied(mut occupied) => {
*occupied.get_mut() = Arc::new(CachedEntry {
reader: Arc::new(promoted_reader),
mmap: Some(mmap_arc),
size_bytes: Arc::new(AtomicU64::new(size)),
accounting: EntryAccounting::Eager,
block_token,
block_source,
fill_spawned: AtomicBool::new(true),
last_access_us: AtomicU64::new(store.now_us()),
});
}
Entry::Vacant(_) => {
let _ = fs::remove_file(&final_path);
}
}
store.coordinators.remove(&uri);
Ok(())
}
.await;
if result.is_err() || background_store_abandoned(&store) {
rollback_lazy_background_fill(&store, &uri, &tmp);
let _ = fs::remove_file(&tmp);
}
let _ = reserved_bytes;
result
}
fn vector_blob_range(reader: &SuperfileReader) -> Option<(u64, u64)> {
let kv_map = footer::extract_kv_map(reader.parquet_metadata()).ok()?;
let off: u64 = kv_map.get(kv::VEC_OFFSET)?.parse().ok()?;
let len: u64 = kv_map.get(kv::VEC_LENGTH)?.parse().ok()?;
(len > 0).then_some((off, len))
}
fn chunk_fetch_ranges(start: u64, end: u64, skip: Option<(u64, u64)>) -> Vec<(u64, u64)> {
debug_assert!(start <= end);
let Some((hole_start, hole_len)) = skip else {
return vec![(start, end)];
};
if hole_len == 0 || start == end {
return vec![(start, end)];
}
let hole_end = hole_start.saturating_add(hole_len);
if end <= hole_start || start >= hole_end {
return vec![(start, end)];
}
let mut out = Vec::with_capacity(2);
if start < hole_start {
out.push((start, hole_start.min(end)));
}
if end > hole_end {
out.push((hole_end.max(start), end));
}
out
}
struct HoleFallbackSource {
local: Arc<dyn LazyByteSource>,
hole_start: u64,
hole_len: u64,
fallback: Arc<BlockCachedSource>,
}
impl HoleFallbackSource {
fn hole_end(&self) -> u64 {
self.hole_start.saturating_add(self.hole_len)
}
fn overlaps_hole(&self, start: u64, len: u64) -> bool {
let end = start.saturating_add(len);
end > self.hole_start && start < self.hole_end()
}
fn fully_in_hole(&self, start: u64, len: u64) -> bool {
let end = start.saturating_add(len);
start >= self.hole_start && end <= self.hole_end()
}
}
#[async_trait]
impl LazyByteSource for HoleFallbackSource {
fn size(&self) -> u64 {
self.local.size()
}
async fn range(&self, start: u64, len: u64) -> Result<Bytes, LazyByteSourceError> {
if len == 0 {
return Ok(Bytes::new());
}
if !self.overlaps_hole(start, len) {
return self.local.range(start, len).await;
}
if self.fully_in_hole(start, len) {
return self.fallback.range(start, len).await;
}
let end = start + len;
let hole_end = self.hole_end();
let mut pieces = Vec::with_capacity(3);
let mut cursor = start;
if cursor < self.hole_start {
let piece_end = self.hole_start.min(end);
pieces.push(self.local.range(cursor, piece_end - cursor).await?);
cursor = piece_end;
}
if cursor < end && cursor < hole_end {
let piece_end = hole_end.min(end);
pieces.push(self.fallback.range(cursor, piece_end - cursor).await?);
cursor = piece_end;
}
if cursor < end {
pieces.push(self.local.range(cursor, end - cursor).await?);
}
if pieces.len() == 1 {
return Ok(pieces.pop().expect("one piece"));
}
let mut out = Vec::with_capacity(len as usize);
for piece in pieces {
out.extend_from_slice(&piece);
}
Ok(Bytes::from(out))
}
fn try_get_range_sync(&self, start: u64, len: u64) -> Option<Bytes> {
if len == 0 {
return Some(Bytes::new());
}
if !self.overlaps_hole(start, len) {
return self.local.try_get_range_sync(start, len);
}
if self.fully_in_hole(start, len) {
return self.fallback.try_get_range_sync(start, len);
}
None
}
}
pub(crate) struct ArcMmapOwner(pub(crate) Arc<Mmap>);
impl AsRef<[u8]> for ArcMmapOwner {
fn as_ref(&self) -> &[u8] {
self.0.as_ref()
}
}
fn open_readonly_mmap(path: &Path) -> io::Result<Mmap> {
let file = fs::File::open(path)?;
unsafe { Mmap::map(&file) }
}
pub(crate) fn mmap_readonly_bytes(path: &Path) -> io::Result<Bytes> {
let mmap = Arc::new(open_readonly_mmap(path)?);
Ok(Bytes::from_owner(ArcMmapOwner(mmap)))
}
#[cfg(test)]
mod tests {
use std::io::Error as IoError;
use arrow_array::{LargeStringArray, RecordBatch};
use arrow_schema::{DataType, Field, Schema};
use tempfile::TempDir;
use tokio::{spawn, task::yield_now, time::timeout};
use super::*;
use crate::{
storage::LocalFsStorageProvider,
superfile::builder::{BuilderOptions, SuperfileBuilder},
test_helpers::{decimal128_id_field, decimal128_ids},
};
const PROMOTE_TIMEOUT: Duration = Duration::from_secs(10);
const FOREGROUND_GUARD_HOLD: Duration = Duration::from_millis(50);
const PREEMPT_TEST_BYTES: usize = 1 << 20;
pub(super) fn tiny_superfile_bytes() -> Bytes {
let schema = Arc::new(Schema::new(vec![
decimal128_id_field("doc_id"),
Field::new("title", DataType::LargeUtf8, false),
]));
let opts = BuilderOptions::new(schema.clone(), "doc_id", vec![], vec![], None);
let mut b = SuperfileBuilder::new(opts).expect("builder");
let ids = decimal128_ids(vec![1u64]);
let titles = LargeStringArray::from(vec!["alpha"]);
let batch =
RecordBatch::try_new(schema, vec![Arc::new(ids), Arc::new(titles)]).expect("batch");
b.add_batch(&batch, &[]).expect("add_batch");
Bytes::from(b.finish().expect("finish"))
}
fn test_store() -> (TempDir, Arc<DiskCacheStore>) {
test_store_with(|cfg| {
cfg.mmap_cold_threshold_secs = 0;
})
}
fn test_store_with(
mutate: impl FnOnce(&mut DiskCacheConfig),
) -> (TempDir, Arc<DiskCacheStore>) {
let dir = TempDir::new().expect("tempdir");
let storage: Arc<dyn StorageProvider> =
Arc::new(LocalFsStorageProvider::new(dir.path()).expect("localfs"));
let mut cfg = DiskCacheConfig {
cache_root: dir.path().join("cache"),
mmap_cold_threshold_secs: 0,
..Default::default()
};
mutate(&mut cfg);
let store = DiskCacheStore::new_unpinned(storage, cfg).expect("store");
(dir, store)
}
async fn put_superfile(store: &Arc<DiskCacheStore>, uri: &SuperfileUri, bytes: Bytes) {
store
.storage
.put_atomic(&uri.storage_path(), bytes)
.await
.expect("put superfile");
}
#[tokio::test]
async fn new_creates_cache_root() {
let (dir, store) = test_store();
assert!(dir.path().join("cache").is_dir(), "cache_root created");
let dbg = format!("{store:?}");
assert!(dbg.contains("DiskCacheStore"));
assert!(dbg.contains("n_cold_fetches"));
}
#[tokio::test]
async fn new_with_sweep_thread_enabled_spawns_and_drops_cleanly() {
let (_dir, store) = test_store_with(|cfg| {
cfg.mmap_cold_threshold_secs = 1;
cfg.mmap_sweep_interval_secs = 0; });
drop(store); }
#[tokio::test]
async fn new_unpinned_installs_empty_pinned_set() {
let (_dir, store) = test_store();
assert!(store.current_pinned_uris().is_empty());
}
#[tokio::test]
async fn stats_reflect_config_and_counters() {
let (_dir, store) = test_store_with(|cfg| {
cfg.disk_budget_bytes = 12345;
});
let s = store.stats();
assert_eq!(s.budget_bytes, 12345);
assert_eq!(s.n_entries, 0);
assert_eq!(s.current_bytes, 0);
assert_eq!(s.n_cold_fetches, 0);
assert_eq!(s.n_evictions, 0);
assert_eq!(s.n_madvise_calls, 0);
let _ = format!("{:?}", s.clone());
assert_eq!(CacheStats::default().n_entries, 0);
}
#[tokio::test]
async fn set_and_read_pinned_fn() {
let (_dir, store) = test_store();
let uri = SuperfileUri::new_v4();
store.set_pinned_fn(Arc::new(move || {
let mut s = HashSet::new();
s.insert(uri);
s
}));
let pinned = store.current_pinned_uris();
assert!(pinned.contains(&uri));
assert_eq!(pinned.len(), 1);
}
#[tokio::test]
async fn is_mmap_promoted_false_for_unknown_uri() {
let (_dir, store) = test_store();
assert!(!store.is_mmap_promoted(&SuperfileUri::new_v4()));
}
#[tokio::test]
async fn rollback_lazy_background_fill_evicts_entry_and_tmp() {
let (_dir, store) = test_store();
let uri = SuperfileUri::new_v4();
store.install_block_entry_for_test(uri, Arc::new(AtomicU64::new(0)), Arc::new(()));
assert!(
store.is_cached(&uri),
"entry must be cached before rollback"
);
let tmp = store.tmp_path(&uri);
std::fs::write(&tmp, b"partial-download-bytes").expect("seed tmp scratch file");
assert!(tmp.exists(), "tmp scratch file must exist before rollback");
rollback_lazy_background_fill(&store, &uri, &tmp);
assert!(
!store.is_cached(&uri),
"cached entry must be gone after rollback"
);
assert!(
!tmp.exists(),
"tmp scratch file must be deleted after rollback"
);
}
#[tokio::test]
async fn insert_warm_caches_and_serves_reader() {
let (_dir, store) = test_store();
let uri = SuperfileUri::new_v4();
let bytes = tiny_superfile_bytes();
let size = bytes.len() as u64;
store.insert_warm(&uri, bytes).await.expect("insert_warm");
assert!(store.is_mmap_promoted(&uri));
let s = store.stats();
assert_eq!(s.n_entries, 1);
assert_eq!(s.current_bytes, size);
assert_eq!(s.n_cold_fetches, 0);
assert_eq!(store.current_mmap_size_bytes(), size);
assert!(store.cache_path(&uri).is_file());
let _r = store.reader(&uri).await.expect("reader");
assert_eq!(store.stats().n_cold_fetches, 0);
}
#[tokio::test]
async fn insert_warm_is_idempotent() {
let (_dir, store) = test_store();
let uri = SuperfileUri::new_v4();
store
.insert_warm(&uri, tiny_superfile_bytes())
.await
.expect("first");
let before = store.stats().current_bytes;
store
.insert_warm(&uri, tiny_superfile_bytes())
.await
.expect("second");
assert_eq!(store.stats().current_bytes, before);
assert_eq!(store.stats().n_entries, 1);
}
#[tokio::test]
async fn insert_warm_rejects_unparseable_bytes() {
let (_dir, store) = test_store();
let uri = SuperfileUri::new_v4();
let err = store
.insert_warm(&uri, Bytes::from_static(b"not a superfile"))
.await
.expect_err("garbage must fail to open");
assert_eq!(store.stats().current_bytes, 0);
assert_eq!(store.stats().n_entries, 0);
let _ = format!("{err}");
let _ = format!("{err:?}");
}
#[tokio::test]
async fn insert_warm_budget_exceeded_when_too_big() {
let (_dir, store) = test_store_with(|cfg| {
cfg.disk_budget_bytes = 4; });
let uri = SuperfileUri::new_v4();
let err = store
.insert_warm(&uri, tiny_superfile_bytes())
.await
.expect_err("must exceed budget");
assert!(matches!(err, DiskCacheError::BudgetExceeded));
assert_eq!(store.stats().current_bytes, 0);
}
const TEST_TINY_BUDGET_BYTES: u64 = 4;
const TEST_RAISED_FLOOR_BYTES: u64 = 1 << 20;
#[tokio::test]
async fn auto_budget_is_raised_and_admits_previously_oversized_entry() {
let (_dir, store) = test_store_with(|cfg| {
cfg.disk_budget_bytes = TEST_TINY_BUDGET_BYTES;
});
store.mark_budget_auto_sized();
let uri = SuperfileUri::new_v4();
let err = store
.insert_warm(&uri, tiny_superfile_bytes())
.await
.expect_err("undersized budget must reject");
assert!(matches!(err, DiskCacheError::BudgetExceeded));
store.reconcile_budget_floor(TEST_RAISED_FLOOR_BYTES, TEST_RAISED_FLOOR_BYTES);
assert_eq!(store.disk_budget_bytes(), TEST_RAISED_FLOOR_BYTES);
store
.insert_warm(&uri, tiny_superfile_bytes())
.await
.expect("raised budget admits the entry");
store.reconcile_budget_floor(TEST_TINY_BUDGET_BYTES, TEST_TINY_BUDGET_BYTES);
assert_eq!(store.disk_budget_bytes(), TEST_RAISED_FLOOR_BYTES);
}
#[tokio::test]
async fn explicit_budget_is_never_changed_by_reconcile() {
let (_dir, store) = test_store_with(|cfg| {
cfg.disk_budget_bytes = TEST_TINY_BUDGET_BYTES;
});
store.reconcile_budget_floor(TEST_RAISED_FLOOR_BYTES, TEST_RAISED_FLOOR_BYTES);
store.reconcile_budget_floor(TEST_RAISED_FLOOR_BYTES, TEST_RAISED_FLOOR_BYTES);
assert_eq!(store.disk_budget_bytes(), TEST_TINY_BUDGET_BYTES);
assert_eq!(store.stats().budget_bytes, TEST_TINY_BUDGET_BYTES);
}
#[tokio::test]
async fn rebuild_index_from_cache_root_on_open() {
let dir = TempDir::new().expect("tempdir");
let cache_root = dir.path().join("cache");
let storage: Arc<dyn StorageProvider> =
Arc::new(LocalFsStorageProvider::new(dir.path()).expect("localfs"));
let uri = SuperfileUri::new_v4();
let bytes = tiny_superfile_bytes();
let size = bytes.len() as u64;
{
let cfg = DiskCacheConfig {
cache_root: cache_root.clone(),
mmap_cold_threshold_secs: 0,
..Default::default()
};
let store = DiskCacheStore::new_unpinned(Arc::clone(&storage), cfg).expect("store1");
store.insert_warm(&uri, bytes).await.expect("insert_warm");
assert!(store.cache_path(&uri).is_file());
}
let cfg2 = DiskCacheConfig {
cache_root: cache_root.clone(),
mmap_cold_threshold_secs: 0,
..Default::default()
};
let store2 = DiskCacheStore::new_unpinned(Arc::clone(&storage), cfg2).expect("store2");
let s = store2.stats();
assert_eq!(s.n_entries, 1, "rebuilt index has the cached superfile");
assert_eq!(s.current_bytes, size, "rebuilt byte accounting matches");
assert_eq!(
s.n_cold_fetches, 0,
"rebuild mmaps locally, never cold-fetches"
);
let _r = store2
.reader(&uri)
.await
.expect("reader from rebuilt index");
assert_eq!(
store2.stats().n_cold_fetches,
0,
"read served from NVMe via rebuilt index, no object-store GET"
);
}
#[tokio::test]
async fn reader_synchronous_cold_then_warm_hit() {
let (_dir, store) = test_store();
let uri = SuperfileUri::new_v4();
let bytes = tiny_superfile_bytes();
let size = bytes.len() as u64;
put_superfile(&store, &uri, bytes).await;
let _r = store.reader_synchronous(&uri).await.expect("cold");
let s = store.stats();
assert_eq!(s.n_cold_fetches, 1);
assert_eq!(s.n_entries, 1);
assert_eq!(s.current_bytes, size);
assert!(store.is_mmap_promoted(&uri));
let _r2 = store.reader_synchronous(&uri).await.expect("warm");
assert_eq!(store.stats().n_cold_fetches, 1);
}
#[tokio::test]
async fn reader_synchronous_missing_object_errors() {
let (_dir, store) = test_store();
let uri = SuperfileUri::new_v4();
let err = store.reader_synchronous(&uri).await.expect_err("no object");
let _ = format!("{err}");
assert!(store.coordinators.is_empty());
}
#[tokio::test]
async fn reader_hybrid_cold_then_stays_lazy_without_full_promotion() {
let (_dir, store) = test_store();
let uri = SuperfileUri::new_v4();
put_superfile(&store, &uri, tiny_superfile_bytes()).await;
let r = store.reader(&uri).await.expect("cold hybrid");
assert_eq!(r.n_docs(), 1);
assert_eq!(store.stats().n_cold_fetches, 1);
assert_eq!(store.stats().n_entries, 1);
assert!(!store.is_mmap_promoted(&uri));
let _r2 = store.reader(&uri).await.expect("warm");
assert_eq!(store.stats().n_cold_fetches, 1);
}
#[tokio::test]
async fn reader_hybrid_empty_object_zero_chunks() {
let (_dir, store) = test_store();
let uri = SuperfileUri::new_v4();
put_superfile(&store, &uri, Bytes::new()).await;
let err = store.reader(&uri).await.expect_err("empty not a superfile");
let _ = format!("{err}");
}
#[tokio::test]
async fn cold_fetch_uses_caller_storage_not_cache_embedded_storage() {
use crate::storage::{LocalFsStorageProvider, PrefixedStorageProvider};
let dir = TempDir::new().expect("tempdir");
let user_storage: Arc<dyn StorageProvider> =
Arc::new(LocalFsStorageProvider::new(dir.path()).expect("user root"));
let hidden_root = dir.path().join("hidden_prefix");
std::fs::create_dir_all(&hidden_root).expect("hidden root");
let hidden_storage: Arc<dyn StorageProvider> = Arc::new(PrefixedStorageProvider::new(
Arc::clone(&user_storage),
"hidden_prefix",
));
let cache = DiskCacheStore::new_unpinned(
Arc::clone(&user_storage),
DiskCacheConfig {
cache_root: dir.path().join("cache"),
mmap_cold_threshold_secs: 0,
..Default::default()
},
)
.expect("cache");
let uri = SuperfileUri::new_v4();
let bytes = tiny_superfile_bytes();
hidden_storage
.put_atomic(&uri.storage_path(), bytes.clone())
.await
.expect("put at hidden prefix");
let reader = cache
.reader_with_hints(&uri, None, Some(&hidden_storage), true)
.await
.expect("cold fetch via caller storage");
assert_eq!(reader.n_docs(), 1);
assert_eq!(cache.stats().n_cold_fetches, 1);
}
#[tokio::test]
async fn lazy_cold_fetch_uses_caller_storage_not_cache_embedded_storage() {
use crate::storage::{LocalFsStorageProvider, PrefixedStorageProvider};
let dir = TempDir::new().expect("tempdir");
let user_storage: Arc<dyn StorageProvider> =
Arc::new(LocalFsStorageProvider::new(dir.path()).expect("user root"));
let hidden_root = dir.path().join("hidden_prefix");
std::fs::create_dir_all(&hidden_root).expect("hidden root");
let hidden_storage: Arc<dyn StorageProvider> = Arc::new(PrefixedStorageProvider::new(
Arc::clone(&user_storage),
"hidden_prefix",
));
let cache = DiskCacheStore::new_unpinned(
Arc::clone(&user_storage),
DiskCacheConfig {
cache_root: dir.path().join("cache"),
cold_fetch_mode: ColdFetchMode::LazyForegroundWithBackgroundFill,
mmap_cold_threshold_secs: 0,
..Default::default()
},
)
.expect("cache");
let uri = SuperfileUri::new_v4();
let bytes = tiny_superfile_bytes();
hidden_storage
.put_atomic(&uri.storage_path(), bytes.clone())
.await
.expect("put at hidden prefix");
let reader = cache
.reader_with_hints(&uri, None, Some(&hidden_storage), true)
.await
.expect("lazy cold fetch via caller storage");
assert_eq!(reader.n_docs(), 1);
assert_eq!(cache.stats().n_cold_fetches, 1);
}
#[tokio::test]
async fn reader_synchronous_with_storage_upgrades_lazy_hidden_entry() {
use crate::storage::{LocalFsStorageProvider, PrefixedStorageProvider};
let dir = TempDir::new().expect("tempdir");
let user_storage: Arc<dyn StorageProvider> =
Arc::new(LocalFsStorageProvider::new(dir.path()).expect("user root"));
let hidden_root = dir.path().join("hidden_prefix");
std::fs::create_dir_all(&hidden_root).expect("hidden root");
let hidden_storage: Arc<dyn StorageProvider> = Arc::new(PrefixedStorageProvider::new(
Arc::clone(&user_storage),
"hidden_prefix",
));
let cache = DiskCacheStore::new_unpinned(
Arc::clone(&user_storage),
DiskCacheConfig {
cache_root: dir.path().join("cache"),
cold_fetch_mode: ColdFetchMode::LazyForegroundWithBackgroundFill,
mmap_cold_threshold_secs: 0,
..Default::default()
},
)
.expect("cache");
let uri = SuperfileUri::new_v4();
hidden_storage
.put_atomic(&uri.storage_path(), tiny_superfile_bytes())
.await
.expect("put at hidden prefix");
let lazy = cache
.reader_with_hints(&uri, None, Some(&hidden_storage), true)
.await
.expect("lazy cold fetch via caller storage");
assert!(
lazy.parquet_bytes().is_none(),
"lazy mode should not materialize full parquet bytes"
);
let eager = cache
.reader_synchronous_with_storage(&uri, Arc::clone(&hidden_storage))
.await
.expect("synchronous compaction open");
assert!(
eager.parquet_bytes().is_some(),
"compaction input must have resident parquet bytes"
);
let batch = eager
.get_record_batch(None)
.expect("compaction should read full RecordBatch");
assert_eq!(batch.num_rows(), 1);
}
#[test]
fn reader_range_only_mode_is_rejected() {
let dir = TempDir::new().expect("tempdir");
let storage: Arc<dyn StorageProvider> =
Arc::new(LocalFsStorageProvider::new(dir.path()).expect("localfs"));
let cfg = DiskCacheConfig {
cache_root: dir.path().join("cache"),
cold_fetch_mode: ColdFetchMode::RangeOnly,
..Default::default()
};
let err = DiskCacheStore::new_unpinned(storage, cfg)
.expect_err("range_only + disk cache must be rejected");
assert!(matches!(err, DiskCacheError::Config(_)));
}
#[tokio::test]
async fn open_range_only_unknown_size_reads_directly() {
let (_dir, store) = test_store();
let uri = SuperfileUri::new_v4();
put_superfile(&store, &uri, tiny_superfile_bytes()).await;
let r = store
.open_range_only(&uri, None, None)
.await
.expect("range open");
assert_eq!(r.n_docs(), 1);
assert_eq!(store.stats().n_entries, 0);
assert_eq!(store.stats().current_bytes, 0);
}
#[tokio::test]
async fn open_range_only_known_size_reads_directly() {
let (_dir, store) = test_store();
let uri = SuperfileUri::new_v4();
let bytes = tiny_superfile_bytes();
let total = bytes.len() as u64;
put_superfile(&store, &uri, bytes).await;
let offsets = SubsectionOffsets {
total_size: total,
vec: None,
fts: None,
vec_open_ranges: Vec::new(),
fts_open_ranges: Vec::new(),
open_blob: Vec::new(),
};
let r = store
.open_range_only(&uri, Some(&offsets), None)
.await
.expect("known-size range open");
assert_eq!(r.n_docs(), 1);
}
#[tokio::test]
async fn reader_lazy_unknown_size_promotes_after_release() {
let (_dir, store) = test_store_with(|cfg| {
cfg.cold_fetch_mode = ColdFetchMode::LazyForegroundWithBackgroundFill;
});
let uri = SuperfileUri::new_v4();
put_superfile(&store, &uri, tiny_superfile_bytes()).await;
let r = store.reader(&uri).await.expect("lazy cold");
assert_eq!(r.n_docs(), 1);
assert_eq!(store.stats().n_cold_fetches, 1);
drop(r);
store
.wait_until_mmap_promoted(&uri, PROMOTE_TIMEOUT)
.await
.expect("background promotion");
let r2 = store.reader(&uri).await.expect("warm mmap");
assert_eq!(store.stats().n_cold_fetches, 1);
assert!(store.is_mmap_promoted(&uri));
assert!(r2.parquet_bytes().is_some());
}
#[tokio::test]
async fn reader_lazy_with_hints_known_size_promotes_after_release() {
let (_dir, store) = test_store_with(|cfg| {
cfg.cold_fetch_mode = ColdFetchMode::LazyForegroundWithBackgroundFill;
});
let uri = SuperfileUri::new_v4();
let bytes = tiny_superfile_bytes();
let total = bytes.len() as u64;
put_superfile(&store, &uri, bytes).await;
let offsets = SubsectionOffsets {
total_size: total,
vec: None,
fts: None,
vec_open_ranges: Vec::new(),
fts_open_ranges: Vec::new(),
open_blob: Vec::new(),
};
let r = store
.reader_with_hints(&uri, Some(&offsets), None, true)
.await
.expect("lazy hinted cold");
assert_eq!(r.n_docs(), 1);
assert_eq!(store.stats().n_cold_fetches, 1);
drop(r);
store
.wait_until_mmap_promoted(&uri, PROMOTE_TIMEOUT)
.await
.expect("background promotion");
let r2 = store
.reader_with_hints(&uri, Some(&offsets), None, true)
.await
.expect("warm hinted mmap");
assert_eq!(store.stats().n_cold_fetches, 1);
assert!(store.is_mmap_promoted(&uri));
assert!(r2.parquet_bytes().is_some());
}
#[tokio::test]
async fn vector_open_skips_fill_fts_open_starts_it() {
let (_dir, store) = test_store_with(|cfg| {
cfg.cold_fetch_mode = ColdFetchMode::LazyForegroundWithBackgroundFill;
});
let uri = SuperfileUri::new_v4();
put_superfile(&store, &uri, tiny_superfile_bytes()).await;
let vector_reader = store
.reader_with_hints(&uri, None, None, false)
.await
.expect("vector lazy open");
drop(vector_reader);
tokio::time::sleep(FOREGROUND_GUARD_HOLD).await;
assert!(
!store.is_mmap_promoted(&uri),
"vector open must not spawn background fill"
);
let fts_reader = store
.reader_with_hints(&uri, None, None, true)
.await
.expect("fts lazy open");
drop(fts_reader);
store
.wait_until_mmap_promoted(&uri, PROMOTE_TIMEOUT)
.await
.expect("FTS open must start background fill");
assert!(store.is_mmap_promoted(&uri));
}
#[tokio::test]
async fn background_fill_waits_for_same_uri_reader() {
let (_dir, store) = test_store_with(|cfg| {
cfg.cold_fetch_mode = ColdFetchMode::LazyForegroundWithBackgroundFill;
});
let uri = SuperfileUri::new_v4();
put_superfile(&store, &uri, tiny_superfile_bytes()).await;
let reader = store.reader(&uri).await.expect("lazy cold");
let _foreground = ForegroundQueryGuard::enter();
tokio::time::sleep(FOREGROUND_GUARD_HOLD).await;
assert!(
!store.is_mmap_promoted(&uri),
"background promotion must yield while this URI's lazy reader is held"
);
drop(reader);
store
.wait_until_mmap_promoted(&uri, PROMOTE_TIMEOUT)
.await
.expect("promotion resumes after the URI reader is released");
}
#[tokio::test]
async fn reader_for_one_uri_does_not_pause_another_uri_fill() {
let (_dir, store) = test_store_with(|cfg| {
cfg.cold_fetch_mode = ColdFetchMode::LazyForegroundWithBackgroundFill;
cfg.prefetch_concurrency = 2;
});
let held_uri = SuperfileUri::new_v4();
let fill_uri = SuperfileUri::new_v4();
put_superfile(&store, &held_uri, tiny_superfile_bytes()).await;
put_superfile(&store, &fill_uri, tiny_superfile_bytes()).await;
let held_reader = store.reader(&held_uri).await.expect("held lazy reader");
let _fill_reader = store.reader(&fill_uri).await.expect("fill lazy reader");
drop(_fill_reader);
let _foreground = ForegroundQueryGuard::enter();
store
.wait_until_mmap_promoted(&fill_uri, PROMOTE_TIMEOUT)
.await
.expect("unrelated URI fill must proceed while another URI is held");
assert!(
!store.is_mmap_promoted(&held_uri),
"held URI must still wait for its own reader release"
);
drop(held_reader);
}
#[tokio::test]
async fn hole_fallback_source_routes_local_and_fallback_by_hole() {
use crate::superfile::lazy_source::BytesLazyByteSource;
let (_dir, store) = test_store();
let uri = SuperfileUri::new_v4();
let local: Arc<dyn LazyByteSource> =
Arc::new(BytesLazyByteSource::new(Bytes::from(vec![0xAAu8; 100])));
let remote: Arc<dyn LazyByteSource> =
Arc::new(BytesLazyByteSource::new(Bytes::from(vec![0xBBu8; 100])));
let fallback = BlockCachedSource::new_pre_reserved(
remote,
Arc::downgrade(&store),
uri,
store.blocks_path(&uri),
None,
);
let hfs = HoleFallbackSource {
local,
hole_start: 40,
hole_len: 20,
fallback,
};
assert_eq!(hfs.size(), 100, "size reflects the local (full) source");
assert_eq!(
&hfs.range(0, 10).await.expect("pre-hole")[..],
&[0xAAu8; 10]
);
assert_eq!(
&hfs.range(40, 20).await.expect("in-hole")[..],
&[0xBBu8; 20]
);
let mut want = vec![0xAAu8; 10];
want.extend_from_slice(&[0xBBu8; 20]);
want.extend_from_slice(&[0xAAu8; 10]);
assert_eq!(
&hfs.range(30, 40).await.expect("spanning")[..],
&want[..],
"spanning read stitches local + fallback + local in order",
);
assert_eq!(
hfs.try_get_range_sync(0, 10).as_deref(),
Some(&[0xAAu8; 10][..]),
"sync read outside the hole comes from local",
);
assert!(
hfs.try_get_range_sync(30, 40).is_none(),
"spanning sync read forces the async path",
);
}
#[test]
fn chunk_fetch_ranges_skips_vector_hole() {
assert_eq!(
chunk_fetch_ranges(0, 100, None),
vec![(0, 100)],
"no hole ⇒ full chunk"
);
assert_eq!(
chunk_fetch_ranges(0, 100, Some((100, 50))),
vec![(0, 100)],
"hole after chunk ⇒ full chunk"
);
assert_eq!(
chunk_fetch_ranges(0, 100, Some((0, 100))),
Vec::<(u64, u64)>::new(),
"chunk fully inside hole ⇒ no GET"
);
assert_eq!(
chunk_fetch_ranges(50, 150, Some((0, 200))),
Vec::<(u64, u64)>::new(),
"chunk fully inside larger hole ⇒ no GET"
);
assert_eq!(
chunk_fetch_ranges(0, 100, Some((40, 20))),
vec![(0, 40), (60, 100)],
"hole splits chunk into two fetch ranges"
);
assert_eq!(
chunk_fetch_ranges(0, 100, Some((80, 40))),
vec![(0, 80)],
"hole overlapping chunk end ⇒ leading fetch only"
);
assert_eq!(
chunk_fetch_ranges(0, 100, Some((0, 40))),
vec![(40, 100)],
"hole overlapping chunk start ⇒ trailing fetch only"
);
}
#[tokio::test]
async fn same_uri_reader_pauses_in_flight_background_ranges() {
let (dir, store) = test_store_with(|cfg| {
cfg.cold_fetch_streams = 1;
cfg.cold_fetch_chunk_bytes = 1;
});
let uri = SuperfileUri::new_v4();
let storage_uri = uri.storage_path();
store
.storage
.put_atomic(&storage_uri, Bytes::from(vec![7u8; PREEMPT_TEST_BYTES]))
.await
.expect("put background-fill payload");
let destination = dir.path().join("preempt.tmp");
let fill_store = Arc::clone(&store);
let fill_storage = Arc::clone(&store.storage);
let fill_destination = destination.clone();
let signal_reader = Arc::new(
SuperfileReader::open(tiny_superfile_bytes()).expect("foreground signal reader"),
);
let signal_weak = Arc::downgrade(&signal_reader);
let fill = spawn(async move {
let mut filled = Vec::new();
let outcome = cold_fetch_to_disk_cancelable(
&fill_store,
&signal_weak,
&fill_storage,
&storage_uri,
&fill_destination,
PREEMPT_TEST_BYTES as u64,
&mut filled,
None,
)
.await;
(outcome, filled)
});
timeout(PROMOTE_TIMEOUT, async {
while !destination.exists() {
yield_now().await;
}
})
.await
.expect("background fill started");
let foreground = Arc::clone(&signal_reader);
let _ = ForegroundQueryGuard::enter();
let (outcome, filled) = fill.await.expect("background task joined");
let outcome = outcome.expect("background fill returned an outcome");
assert_eq!(outcome, BackgroundFillOutcome::Paused);
assert_eq!(filled.len(), PREEMPT_TEST_BYTES);
assert!(
filled.iter().any(|&done| !done),
"a same-URI pause must leave unfinished chunks for the resume"
);
drop(foreground);
}
#[tokio::test]
async fn wait_until_mmap_promoted_times_out_for_unpromoted() {
let (_dir, store) = test_store();
let uri = SuperfileUri::new_v4();
let err = store
.wait_until_mmap_promoted(&uri, Duration::from_millis(30))
.await
.expect_err("must time out");
assert!(matches!(err, DiskCacheError::SuperfileOpen(_)));
assert_eq!(store.n_promotion_waiters.load(Ordering::Acquire), 0);
}
#[tokio::test]
async fn cold_fetch_evicts_lru_when_over_budget() {
let one = tiny_superfile_bytes();
let entry_size = one.len() as u64;
let (_dir, store) = test_store_with(move |cfg| {
cfg.disk_budget_bytes = entry_size + entry_size / 2;
});
let uri_a = SuperfileUri::new_v4();
let uri_b = SuperfileUri::new_v4();
put_superfile(&store, &uri_a, tiny_superfile_bytes()).await;
put_superfile(&store, &uri_b, tiny_superfile_bytes()).await;
store.reader_synchronous(&uri_a).await.expect("a");
store.reader_synchronous(&uri_b).await.expect("b");
assert_eq!(store.stats().n_evictions, 1);
assert!(store.cached.contains_key(&uri_b));
assert!(!store.cached.contains_key(&uri_a));
assert!(!store.cache_path(&uri_a).exists());
assert_eq!(store.stats().current_bytes, entry_size);
}
#[tokio::test]
async fn cold_fetch_budget_exceeded_with_all_pinned() {
let one = tiny_superfile_bytes();
let entry_size = one.len() as u64;
let (_dir, store) = test_store_with(move |cfg| {
cfg.disk_budget_bytes = entry_size + entry_size / 2;
});
let uri_a = SuperfileUri::new_v4();
let uri_b = SuperfileUri::new_v4();
put_superfile(&store, &uri_a, tiny_superfile_bytes()).await;
put_superfile(&store, &uri_b, tiny_superfile_bytes()).await;
store.reader_synchronous(&uri_a).await.expect("a");
store.set_pinned_fn(Arc::new(move || {
let mut s = HashSet::new();
s.insert(uri_a);
s
}));
let err = store
.reader_synchronous(&uri_b)
.await
.expect_err("no eligible victims");
assert!(matches!(err, DiskCacheError::BudgetExceeded));
assert!(store.cached.contains_key(&uri_a));
}
#[tokio::test]
async fn sweep_once_advises_idle_mmap_entries() {
let (_dir, store) = test_store();
let uri = SuperfileUri::new_v4();
store
.insert_warm(&uri, tiny_superfile_bytes())
.await
.expect("warm");
let advised = store.sweep_once();
assert_eq!(advised, 1);
assert_eq!(store.stats().n_madvise_calls, 1);
assert_eq!(store.sweep_once(), 1);
assert_eq!(store.stats().n_madvise_calls, 2);
}
#[tokio::test]
async fn sweep_once_skips_when_threshold_not_reached() {
let (_dir, store) = test_store_with(|cfg| {
cfg.mmap_cold_threshold_secs = 1_000_000;
});
let uri = SuperfileUri::new_v4();
store
.insert_warm(&uri, tiny_superfile_bytes())
.await
.expect("warm");
assert_eq!(store.sweep_once(), 0);
assert_eq!(store.stats().n_madvise_calls, 0);
}
#[tokio::test]
async fn sweep_for_budget_noop_under_budget() {
let (_dir, store) = test_store();
let uri = SuperfileUri::new_v4();
store
.insert_warm(&uri, tiny_superfile_bytes())
.await
.expect("warm");
assert_eq!(store.sweep_for_budget(u64::MAX), 0);
assert_eq!(store.stats().n_madvise_calls, 0);
}
#[tokio::test]
async fn sweep_for_budget_reclaims_oldest_first() {
let (_dir, store) = test_store();
let uri = SuperfileUri::new_v4();
store
.insert_warm(&uri, tiny_superfile_bytes())
.await
.expect("warm");
let resident = store.current_mmap_size_bytes();
assert!(resident > 0);
let advised = store.sweep_for_budget(0);
assert_eq!(advised, 1);
assert_eq!(store.stats().n_madvise_calls, 1);
}
#[tokio::test]
async fn current_mmap_size_bytes_zero_when_empty() {
let (_dir, store) = test_store();
assert_eq!(store.current_mmap_size_bytes(), 0);
}
#[tokio::test]
async fn disk_cache_error_displays_all_variants() {
let variants = [
DiskCacheError::SuperfileOpen("x".into()),
DiskCacheError::BudgetExceeded,
DiskCacheError::Io(IoError::other("boom")),
];
for v in variants {
assert!(!format!("{v}").is_empty());
assert!(!format!("{v:?}").is_empty());
}
}
}