use std::collections::BTreeMap;
use std::io::{self, BufReader, BufWriter, Seek, Write};
use std::path::{Path, PathBuf};
use byteorder::{ReadBytesExt, WriteBytesExt};
use crate::common::types::PointOffsetType;
use fs_err::File;
use parking_lot::Mutex;
use crate::segment::common::operation_error::{OperationError, OperationResult};
use crate::segment::id_tracker::point_mappings::FileEndianess;
use crate::segment::types::SeqNumberType;
const FILE_VERSIONS: &str = "mutable_id_tracker.versions";
pub(super) const VERSION_ELEMENT_SIZE: u64 = size_of::<SeqNumberType>() as u64;
pub(super) fn versions_path(segment_path: &Path) -> PathBuf {
segment_path.join(FILE_VERSIONS)
}
pub(super) fn load_versions(versions_path: &Path) -> OperationResult<Vec<SeqNumberType>> {
let file = File::open(versions_path)?;
let file_len = file.metadata()?.len();
if file_len % VERSION_ELEMENT_SIZE != 0 {
log::warn!(
"Mutable ID tracker versions file has partial trailing entry, ignoring last {} bytes (will be cleaned up on next flush)",
file_len % VERSION_ELEMENT_SIZE,
);
}
let version_count = file_len / VERSION_ELEMENT_SIZE;
let mut reader = BufReader::new(file);
Ok((0..version_count)
.map(|_| reader.read_u64::<FileEndianess>())
.collect::<Result<_, _>>()?)
}
pub(super) fn store_version_changes(
versions_path: &Path,
changes: &BTreeMap<PointOffsetType, SeqNumberType>,
) -> OperationResult<()> {
if changes.is_empty() {
return Ok(());
}
let file = File::options()
.create(true)
.write(true)
.truncate(false)
.open(versions_path)?;
let file_len = file.metadata()?.len();
let valid_len = (file_len / VERSION_ELEMENT_SIZE) * VERSION_ELEMENT_SIZE;
if file_len != valid_len {
log::warn!(
"Mutable ID tracker versions file has partial trailing entry ({} extra bytes), truncating",
file_len - valid_len,
);
file.set_len(valid_len)?;
}
let mut writer = BufWriter::new(file);
write_version_changes(&mut writer, changes).map_err(|err| {
OperationError::service_error(format!(
"Failed to persist ID tracker point versions ({}): {err}",
versions_path.display(),
))
})?;
writer.flush().map_err(|err| {
OperationError::service_error(format!(
"Failed to flush ID tracker point versions write buffer: {err}",
))
})?;
let file = writer.into_inner().map_err(|err| err.into_error())?;
file.sync_all().map_err(|err| {
OperationError::service_error(format!("Failed to fsync ID tracker point versions: {err}"))
})?;
Ok(())
}
fn write_version_changes<W>(
mut writer: W,
changes: &BTreeMap<PointOffsetType, SeqNumberType>,
) -> OperationResult<()>
where
W: Write + Seek,
{
let mut position = writer.stream_position()?;
for (&internal_id, &version) in changes {
let offset = u64::from(internal_id) * VERSION_ELEMENT_SIZE;
if offset != position {
position = writer.seek(io::SeekFrom::Start(offset))?;
}
writer.write_u64::<FileEndianess>(version)?;
position += VERSION_ELEMENT_SIZE;
}
writer.flush()?;
Ok(())
}
pub(super) fn reconcile_persisted_version_changes(
pending: &Mutex<BTreeMap<PointOffsetType, SeqNumberType>>,
changes: BTreeMap<PointOffsetType, SeqNumberType>,
) {
pending.lock().retain(|point_offset, pending_version| {
changes
.get(point_offset)
.is_none_or(|persisted_version| pending_version != persisted_version)
});
}