use std::fmt;
use std::ops::RangeInclusive;
use std::path::Path;
use std::time::Instant;
use quickwit_actors::{KillSwitch, Progress};
use quickwit_config::IndexingResources;
use quickwit_metastore::checkpoint::CheckpointDelta;
use tantivy::directory::MmapDirectory;
use tantivy::merge_policy::NoMergePolicy;
use tantivy::IndexBuilder;
use crate::controlled_directory::ControlledDirectory;
use crate::models::ScratchDirectory;
use crate::new_split_id;
pub struct IndexedSplit {
pub index_id: String,
pub split_id: String,
pub replaced_split_ids: Vec<String>,
pub time_range: Option<RangeInclusive<i64>>,
pub num_docs: u64,
pub docs_size_in_bytes: u64,
pub split_date_of_birth: Instant,
pub demux_num_ops: usize,
pub checkpoint_delta: CheckpointDelta,
pub index: tantivy::Index,
pub index_writer: tantivy::IndexWriter,
pub split_scratch_directory: ScratchDirectory,
pub controlled_directory_opt: Option<ControlledDirectory>,
}
impl fmt::Debug for IndexedSplit {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("IndexedSplit")
.field("id", &self.split_id)
.field("dir", &self.split_scratch_directory.path())
.field("num_docs", &self.num_docs)
.finish()
}
}
impl IndexedSplit {
pub fn new_in_dir(
index_id: String,
scratch_directory: ScratchDirectory,
indexing_resources: IndexingResources,
index_builder: IndexBuilder,
progress: Progress,
kill_switch: KillSwitch,
) -> anyhow::Result<Self> {
let split_id = new_split_id();
let split_scratch_directory_prefix = format!("split-{}-", split_id);
let split_scratch_directory =
scratch_directory.named_temp_child(split_scratch_directory_prefix)?;
let mmap_directory = MmapDirectory::open(split_scratch_directory.path())?;
let box_mmap_directory = Box::new(mmap_directory);
let controlled_directory =
ControlledDirectory::new(box_mmap_directory, progress, kill_switch);
let index = index_builder.open_or_create(controlled_directory.clone())?;
let index_writer = index.writer_with_num_threads(
indexing_resources.num_threads,
indexing_resources.heap_size.get_bytes() as usize,
)?;
index_writer.set_merge_policy(Box::new(NoMergePolicy));
Ok(IndexedSplit {
index_id,
split_id,
replaced_split_ids: Vec::new(),
time_range: None,
demux_num_ops: 0,
docs_size_in_bytes: 0,
num_docs: 0,
split_date_of_birth: Instant::now(),
index,
index_writer,
split_scratch_directory,
checkpoint_delta: CheckpointDelta::default(),
controlled_directory_opt: Some(controlled_directory),
})
}
pub fn path(&self) -> &Path {
self.split_scratch_directory.path()
}
}
#[derive(Debug)]
pub struct IndexedSplitBatch {
pub splits: Vec<IndexedSplit>,
}