use std::io::Write;
use crate::common::types::PointOffsetType;
use crate::common::universal_io::{MmapFile, MmapFs};
use fs_err as fs;
use tempfile::Builder;
use super::{LiveReloadResult, ReadOnlyAppendableIdTracker};
use crate::segment::id_tracker::mutable_id_tracker::MutableIdTracker;
use crate::segment::id_tracker::mutable_id_tracker::mappings_storage::mappings_path;
use crate::segment::id_tracker::mutable_id_tracker::versions_storage::versions_path;
use crate::segment::id_tracker::{IdTracker, IdTrackerRead};
use crate::segment::types::{PointIdType, SeqNumberType};
type ReadOnlyTracker = ReadOnlyAppendableIdTracker<MmapFile>;
fn flush(tracker: &MutableIdTracker) {
tracker.mapping_flusher()().unwrap();
tracker.versions_flusher()().unwrap();
}
fn insert(
tracker: &mut MutableIdTracker,
external: PointIdType,
internal: PointOffsetType,
version: SeqNumberType,
) {
tracker.set_link(external, internal).unwrap();
tracker.set_internal_version(internal, version).unwrap();
}
fn assert_in_sync(read_only: &ReadOnlyTracker, mutable: &MutableIdTracker) {
assert_eq!(
read_only.available_point_count(),
mutable.available_point_count(),
);
for internal_id in 0..mutable.total_point_count() as PointOffsetType {
assert_eq!(
read_only.is_deleted_point(internal_id),
mutable.is_deleted_point(internal_id),
"deleted state mismatch at offset {internal_id}",
);
assert_eq!(
read_only.external_id(internal_id),
mutable.external_id(internal_id),
"external id mismatch at offset {internal_id}",
);
if !mutable.is_deleted_point(internal_id) {
assert_eq!(
read_only.internal_version(internal_id),
mutable.internal_version(internal_id),
"version mismatch at offset {internal_id}",
);
}
}
}
#[test]
fn test_open_matches_mutable_tracker() {
let segment_dir = Builder::new().prefix("segment_dir").tempdir().unwrap();
let mut mutable = MutableIdTracker::open(segment_dir.path(), None).unwrap();
insert(&mut mutable, 100.into(), 0, 10);
insert(&mut mutable, 200.into(), 1, 11);
insert(&mut mutable, 300.into(), 2, 12);
mutable.drop(200.into()).unwrap();
flush(&mutable);
let read_only = ReadOnlyTracker::open(&MmapFs, segment_dir.path(), None).unwrap();
assert_in_sync(&read_only, &mutable);
assert_eq!(
read_only
.internal_id_with_behavior(100.into(), common::types::DeferredBehavior::VisibleOnly),
Some(0)
);
assert_eq!(
read_only
.internal_id_with_behavior(200.into(), common::types::DeferredBehavior::VisibleOnly),
None
);
assert_eq!(read_only.external_id(2), Some(300.into()));
}
#[test]
fn test_open_without_storage_is_empty() {
let segment_dir = Builder::new().prefix("segment_dir").tempdir().unwrap();
let read_only = ReadOnlyTracker::open(&MmapFs, segment_dir.path(), None).unwrap();
assert_eq!(read_only.available_point_count(), 0);
assert_eq!(read_only.total_point_count(), 0);
}
#[test]
fn test_open_with_missing_versions_is_empty() {
let segment_dir = Builder::new().prefix("segment_dir").tempdir().unwrap();
let mut mutable = MutableIdTracker::open(segment_dir.path(), None).unwrap();
insert(&mut mutable, 100.into(), 0, 10);
flush(&mutable);
fs::remove_file(versions_path(segment_dir.path())).unwrap();
let read_only = ReadOnlyTracker::open(&MmapFs, segment_dir.path(), None).unwrap();
assert_eq!(read_only.available_point_count(), 0);
}
#[test]
fn test_live_reload_reports_inserts_and_deletes() {
let segment_dir = Builder::new().prefix("segment_dir").tempdir().unwrap();
let mut mutable = MutableIdTracker::open(segment_dir.path(), None).unwrap();
insert(&mut mutable, 100.into(), 0, 10);
insert(&mut mutable, 200.into(), 1, 11);
insert(&mut mutable, 300.into(), 2, 12);
flush(&mutable);
let mut read_only = ReadOnlyTracker::open(&MmapFs, segment_dir.path(), None).unwrap();
assert_in_sync(&read_only, &mutable);
assert_eq!(
read_only.live_reload().unwrap(),
LiveReloadResult::default(),
);
insert(&mut mutable, 400.into(), 3, 13);
insert(&mut mutable, 500.into(), 4, 14);
mutable.drop(200.into()).unwrap();
flush(&mutable);
let result = read_only.live_reload().unwrap();
assert_eq!(result.inserted, vec![3, 4]);
assert_eq!(result.deleted, vec![1]);
assert_in_sync(&read_only, &mutable);
assert!(read_only.is_deleted_point(1));
assert_eq!(read_only.internal_version(3), Some(13));
assert_eq!(read_only.internal_version(4), Some(14));
}
#[test]
fn test_live_reload_insert_then_delete_within_batch() {
let segment_dir = Builder::new().prefix("segment_dir").tempdir().unwrap();
let mut mutable = MutableIdTracker::open(segment_dir.path(), None).unwrap();
insert(&mut mutable, 100.into(), 0, 10);
flush(&mutable);
let mut read_only = ReadOnlyTracker::open(&MmapFs, segment_dir.path(), None).unwrap();
insert(&mut mutable, 200.into(), 1, 11);
mutable.drop(200.into()).unwrap();
flush(&mutable);
let result = read_only.live_reload().unwrap();
assert_eq!(result.inserted, Vec::<PointOffsetType>::new());
assert_eq!(result.deleted, Vec::<PointOffsetType>::new());
assert_in_sync(&read_only, &mutable);
assert_eq!(
read_only
.internal_id_with_behavior(100.into(), common::types::DeferredBehavior::VisibleOnly),
Some(0)
);
assert!(read_only.is_deleted_point(1));
}
#[test]
fn test_live_reload_upsert_relinks_to_new_offset() {
let segment_dir = Builder::new().prefix("segment_dir").tempdir().unwrap();
let mut mutable = MutableIdTracker::open(segment_dir.path(), None).unwrap();
insert(&mut mutable, 100.into(), 0, 10);
flush(&mutable);
let mut read_only = ReadOnlyTracker::open(&MmapFs, segment_dir.path(), None).unwrap();
assert_eq!(
read_only
.internal_id_with_behavior(100.into(), common::types::DeferredBehavior::VisibleOnly),
Some(0)
);
mutable.set_link(100.into(), 1).unwrap();
mutable.set_internal_version(1, 20).unwrap();
flush(&mutable);
let result = read_only.live_reload().unwrap();
assert_eq!(result.inserted, vec![1]);
assert_eq!(result.deleted, vec![0]);
assert_eq!(
read_only
.internal_id_with_behavior(100.into(), common::types::DeferredBehavior::VisibleOnly),
Some(1)
);
assert_eq!(read_only.internal_version(1), Some(20));
assert!(read_only.is_deleted_point(0));
assert_in_sync(&read_only, &mutable);
}
#[test]
fn test_live_reload_withholds_insert_until_version_present() {
let segment_dir = Builder::new().prefix("segment_dir").tempdir().unwrap();
let mut mutable = MutableIdTracker::open(segment_dir.path(), None).unwrap();
insert(&mut mutable, 100.into(), 0, 10);
flush(&mutable);
let mut read_only = ReadOnlyTracker::open(&MmapFs, segment_dir.path(), None).unwrap();
insert(&mut mutable, 200.into(), 1, 11);
mutable.mapping_flusher()().unwrap();
let result = read_only.live_reload().unwrap();
assert_eq!(result, LiveReloadResult::default());
assert_eq!(
read_only
.internal_id_with_behavior(200.into(), common::types::DeferredBehavior::VisibleOnly),
None
);
assert_eq!(read_only.internal_version(1), None);
mutable.versions_flusher()().unwrap();
let result = read_only.live_reload().unwrap();
assert_eq!(result.inserted, vec![1]);
assert_eq!(result.deleted, Vec::<PointOffsetType>::new());
assert_eq!(
read_only
.internal_id_with_behavior(200.into(), common::types::DeferredBehavior::VisibleOnly),
Some(1)
);
assert_eq!(read_only.internal_version(1), Some(11));
assert_in_sync(&read_only, &mutable);
}
#[test]
fn test_live_reload_ignores_partial_trailing_mapping_entry() {
let segment_dir = Builder::new().prefix("segment_dir").tempdir().unwrap();
let mut mutable = MutableIdTracker::open(segment_dir.path(), None).unwrap();
insert(&mut mutable, 100.into(), 0, 10);
insert(&mut mutable, 200.into(), 1, 11);
flush(&mutable);
let mut read_only = ReadOnlyTracker::open(&MmapFs, segment_dir.path(), None).unwrap();
assert_in_sync(&read_only, &mutable);
let complete_len = read_only.mappings_read_to;
assert!(complete_len > 0);
let mappings_file_path = mappings_path(segment_dir.path());
{
let mut file = fs::OpenOptions::new()
.append(true)
.open(&mappings_file_path)
.unwrap();
file.write_all(&[1, 0, 0, 0, 0]).unwrap();
}
let result = read_only.live_reload().unwrap();
assert_eq!(result, LiveReloadResult::default());
assert_eq!(
read_only.mappings_read_to, complete_len,
"must not consume the partial trailing entry",
);
#[cfg(not(target_os = "windows"))]
{
insert(&mut mutable, 300.into(), 2, 12);
flush(&mutable);
let result = read_only.live_reload().unwrap();
assert_eq!(result.inserted, vec![2]);
assert_eq!(result.deleted, Vec::<PointOffsetType>::new());
assert_in_sync(&read_only, &mutable);
}
}
#[test]
fn test_open_and_reload_ignore_partial_trailing_version() {
let segment_dir = Builder::new().prefix("segment_dir").tempdir().unwrap();
let mut mutable = MutableIdTracker::open(segment_dir.path(), None).unwrap();
insert(&mut mutable, 100.into(), 0, 10);
insert(&mut mutable, 200.into(), 1, 11);
flush(&mutable);
let versions_file_path = versions_path(segment_dir.path());
{
let mut file = fs::OpenOptions::new()
.append(true)
.open(&versions_file_path)
.unwrap();
file.write_all(&[1, 2, 3]).unwrap();
}
let read_only = ReadOnlyTracker::open(&MmapFs, segment_dir.path(), None).unwrap();
assert_eq!(read_only.internal_version(0), Some(10));
assert_eq!(read_only.internal_version(1), Some(11));
assert_in_sync(&read_only, &mutable);
}
#[test]
fn test_live_reload_withholds_partially_written_version() {
let segment_dir = Builder::new().prefix("segment_dir").tempdir().unwrap();
let mut mutable = MutableIdTracker::open(segment_dir.path(), None).unwrap();
insert(&mut mutable, 100.into(), 0, 10);
insert(&mut mutable, 200.into(), 1, 11);
flush(&mutable);
let mut read_only = ReadOnlyTracker::open(&MmapFs, segment_dir.path(), None).unwrap();
insert(&mut mutable, 300.into(), 2, 12);
mutable.mapping_flusher()().unwrap();
{
let mut file = fs::OpenOptions::new()
.append(true)
.open(versions_path(segment_dir.path()))
.unwrap();
file.write_all(&[1, 2, 3]).unwrap();
}
let result = read_only.live_reload().unwrap();
assert_eq!(result.inserted, Vec::<PointOffsetType>::new());
assert_eq!(read_only.internal_version(2), None);
#[cfg(not(target_os = "windows"))]
{
mutable.versions_flusher()().unwrap();
let result = read_only.live_reload().unwrap();
assert_eq!(result.inserted, vec![2]);
assert_eq!(
read_only.internal_id_with_behavior(
300.into(),
common::types::DeferredBehavior::VisibleOnly
),
Some(2)
);
assert_eq!(read_only.internal_version(2), Some(12));
assert_in_sync(&read_only, &mutable);
}
}
#[test]
fn test_merge_accumulates_unapplied_delta() {
let mut pending = LiveReloadResult {
inserted: vec![5, 1, 3],
deleted: vec![10, 2],
};
pending.merge(LiveReloadResult {
inserted: vec![7],
deleted: vec![8, 2],
});
assert_eq!(pending.inserted, vec![1, 3, 5, 7]);
assert_eq!(pending.deleted, vec![2, 8, 10]);
}
#[test]
fn test_merge_cancels_insert_then_delete() {
let mut pending = LiveReloadResult {
inserted: vec![3, 4],
deleted: vec![],
};
pending.merge(LiveReloadResult {
inserted: vec![],
deleted: vec![4],
});
assert_eq!(pending.inserted, vec![3]);
assert_eq!(pending.deleted, vec![4]);
assert!(!pending.is_empty());
}
#[test]
fn test_merge_empty_is_noop() {
let mut pending = LiveReloadResult {
inserted: vec![1],
deleted: vec![2],
};
pending.merge(LiveReloadResult::default());
assert_eq!(pending.inserted, vec![1]);
assert_eq!(pending.deleted, vec![2]);
let mut empty = LiveReloadResult::default();
assert!(empty.is_empty());
empty.merge(LiveReloadResult::default());
assert!(empty.is_empty());
}