use std::io::Write;
use std::path::{Path, PathBuf};
use crate::content::node::NodeState;
use crate::error::Result;
use crate::external_sort::{RunLocation, SortBudget};
use crate::index::definition::IndexDefinition;
use crate::index::lucene::codec::segment_info::SegmentDirectory;
use crate::index::lucene::documents::binaries::BinaryTextPolicy;
use crate::index::lucene::documents::document_maker::DocumentMaker;
use crate::index::lucene::documents::facets::FacetDimension;
use crate::index::lucene::documents::rules::IndexingRules;
use crate::index::lucene::writer::{LuceneIndexWriter, WrittenIndex};
use crate::index::path_filter::PathVerdict;
use crate::progress::{ProgressObserver, Step, WorkUnit, count, observe};
use crate::segment::record::RecordIdentifier;
use crate::writer::index::lucene_directory::{DirectoryListing, OakDirectoryWriter};
use crate::writer::record_writer::{
PropertyToWrite, PropertyValuesToWrite, RecordWriter, SegmentSink,
};
const DOCUMENTS_STEP: &str = "making index documents";
const SEGMENT_STEP: &str = "writing index segments";
const COPY_STEP: &str = "copying index files";
const AFTER_LAST_DOCUMENT: &str = "lucene-reindex.after-last-document";
const AFTER_SEGMENT_FINISHED: &str = "lucene-reindex.after-segment-finished";
const MID_FILE_COPY: &str = "lucene-reindex.mid-file-copy";
#[cfg(test)]
fn probe(cutpoint: &str) -> Result<()> {
crate::writer::fault_injection::fail_if_armed(cutpoint)?;
crate::writer::fault_injection::crash_if_armed(cutpoint);
Ok(())
}
#[cfg(not(test))]
#[inline]
#[expect(
clippy::unnecessary_wraps,
reason = "the sibling this stands in for can fail, and the caller is the same either way"
)]
fn probe(_cutpoint: &str) -> Result<()> {
Ok(())
}
#[derive(Clone, Copy, PartialEq, Eq, Debug, Default)]
pub struct LuceneEstimate {
pub documents: u64,
pub nodes_visited: u64,
pub stored_bytes: u64,
pub indexed_bytes: u64,
}
#[derive(Clone, Debug)]
pub struct LuceneRebuild {
pub data_record: RecordIdentifier,
pub documents: u64,
pub nodes_visited: u64,
pub files: Vec<String>,
pub segment_bytes: u64,
pub facet_dimensions: Vec<FacetDimension>,
}
pub struct LuceneRebuildSubject<'subject> {
pub state_root: &'subject NodeState<'subject>,
pub definition: &'subject IndexDefinition,
pub rules: &'subject IndexingRules,
pub policy: &'subject BinaryTextPolicy,
pub segment_directory: &'subject Path,
pub runs: &'subject RunLocation,
pub budget: &'subject SortBudget,
}
struct WorkingDirectory {
path: PathBuf,
}
impl SegmentDirectory for WorkingDirectory {
fn write_file(
&mut self,
name: &str,
write: &mut dyn FnMut(&mut dyn Write) -> Result<()>,
) -> Result<()> {
let mut file = std::fs::File::create(self.path.join(name))?;
write(&mut file)?;
file.sync_all()?;
Ok(())
}
}
pub fn estimate_lucene_index(
state_root: &NodeState<'_>,
definition: &IndexDefinition,
rules: &IndexingRules,
observer: &mut dyn ProgressObserver,
) -> Result<LuceneEstimate> {
let mut estimate = LuceneEstimate::default();
let step = Step::new(DOCUMENTS_STEP, WorkUnit::Nodes);
observe(observer, &step, |observer| {
super::property_collector::walk_visible(state_root, |node, path| {
estimate.nodes_visited += 1;
observer.step_advanced(count(estimate.nodes_visited as usize));
if definition.path_filter.filter(path) != PathVerdict::Include {
return Ok(());
}
let Some(rule) = rules
.applicable_rule(node)
.map_err(super::property_collector::index_error_to_store_error)?
else {
return Ok(());
};
estimate.documents += 1;
for property in node.properties()? {
if property.name.starts_with(':') {
continue;
}
let Some(definition) = rule.config_of(&property.name) else {
continue;
};
if !definition.index {
continue;
}
let bytes = value_bytes(&property);
if definition.use_in_excerpt {
estimate.stored_bytes += bytes;
}
estimate.indexed_bytes += bytes;
}
Ok(())
})
})?;
Ok(estimate)
}
fn value_bytes(property: &crate::content::node::PropertyState) -> u64 {
let values = match &property.values {
crate::content::node::PropertyValues::Single(value) => std::slice::from_ref(value),
crate::content::node::PropertyValues::Multiple(values) => values,
};
values
.iter()
.map(|value| value.as_text().map_or(0, |text| text.len() as u64))
.sum()
}
pub fn rebuild_lucene_index<Sink: SegmentSink>(
subject: &LuceneRebuildSubject<'_>,
writer: &mut RecordWriter<Sink>,
observer: &mut dyn ProgressObserver,
) -> Result<LuceneRebuild> {
let &LuceneRebuildSubject {
definition,
segment_directory,
..
} = subject;
if segment_directory.exists() {
std::fs::remove_dir_all(segment_directory)?;
}
std::fs::create_dir_all(segment_directory)?;
let written = make_and_write_documents(subject, observer)?;
probe(AFTER_SEGMENT_FINISHED)?;
plant_a_stray_file(segment_directory)?;
let copied = copy_segment_into_store(
definition,
segment_directory,
&written.index,
writer,
observer,
)?;
Ok(LuceneRebuild {
data_record: copied.0,
documents: u64::try_from(written.index.document_count).unwrap_or(0),
nodes_visited: written.nodes_visited,
files: written.index.files.clone(),
segment_bytes: copied.1,
facet_dimensions: written.facet_dimensions,
})
}
#[cfg(test)]
std::thread_local! {
static STRAY_FILE: std::cell::RefCell<Option<String>> =
const { std::cell::RefCell::new(None) };
}
#[cfg(all(test, unix))]
pub(crate) fn plant_stray_file(name: Option<String>) {
STRAY_FILE.with(|cell| *cell.borrow_mut() = name);
}
#[cfg(test)]
fn plant_a_stray_file(segment_directory: &Path) -> Result<()> {
let planted = STRAY_FILE.with(|cell| cell.borrow().clone());
if let Some(name) = planted {
std::fs::write(segment_directory.join(name), b"not a segment file\n")?;
}
Ok(())
}
#[cfg(not(test))]
#[inline]
#[expect(
clippy::unnecessary_wraps,
reason = "the sibling this stands in for can fail, and the caller is the same either way"
)]
fn plant_a_stray_file(_segment_directory: &Path) -> Result<()> {
Ok(())
}
struct WrittenSegment {
index: WrittenIndex,
nodes_visited: u64,
facet_dimensions: Vec<FacetDimension>,
}
fn make_and_write_documents(
subject: &LuceneRebuildSubject<'_>,
observer: &mut dyn ProgressObserver,
) -> Result<WrittenSegment> {
let LuceneRebuildSubject {
state_root,
definition,
rules,
policy,
segment_directory,
runs,
budget,
} = subject;
let mut index_writer = LuceneIndexWriter::new(
WorkingDirectory {
path: segment_directory.to_path_buf(),
},
(*runs).clone(),
(*budget).clone(),
);
let maker = DocumentMaker::new(&definition.path, rules, (*policy).clone());
let mut nodes_visited = 0u64;
let mut documents = 0u64;
let mut dimensions: Vec<FacetDimension> = Vec::new();
let step = Step::new(DOCUMENTS_STEP, WorkUnit::Nodes);
observe(observer, &step, |observer| {
super::property_collector::walk_visible(state_root, |node, path| {
nodes_visited += 1;
observer.step_advanced(count(nodes_visited as usize));
if definition.path_filter.filter(path) != PathVerdict::Include {
return Ok(());
}
let Some(rule) = rules
.applicable_rule(node)
.map_err(super::property_collector::index_error_to_store_error)?
else {
return Ok(());
};
let made = maker
.make(node, absolute(path), rule)
.map_err(super::property_collector::index_error_to_store_error)?;
let Some(made) = made else {
return Ok(());
};
documents += 1;
for dimension in made.facet_dimensions {
if !dimensions.iter().any(|seen| seen.name == dimension.name) {
dimensions.push(dimension);
} else if dimension.multi_valued
&& let Some(seen) = dimensions
.iter_mut()
.find(|seen| seen.name == dimension.name)
{
seen.multi_valued = true;
}
}
index_writer.add_document(&made.document)
})
})?;
probe(AFTER_LAST_DOCUMENT)?;
let step = Step::new(SEGMENT_STEP, WorkUnit::IndexDocuments).with_total(documents);
let (_directory, index) = observe(observer, &step, |_| index_writer.finish())?;
Ok(WrittenSegment {
index,
nodes_visited,
facet_dimensions: dimensions,
})
}
fn absolute(path: &str) -> &str {
if path.is_empty() { "/" } else { path }
}
fn copy_segment_into_store<Sink: SegmentSink>(
definition: &IndexDefinition,
segment_directory: &Path,
written: &WrittenIndex,
writer: &mut RecordWriter<Sink>,
observer: &mut dyn ProgressObserver,
) -> Result<(RecordIdentifier, u64)> {
let listing = if definition.lucene.save_directory_listing {
DirectoryListing::Saved
} else {
DirectoryListing::Omitted
};
let mut directory = OakDirectoryWriter::new(writer, definition.lucene.blob_size, listing);
let mut bytes = 0u64;
let step = Step::new(COPY_STEP, WorkUnit::Files).with_total(written.files.len() as u64);
let copied: Result<()> = observe(observer, &step, |observer| {
for (at, name) in written.files.iter().enumerate() {
if at > 0 {
probe(MID_FILE_COPY)?;
}
let path = segment_directory.join(name);
let file = std::fs::File::open(&path)?;
bytes += file.metadata()?.len();
directory.add_file(name, file)?;
observer.step_advanced(count(at + 1));
}
Ok(())
});
copied?;
Ok((directory.finish()?, bytes))
}
pub fn lucene_definition_edits<Sink: SegmentSink>(
provider: &dyn crate::content::provider::SegmentProvider,
writer: &mut RecordWriter<Sink>,
definition_node: &NodeState<'_>,
built: &LuceneRebuild,
last_updated: Option<&str>,
edits: &mut super::definition_update::DefinitionEdits,
) -> Result<()> {
edits.property_removals.push("refresh".to_owned());
edits.property_removals.push("indexImportState".to_owned());
let version = super::definition_update::fresh_index_format_version(definition_node)?;
let value = writer.write_string(&version.to_string())?;
edits.property_replacements.push(PropertyToWrite {
name: super::definition_update::INDEX_VERSION_PROPERTY.to_owned(),
property_type: crate::PropertyType::Long,
values: PropertyValuesToWrite::Single(value),
});
if !built.facet_dimensions.is_empty() {
let facets = write_facet_configuration(writer, &built.facet_dimensions)?;
edits
.visible_children
.insert(FACETS_CHILD.to_owned(), Some(facets));
}
let status = write_status_node(writer, built, last_updated)?;
edits
.hidden_children
.push((STATUS_CHILD.to_owned(), status));
let stored = super::definition_update::clone_visible_state(provider, writer, definition_node)?;
edits.hidden_children.push((
super::definition_update::STORED_DEFINITION_CHILD.to_owned(),
stored,
));
Ok(())
}
const STATUS_CHILD: &str = ":status";
const FACETS_CHILD: &str = "facets";
fn write_status_node<Sink: SegmentSink>(
writer: &mut RecordWriter<Sink>,
built: &LuceneRebuild,
last_updated: Option<&str>,
) -> Result<RecordIdentifier> {
let now = epoch_milliseconds();
let unique = writer.write_string(&now.to_string())?;
let indexed = writer.write_string(&built.documents.to_string())?;
let completion = writer.write_string(&iso8601_of(now))?;
let updated = writer.write_string(last_updated.unwrap_or(&iso8601_of(now)))?;
let properties = vec![
PropertyToWrite {
name: "uid".to_owned(),
property_type: crate::PropertyType::String,
values: PropertyValuesToWrite::Single(unique),
},
PropertyToWrite {
name: "lastUpdated".to_owned(),
property_type: crate::PropertyType::Date,
values: PropertyValuesToWrite::Single(updated),
},
PropertyToWrite {
name: "indexedNodes".to_owned(),
property_type: crate::PropertyType::Long,
values: PropertyValuesToWrite::Single(indexed),
},
PropertyToWrite {
name: "reindexCompletionTimestamp".to_owned(),
property_type: crate::PropertyType::Date,
values: PropertyValuesToWrite::Single(completion),
},
];
writer.write_node(
None,
&[],
&crate::writer::record_writer::ChildNodesToWrite::Zero,
&properties,
)
}
fn write_facet_configuration<Sink: SegmentSink>(
writer: &mut RecordWriter<Sink>,
dimensions: &[FacetDimension],
) -> Result<RecordIdentifier> {
let mut children = Vec::new();
for dimension in dimensions {
if !dimension.multi_valued {
continue;
}
let elements: Vec<&str> = dimension
.name
.split('/')
.filter(|element| !element.is_empty())
.collect();
let Some((first, below)) = elements.split_first() else {
continue;
};
let mut record = write_multi_valued_element(writer, None)?;
for element in below.iter().rev() {
record = write_multi_valued_element(writer, Some(((*element).to_owned(), record)))?;
}
children.push(((*first).to_owned(), record));
}
writer.write_node(
Some(UNSTRUCTURED_TYPE),
&[],
&match children.as_slice() {
[] => crate::writer::record_writer::ChildNodesToWrite::Zero,
[(name, node)] => crate::writer::record_writer::ChildNodesToWrite::One {
name: name.clone(),
node: *node,
},
many => crate::writer::record_writer::ChildNodesToWrite::Many(many.to_vec()),
},
&[],
)
}
fn write_multi_valued_element<Sink: SegmentSink>(
writer: &mut RecordWriter<Sink>,
below: Option<(String, RecordIdentifier)>,
) -> Result<RecordIdentifier> {
let truth = writer.write_string("true")?;
let properties = vec![PropertyToWrite {
name: "multivalued".to_owned(),
property_type: crate::PropertyType::Boolean,
values: PropertyValuesToWrite::Single(truth),
}];
writer.write_node(
Some(UNSTRUCTURED_TYPE),
&[],
&match below {
None => crate::writer::record_writer::ChildNodesToWrite::Zero,
Some((name, node)) => {
crate::writer::record_writer::ChildNodesToWrite::One { name, node }
}
},
&properties,
)
}
const UNSTRUCTURED_TYPE: &str = "nt:unstructured";
fn epoch_milliseconds() -> i64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_or(0, |elapsed| i64::try_from(elapsed.as_millis()).unwrap_or(0))
}
fn iso8601_of(milliseconds: i64) -> String {
let (days, remainder) = (
milliseconds.div_euclid(86_400_000),
milliseconds.rem_euclid(86_400_000),
);
let (year, month, day) = civil_from_days(days);
format!(
"{year:04}-{month:02}-{day:02}T{:02}:{:02}:{:02}.{:03}Z",
remainder / 3_600_000,
remainder / 60_000 % 60,
remainder / 1_000 % 60,
remainder % 1_000
)
}
fn civil_from_days(days: i64) -> (i64, i64, i64) {
let shifted = days + 719_468;
let era = shifted.div_euclid(146_097);
let day_of_era = shifted - era * 146_097;
let year_of_era =
(day_of_era - day_of_era / 1_460 + day_of_era / 36_524 - day_of_era / 146_096) / 365;
let year = year_of_era + era * 400;
let day_of_year = day_of_era - (365 * year_of_era + year_of_era / 4 - year_of_era / 100);
let month_shifted = (5 * day_of_year + 2) / 153;
let day = day_of_year - (153 * month_shifted + 2) / 5 + 1;
let month = if month_shifted < 10 {
month_shifted + 3
} else {
month_shifted - 9
};
(year + i64::from(month <= 2), month, day)
}
#[cfg(all(test, unix))]
mod tests {
use crate::index::lucene::documents::binaries::{BinaryTextFallback, BinaryTextPolicy};
use crate::writer::fault_injection::lucene_fixture::write_lucene_reindex_fixture;
use crate::writer::fault_injection::test_support::{TestDirectory, reindex_work_directory};
use crate::writer::index::{ReindexOptions, WorkDirectory, reindex};
#[test]
fn a_stray_file_in_the_segment_directory_is_not_copied() {
let directory = TestDirectory::new("lucene-stray-file");
let store = write_lucene_reindex_fixture(&directory.path);
super::plant_stray_file(Some("stray-run-0.tmp".to_owned()));
let options = ReindexOptions::new()
.with_work_directory(WorkDirectory::OperatorNamed(reindex_work_directory(&store)))
.with_binary_text_policy(BinaryTextPolicy::new(BinaryTextFallback::Marker));
let outcome = reindex(&store, options);
super::plant_stray_file(None);
outcome.expect("the rebuild runs");
let repository = crate::store::Repository::open(&store).expect("open the repository");
let data = repository
.node_at_path("/oak:index/lucene/:data")
.expect("resolve :data")
.expect(":data exists");
let names: Vec<String> = data
.child_node_entries()
.expect("the files")
.into_iter()
.map(|(name, _)| name)
.collect();
assert!(
!names.iter().any(|name| name.starts_with("stray")),
"a file no segment names reached the store: {names:?}"
);
assert_eq!(names.len(), 5, "{names:?}");
}
}