use std::io::{self, ErrorKind};
use std::path::Path;
use super::{DiskCache, ScheduledReopen, State};
use crate::common::generic_consts::Sequential;
use crate::common::universal_io::cached_fs::FileInfo;
use crate::common::universal_io::simple_disk_cache::pipeline::REMOTE_READ_ALIGNMENT;
use crate::common::universal_io::simple_disk_cache::{DiskCacheRemote, block_aligned_fetch};
use crate::common::universal_io::{OwnedPipeline, Populate, UioResult, UniversalIoError, UniversalRead};
impl<R> DiskCache<R>
where
R: DiskCacheRemote,
{
pub(super) fn reopen_impl(&mut self) -> UioResult<()> {
if self.resolve_pending_reopen()? {
return Ok(());
}
self.schedule_reopen_with_len(None)?;
self.resolve_pending_reopen()?;
Ok(())
}
fn resolve_pending_reopen(&mut self) -> UioResult<bool> {
let State::Ready {
remote,
local,
scheduled_reopen,
} = self.state.get_mut()
else {
return Ok(false);
};
match std::mem::replace(scheduled_reopen, ScheduledReopen::No) {
ScheduledReopen::No => Ok(false),
ScheduledReopen::Resize { target_len } => {
remote.reopen()?;
local.resize(&self.local_path, target_len)?;
Ok(true)
}
ScheduledReopen::Tail {
mut pipeline,
target_len,
} => {
let fetched = pipeline.wait()?;
local.resize(&self.local_path, target_len)?;
match fetched {
Some((blocks_range, bytes)) if !bytes.is_empty() => {
unsafe { local.write_mmap_bytes(&bytes, blocks_range) }
}
Some(_) | None => {}
}
*remote = pipeline.into_inner();
Ok(true)
}
}
}
pub(super) fn schedule_reopen_impl<F: FnOnce(&Path) -> Option<FileInfo>>(
&mut self,
get_file_info: F,
) -> UioResult<()> {
let Some(file_info) = get_file_info(&self.remote_path) else {
return Err(UniversalIoError::NotFound {
path: self.remote_path.clone(),
});
};
self.schedule_reopen_with_len(Some(file_info.size))
}
pub(super) fn schedule_reopen_with_len(&mut self, known_len: Option<u64>) -> UioResult<()> {
self.init_state()?;
let State::Ready {
remote,
local,
scheduled_reopen,
} = self.state.get_mut()
else {
unreachable!("init_state drives state to Ready");
};
let local_len = local.mmap().len::<u8>()?;
let remote_len = match known_len {
Some(known_len) => known_len,
None => {
remote.reopen()?;
remote.len::<u8>()?
}
};
if remote_len < local_len {
return Err(UniversalIoError::Io(io::Error::new(
ErrorKind::UnexpectedEof,
format!(
"Reopen encountered a smaller file than expected; old_len: {local_len}, new_len: {remote_len}"
),
)));
}
if scheduled_reopen.target_len() == Some(remote_len) || remote_len == local_len {
return Ok(());
}
let new_scheduled_reopen = match self.open_options.populate {
Populate::Blocking | Populate::PreferBackground => {
let (blocks_range, byte_range) =
block_aligned_fetch(local_len..remote_len, remote_len)
.expect("the byte range is non-empty");
let new_remote = self.open_remote()?;
let mut pipeline = OwnedPipeline::new(new_remote)?;
pipeline.schedule::<Sequential>(blocks_range, byte_range, REMOTE_READ_ALIGNMENT)?;
ScheduledReopen::Tail {
pipeline,
target_len: remote_len,
}
}
Populate::Auto | Populate::No | Populate::Partial(_) => ScheduledReopen::Resize {
target_len: remote_len,
},
};
let State::Ready {
remote: _,
local: _,
scheduled_reopen,
} = self.state.get_mut()
else {
unreachable!("state was Ready above");
};
*scheduled_reopen = new_scheduled_reopen;
Ok(())
}
}