use std::cmp::Ordering;
use std::collections::BTreeMap;
use std::io::{self, BufReader, BufWriter, Read, Seek, Write};
use std::path::{Path, PathBuf};
use std::sync::atomic::AtomicU64;
use crate::common::bitvec::BitVec;
use crate::common::fs::OneshotFile;
use crate::common::types::PointOffsetType;
use fs_err::File;
use itertools::Itertools;
use parking_lot::Mutex;
use uuid::Uuid;
use super::change::{MappingChange, read_entry, write_entry};
use crate::segment::common::operation_error::{OperationError, OperationResult};
use crate::segment::id_tracker::point_mappings::PointMappings;
use crate::segment::types::PointIdType;
const FILE_MAPPINGS: &str = "mutable_id_tracker.mappings";
pub(crate) fn mappings_path(segment_path: &Path) -> PathBuf {
segment_path.join(FILE_MAPPINGS)
}
pub(super) fn store_mapping_changes(
mappings_path: &Path,
changes: &Vec<MappingChange>,
persisted_mappings_size: &AtomicU64,
) -> OperationResult<()> {
let file = File::options()
.create(true)
.append(true)
.open(mappings_path)?;
let file_len = file
.metadata()
.map_err(|err| {
OperationError::service_error(format!(
"Failed to get ID tracker mappings file size: {err}"
))
})?
.len();
let file_start_appending = persisted_mappings_size.load(std::sync::atomic::Ordering::Relaxed);
match file_len.cmp(&file_start_appending) {
Ordering::Equal => {}
Ordering::Greater => {
file.set_len(file_start_appending)
.map_err(|err| OperationError::service_error(
format!("Failed to truncate mutable ID tracker mappings file that is too large, ignoring: {err}"),
))?;
}
Ordering::Less => {
return Err(OperationError::service_error(format!(
"Mutable ID tracker mappings file size is less than persisted mappings size, cannot append new mappings (file size: {file_len}, persisted mappings size: {file_start_appending})",
)));
}
}
let mut writer = BufWriter::new(file);
log::trace!("writing mapping changes to {mappings_path:?}: {changes:?}");
write_mapping_changes(&mut writer, changes).map_err(|err| {
OperationError::service_error(format!(
"Failed to persist ID tracker point mappings ({}): {err}",
mappings_path.display(),
))
})?;
writer.flush()?;
let mut file = writer.into_inner().map_err(|err| {
OperationError::service_error(format!(
"Failed to flush ID tracker point mappings write buffer: {err}"
))
})?;
let new_persisted_size = file.stream_position().map_err(|err| {
OperationError::service_error(format!(
"Failed to get new persisted size of ID tracker mappings: {err}"
))
})?;
file.sync_all().map_err(|err| {
OperationError::service_error(format!("Failed to fsync ID tracker point mappings: {err}"))
})?;
persisted_mappings_size.store(new_persisted_size, std::sync::atomic::Ordering::Relaxed);
Ok(())
}
fn write_mapping_changes<W: Write>(
mut writer: W,
changes: &Vec<MappingChange>,
) -> OperationResult<()> {
for &change in changes {
write_entry(&mut writer, change)?;
}
writer.flush()?;
Ok(())
}
pub(super) fn load_mappings(
mappings_path: &Path,
deferred_internal_id: Option<PointOffsetType>,
) -> OperationResult<(PointMappings, u64)> {
let file = OneshotFile::open(mappings_path)?;
let file_len = file.metadata()?.len();
let mut reader = BufReader::new(file);
let mappings = read_mappings(&mut reader, deferred_internal_id)?;
let read_to = reader.stream_position()?;
reader.into_inner().drop_cache()?;
debug_assert!(read_to <= file_len, "cannot read past the end of the file");
if read_to < file_len {
log::warn!(
"Mutable ID tracker mappings file ends with incomplete entry, removing last {} bytes and assuming automatic recovery by WAL",
file_len - read_to,
);
let file = File::options()
.write(true)
.truncate(false)
.open(mappings_path)?;
file.set_len(read_to)?;
file.sync_all()?;
}
Ok((mappings, read_to))
}
fn read_mappings_iter<R>(mut reader: R) -> impl Iterator<Item = OperationResult<MappingChange>>
where
R: Read + Seek,
{
let mut position = reader.stream_position().unwrap_or(0);
std::iter::from_fn(move || match read_entry(&mut reader) {
Ok((entry, read_bytes)) => {
position += read_bytes;
Some(Ok(entry))
}
Err(err) if err.kind() == io::ErrorKind::UnexpectedEof => {
match reader.seek(io::SeekFrom::Start(position)) {
Ok(_) => None,
Err(err) => Some(Err(err.into())),
}
}
Err(err) => Some(Err(err.into())),
})
.take_while_inclusive(|item| item.is_ok())
}
pub(super) fn read_mappings<R>(
reader: R,
deferred_internal_id: Option<PointOffsetType>,
) -> OperationResult<PointMappings>
where
R: Read + Seek,
{
let mut deleted = BitVec::new();
let mut internal_to_external: Vec<PointIdType> = Default::default();
let mut external_to_internal_num: BTreeMap<u64, PointOffsetType> = Default::default();
let mut external_to_internal_uuid: BTreeMap<Uuid, PointOffsetType> = Default::default();
for change in read_mappings_iter(reader) {
match change? {
MappingChange::Insert(external_id, internal_id) => {
if internal_id as usize >= internal_to_external.len() {
internal_to_external
.resize(internal_id as usize + 1, PointIdType::NumId(u64::MAX));
}
let replaced_external_id = internal_to_external[internal_id as usize];
internal_to_external[internal_id as usize] = external_id;
if deleted
.get(internal_id as usize)
.is_some_and(|deleted| !deleted)
{
log::warn!(
"removing duplicated external id {external_id} in internal id {replaced_external_id}",
);
debug_assert!(false, "should never have to remove");
match replaced_external_id {
PointIdType::NumId(num) => {
external_to_internal_num.remove(&num);
}
PointIdType::Uuid(uuid) => {
external_to_internal_uuid.remove(&uuid);
}
}
}
if internal_id as usize >= deleted.len() {
deleted.resize(internal_id as usize + 1, true);
}
deleted.set(internal_id as usize, false);
match external_id {
PointIdType::NumId(num) => {
external_to_internal_num.insert(num, internal_id);
}
PointIdType::Uuid(uuid) => {
external_to_internal_uuid.insert(uuid, internal_id);
}
}
}
MappingChange::Delete(external_id) => {
let internal_id = match external_id {
PointIdType::NumId(idx) => external_to_internal_num.remove(&idx),
PointIdType::Uuid(uuid) => external_to_internal_uuid.remove(&uuid),
};
let Some(internal_id) = internal_id else {
continue;
};
if (internal_id as usize) < internal_to_external.len() {
internal_to_external[internal_id as usize] = PointIdType::NumId(u64::MAX);
}
if internal_id as usize >= deleted.len() {
deleted.resize(internal_id as usize + 1, true);
}
deleted.set(internal_id as usize, true);
}
}
}
let mappings = PointMappings::new(
deleted,
internal_to_external,
external_to_internal_num,
external_to_internal_uuid,
deferred_internal_id,
);
Ok(mappings)
}
pub(super) fn reconcile_persisted_mapping_changes(
pending: &Mutex<Vec<MappingChange>>,
changes: &Vec<MappingChange>,
) {
let mut pending = pending.lock();
let count = pending
.iter()
.zip(changes)
.take_while(|(pending, persisted)| pending == persisted)
.count();
pending.drain(0..count);
}