use std::collections::BTreeMap;
use std::path::{Path, PathBuf};
use crate::error::{Error, Result};
use crate::progress::{ProgressObserver, Step, WorkUnit, observe};
use crate::segment::record::RecordIdentifier;
use crate::writer::index::counter_builder::CounterBuilder;
use crate::writer::index::definition_update::{
DefinitionEdits, disabler_verdict, rewrite_definition,
};
use crate::writer::index::plan::{NoWorkReason, ReindexAction};
use crate::writer::index::prepared::{PreparedReindex, RUN_DIRECTORY_PREFIX};
use crate::writer::index::property_builder::{MirrorBuilder, UniqueBuilder};
use crate::writer::index::property_collector::{
CollectedEntries, CollectedReferences, EntrySink, PropertyCollector, ReferenceCollector,
SortedEntries,
};
use crate::writer::index::selection::IndexingState;
use crate::writer::index::{IndexEntry, RunLocation, SortBudget};
use crate::writer::record_writer::{
PropertyToWrite, PropertyValuesToWrite, RecordWriter, SegmentSink,
};
use crate::writer::store_writer::WritableRepository;
const WRITE_STEP: &str = "writing index records";
const VERIFY_STEP: &str = "verifying the rebuilt indexes";
const BEFORE_SPILL_CLEANUP: &str = "index-reindex.before-spill-cleanup";
const BEFORE_HEAD_PUBLISH: &str = "index-reindex.before-head-publish";
const AFTER_HEAD_PUBLISH_BEFORE_FLUSH: &str = "index-reindex.after-head-publish-before-flush";
const BEFORE_APPLIED_VERIFICATION: &str = "index-reindex.before-applied-verification";
#[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<()> {
let _ = cutpoint;
Ok(())
}
pub struct RebuiltDefinition {
pub path: String,
pub previous_record: RecordIdentifier,
pub rebuilt_record: RecordIdentifier,
pub report: DefinitionReport,
pub verification: Verification,
}
#[derive(Clone, PartialEq, Eq, Debug)]
pub enum Verification {
NodeTreeOnly,
Entries {
state_root: RecordIdentifier,
entries: u64,
},
LuceneSegment {
documents: u64,
files: Vec<String>,
},
Counter {
state_root: RecordIdentifier,
credited_by_path: BTreeMap<String, i64>,
},
}
#[derive(Clone, PartialEq, Eq, Debug)]
#[non_exhaustive]
pub enum DefinitionReport {
Rebuilt {
entries: u64,
distinct_keys: u64,
nodes_written: u64,
},
RebuiltIndex {
documents: u64,
nodes_visited: u64,
files: Vec<String>,
segment_bytes: u64,
},
Reset {
removed_hidden_children: Vec<String>,
retained_hidden_children: Vec<String>,
},
NothingToDo {
reason: NoWorkReason,
},
}
#[derive(Clone, PartialEq, Eq, Debug)]
#[non_exhaustive]
pub struct ReindexOutcome {
pub definitions: Vec<(String, DefinitionReport)>,
pub head_before: RecordIdentifier,
pub head_after: RecordIdentifier,
}
impl ReindexOutcome {
#[must_use]
pub fn moved_the_head(&self) -> bool {
self.head_before != self.head_after
}
}
pub(crate) fn apply_prepared(
prepared: &PreparedReindex,
observer: &mut dyn ProgressObserver,
) -> Result<ReindexOutcome> {
if prepared.plan.is_empty() {
let repository = crate::store::Repository::open(&prepared.directory)?;
let head = repository.head_record_identifier();
return Ok(ReindexOutcome {
definitions: prepared
.plan
.actions
.iter()
.filter_map(|action| match action {
ReindexAction::NothingToDo { path, reason } => Some((
path.clone(),
DefinitionReport::NothingToDo { reason: *reason },
)),
_ => None,
})
.collect(),
head_before: head,
head_after: head,
});
}
prepared.recheck_before_mutation()?;
let run = RunDirectory::create(prepared)?;
let outcome = apply_under_the_lock(prepared, &run, observer);
run.remove();
outcome
}
struct RunDirectory {
path: PathBuf,
_lock: crate::writer::repository_lock::RepositoryLock,
}
impl RunDirectory {
fn create(prepared: &PreparedReindex) -> Result<Self> {
let name = format!(
"{RUN_DIRECTORY_PREFIX}{:016x}",
store_name_hash(&prepared.directory)
);
let path = prepared.plan.work_directory.join(name);
std::fs::create_dir_all(&path)?;
let lock = crate::writer::repository_lock::RepositoryLock::acquire(&path)?;
Ok(Self { path, _lock: lock })
}
fn remove(self) {
let path = self.path.clone();
drop(self);
let _ = std::fs::remove_dir_all(path);
}
}
fn store_name_hash(directory: &Path) -> u64 {
let mut hash = 0xcbf2_9ce4_8422_2325_u64;
for byte in directory.as_os_str().as_encoded_bytes() {
hash ^= u64::from(*byte);
hash = hash.wrapping_mul(0x0000_0100_0000_01b3);
}
hash
}
struct PublishedRun {
rebuilt: Vec<RebuiltDefinition>,
head_before: RecordIdentifier,
head_after: RecordIdentifier,
journal_lines_before: usize,
}
fn apply_under_the_lock(
prepared: &PreparedReindex,
run: &RunDirectory,
observer: &mut dyn ProgressObserver,
) -> Result<ReindexOutcome> {
let store = WritableRepository::open_prepared(
&prepared.directory,
prepared.repository_lock.clone(),
prepared.certified_archive_number,
)?;
let head_before = store.head();
let published = build_and_publish(&store, prepared, run, observer, head_before);
if published.is_err() {
let current = store.head();
if current != head_before {
store.compare_and_set_head(current, head_before);
}
}
let closed = store.close();
let published = published?;
closed?;
if published.head_after == published.head_before {
return Ok(ReindexOutcome {
definitions: Vec::new(),
head_before,
head_after: head_before,
});
}
probe(BEFORE_APPLIED_VERIFICATION)?;
verification::verify_after_reopen(
&prepared.directory,
published.head_after,
&published.rebuilt,
published.journal_lines_before,
)?;
Ok(ReindexOutcome {
definitions: published
.rebuilt
.iter()
.map(|definition| (definition.path.clone(), definition.report.clone()))
.collect(),
head_before: published.head_before,
head_after: published.head_after,
})
}
fn build_and_publish(
store: &WritableRepository,
prepared: &PreparedReindex,
run: &RunDirectory,
observer: &mut dyn ProgressObserver,
head_before: RecordIdentifier,
) -> Result<PublishedRun> {
let journal_lines_before =
crate::journal::read_journal(&prepared.directory.join("journal.log"))?.len();
let mut rebuilt = {
let mut writer = store.record_writer(store.writing_generation()?);
let step = Step::new(WRITE_STEP, WorkUnit::Nodes);
let rebuilt = observe(observer, &step, |observer| {
rebuild_every_definition(store, &mut writer, prepared, run, observer)
})?;
writer.finish()?;
rebuilt
};
if rebuilt.is_empty() {
return Ok(PublishedRun {
rebuilt,
head_before,
head_after: head_before,
journal_lines_before,
});
}
verification::perturb_all(store, &mut rebuilt)?;
let mut writer = store.record_writer(store.writing_generation()?);
let super_root = rewrite_the_spine(store, &mut writer, head_before, &rebuilt)?;
let step = Step::new(VERIFY_STEP, WorkUnit::Nodes);
observe(observer, &step, |observer| {
verification::verify_before_publication(store, &rebuilt, observer)
})?;
writer.finish()?;
probe(BEFORE_HEAD_PUBLISH)?;
if !store.compare_and_set_head(head_before, super_root) {
return Err(Error::InvalidFormat {
details: "the head moved while the reindex held the lock, which cannot happen \
and means the lock did not hold"
.to_owned(),
});
}
probe(AFTER_HEAD_PUBLISH_BEFORE_FLUSH)?;
Ok(PublishedRun {
rebuilt,
head_before,
head_after: super_root,
journal_lines_before,
})
}
fn rebuild_every_definition<Sink: SegmentSink>(
store: &WritableRepository,
writer: &mut RecordWriter<Sink>,
prepared: &PreparedReindex,
run: &RunDirectory,
observer: &mut dyn ProgressObserver,
) -> Result<Vec<RebuiltDefinition>> {
let super_root = store.head_node();
let inventory = crate::index::inventory::IndexInventory::collect(store, &super_root)
.map_err(index_error_to_store_error)?;
let selection = crate::writer::index::selection::select(
store,
&super_root,
&inventory,
&crate::writer::index::selection::SelectionOptions {
requested_paths: prepared.options.requested_paths().to_vec(),
from_head: prepared.options.from_head(),
has_binary_text_policy: prepared.options.binary_text_policy().is_some(),
},
)?;
let mut rebuilt = Vec::new();
for selected in &selection.selected {
let definition = crate::content::node::NodeState::new(store, selected.definition_record);
if matches!(selected.state, IndexingState::ResetForReplay { .. }) {
if let Some(built) = reset_one(store, writer, selected, &definition)? {
rebuilt.push(built);
}
continue;
}
let built = rebuild_one(
store,
writer,
selected,
&definition,
prepared,
run,
observer,
)?;
rebuilt.push(built);
}
Ok(rebuilt)
}
fn reset_one<Sink: SegmentSink>(
store: &WritableRepository,
writer: &mut RecordWriter<Sink>,
selected: &crate::writer::index::selection::SelectedIndex,
definition: &crate::content::node::NodeState<'_>,
) -> Result<Option<RebuiltDefinition>> {
let mut removed = Vec::new();
let mut retained = Vec::new();
for (name, child) in definition.child_node_entries()? {
if !name.starts_with(':') {
continue;
}
if crate::index::strict_boolean(child.property("retainNodeInReindex")?.as_ref()) {
retained.push(name);
} else {
removed.push(name);
}
}
if removed.is_empty() {
return Ok(None);
}
let rebuilt_record = rewrite_definition(store, writer, definition, &DefinitionEdits::reset())?;
Ok(Some(RebuiltDefinition {
path: selected.definition.path.clone(),
previous_record: selected.definition_record,
rebuilt_record,
report: DefinitionReport::Reset {
removed_hidden_children: removed,
retained_hidden_children: retained,
},
verification: Verification::NodeTreeOnly,
}))
}
type LuceneArm = (
Vec<(String, RecordIdentifier)>,
DefinitionReport,
Verification,
crate::writer::index::lucene_reindex::LuceneRebuild,
);
#[expect(
clippy::too_many_arguments,
reason = "the dispatch's own inputs, passed through rather than re-derived"
)]
fn rebuild_lucene_one<Sink: SegmentSink>(
writer: &mut RecordWriter<Sink>,
selected: &crate::writer::index::selection::SelectedIndex,
definition: &crate::content::node::NodeState<'_>,
prepared: &PreparedReindex,
state_root: &crate::content::node::NodeState<'_>,
location: &RunLocation,
budget: &SortBudget,
run: &RunDirectory,
observer: &mut dyn ProgressObserver,
) -> Result<LuceneArm> {
let policy = prepared.options.binary_text_policy().ok_or_else(|| {
Error::InvalidFormat {
details:
crate::writer::index::selection::SelectionRefusal::LuceneWithoutBinaryTextPolicy {
path: selected.definition.path.clone(),
}
.to_string(),
}
})?;
let mut warnings = Vec::new();
let rules = crate::index::lucene::documents::rules::IndexingRules::read(
definition,
&selected.definition.path,
state_root,
&mut warnings,
)
.map_err(index_error_to_store_error)?;
let segment_directory = run.path.join(format!(
"{}-segment",
sanitized_prefix(&selected.definition.path)
));
let subject = crate::writer::index::lucene_reindex::LuceneRebuildSubject {
state_root,
definition: &selected.definition,
rules: &rules,
policy,
segment_directory: &segment_directory,
runs: location,
budget,
};
let built =
crate::writer::index::lucene_reindex::rebuild_lucene_index(&subject, writer, observer)?;
let children = vec![(
crate::index::lucene::INDEX_DATA_CHILD_NAME.to_owned(),
built.data_record,
)];
let report = DefinitionReport::RebuiltIndex {
documents: built.documents,
nodes_visited: built.nodes_visited,
files: built.files.clone(),
segment_bytes: built.segment_bytes,
};
let verification = Verification::LuceneSegment {
documents: built.documents,
files: built.files.clone(),
};
Ok((children, report, verification, built))
}
fn checkpoint_creation_time(
store: &WritableRepository,
state: &IndexingState,
) -> Result<Option<String>> {
let IndexingState::LaneCheckpoint { checkpoint, .. } = state else {
return Ok(None);
};
let Some(checkpoints) = store.head_node().child_node("checkpoints")? else {
return Ok(None);
};
let Some(node) = checkpoints.child_node(checkpoint)? else {
return Ok(None);
};
let Some(properties) = node.child_node("properties")? else {
return Ok(None);
};
let Some(created) = properties.property("created")? else {
return Ok(None);
};
let Some(text) = crate::index::values_of(&created)
.first()
.and_then(crate::content::PropertyValue::as_text)
else {
return Ok(None);
};
Ok(crate::java::parse_epoch_milliseconds(&text).map(|_| text))
}
fn rebuild_one<Sink: SegmentSink>(
store: &WritableRepository,
writer: &mut RecordWriter<Sink>,
selected: &crate::writer::index::selection::SelectedIndex,
definition: &crate::content::node::NodeState<'_>,
prepared: &PreparedReindex,
run: &RunDirectory,
observer: &mut dyn ProgressObserver,
) -> Result<RebuiltDefinition> {
let Some(state_record) = selected.state_root else {
return Err(Error::InvalidFormat {
details: format!(
"{} was selected for a rebuild with no state root, which selection does \
not produce",
selected.definition.path
),
});
};
let state_root = crate::content::node::NodeState::new(store, state_record);
let budget = SortBudget::of_bytes(prepared.options.sort_budget_bytes());
let location = RunLocation::new(&run.path, sanitized_prefix(&selected.definition.path));
let mut created_seed = None;
let mut lucene_built = None;
let (hidden_children, report, verification) = match selected.definition.index_type.as_ref() {
Some(crate::index::IndexType::Counter) => {
let builder = CounterBuilder::new(&selected.definition);
let built = builder.build(&state_root, writer)?;
created_seed = built.created_seed;
let children = built
.index_record
.map(|record| vec![(crate::index::INDEX_CONTENT_NODE_NAME.to_owned(), record)])
.unwrap_or_default();
let nodes = built.credited_by_path.len() as u64;
(
children,
DefinitionReport::Rebuilt {
entries: nodes,
distinct_keys: nodes,
nodes_written: nodes,
},
Verification::Counter {
state_root: state_record,
credited_by_path: built.credited_by_path,
},
)
}
Some(crate::index::IndexType::Lucene) => {
let (children, report, verification, built) = rebuild_lucene_one(
writer,
selected,
definition,
prepared,
&state_root,
&location,
&budget,
run,
observer,
)?;
lucene_built = Some(built);
(children, report, verification)
}
Some(crate::index::IndexType::Reference) => {
let (children, report) =
build_reference_index(&state_root, writer, &location, &budget, observer)?;
let verification = entry_verification(state_record, &report);
(children, report, verification)
}
_ => {
let (children, report) = build_property_index(
&state_root,
&selected.definition,
writer,
&location,
&budget,
observer,
)?;
let verification = entry_verification(state_record, &report);
(children, report, verification)
}
};
let rebuilt_record = rewrite_the_rebuilt_definition(
store,
writer,
definition,
hidden_children,
created_seed,
lucene_built.as_ref(),
&selected.state,
)?;
Ok(RebuiltDefinition {
path: selected.definition.path.clone(),
previous_record: selected.definition_record,
rebuilt_record,
report,
verification,
})
}
fn rewrite_the_rebuilt_definition<Sink: SegmentSink>(
store: &WritableRepository,
writer: &mut RecordWriter<Sink>,
definition: &crate::content::node::NodeState<'_>,
hidden_children: Vec<(String, RecordIdentifier)>,
created_seed: Option<i64>,
lucene_built: Option<&crate::writer::index::lucene_reindex::LuceneRebuild>,
state: &IndexingState,
) -> Result<RecordIdentifier> {
let head_root = store
.head_node()
.child_node("root")?
.ok_or_else(|| Error::InvalidFormat {
details: "the super-root has no \"root\" child node".to_owned(),
})?;
let verdict = disabler_verdict(&head_root, definition)?;
let mut edits = DefinitionEdits::reindexed(verdict, hidden_children);
if let Some(built) = lucene_built {
let last_updated = checkpoint_creation_time(store, state)?;
crate::writer::index::lucene_reindex::lucene_definition_edits(
store,
writer,
definition,
built,
last_updated.as_deref(),
&mut edits,
)?;
}
if let Some(seed) = created_seed {
let value = writer.write_string(&seed.to_string())?;
edits.property_replacements.push(PropertyToWrite {
name: "seed".to_owned(),
property_type: crate::PropertyType::Long,
values: PropertyValuesToWrite::Single(value),
});
}
rewrite_definition(store, writer, definition, &edits)
}
fn entry_verification(state_root: RecordIdentifier, report: &DefinitionReport) -> Verification {
match report {
DefinitionReport::Rebuilt { entries, .. } => Verification::Entries {
state_root,
entries: *entries,
},
_ => Verification::NodeTreeOnly,
}
}
fn build_property_index<Sink: SegmentSink>(
state_root: &crate::content::node::NodeState<'_>,
definition: &crate::index::IndexDefinition,
writer: &mut RecordWriter<Sink>,
location: &RunLocation,
budget: &SortBudget,
observer: &mut dyn ProgressObserver,
) -> Result<(Vec<(String, RecordIdentifier)>, DefinitionReport)> {
let (collected, _) = PropertyCollector::collect(
state_root,
definition,
&EntrySink::Runs {
location: location.clone(),
budget: budget.clone(),
},
observer,
)?;
probe(BEFORE_SPILL_CLEANUP)?;
let CollectedEntries::Sorted(sorted) = collected else {
unreachable!("the run sink sorts");
};
let entries = sorted.emitted();
let (record, accounting, distinct_keys) = if definition.property.unique {
let mut builder = UniqueBuilder::new(writer);
let keys = feed(sorted, |key, path| builder.push(key, path))?;
let (record, accounting) = builder.finish()?;
(record, accounting, keys)
} else {
let mut builder = MirrorBuilder::new(writer);
let keys = feed(sorted, |key, path| builder.push(key, path))?;
let (record, accounting) = builder.finish()?;
(record, accounting, keys)
};
Ok((
vec![(crate::index::INDEX_CONTENT_NODE_NAME.to_owned(), record)],
DefinitionReport::Rebuilt {
entries,
distinct_keys,
nodes_written: accounting.nodes_written,
},
))
}
fn build_reference_index<Sink: SegmentSink>(
state_root: &crate::content::node::NodeState<'_>,
writer: &mut RecordWriter<Sink>,
location: &RunLocation,
budget: &SortBudget,
observer: &mut dyn ProgressObserver,
) -> Result<(Vec<(String, RecordIdentifier)>, DefinitionReport)> {
let (collected, _) = ReferenceCollector::collect(
state_root,
&EntrySink::Runs {
location: location.clone(),
budget: budget.clone(),
},
observer,
)?;
probe(BEFORE_SPILL_CLEANUP)?;
let CollectedReferences::Sorted(mut sets) = collected else {
unreachable!("the run sink sorts");
};
let mut children = Vec::new();
let mut entries = 0u64;
let mut nodes_written = 0u64;
let mut distinct_keys = 0u64;
let had_strong = sets.has_strong();
if let Some(sorted) = sets.strong()? {
let mut builder = MirrorBuilder::new(writer);
let keys = feed(sorted, |key, path| builder.push(key, path))?;
let (record, accounting) = builder.finish()?;
if had_strong {
children.push((":references".to_owned(), record));
nodes_written += accounting.nodes_written;
distinct_keys += keys;
entries += keys;
}
}
let had_weak = sets.has_weak();
if let Some(sorted) = sets.weak()? {
let mut builder = MirrorBuilder::new(writer);
let keys = feed(sorted, |key, path| builder.push(key, path))?;
let (record, accounting) = builder.finish()?;
if had_weak {
children.push((":weakreferences".to_owned(), record));
nodes_written += accounting.nodes_written;
distinct_keys += keys;
entries += keys;
}
}
Ok((
children,
DefinitionReport::Rebuilt {
entries,
distinct_keys,
nodes_written,
},
))
}
fn feed(sorted: SortedEntries, mut push: impl FnMut(&str, &str) -> Result<()>) -> Result<u64> {
let mut distinct_keys = 0u64;
let mut previous: Option<String> = None;
for entry in sorted {
let IndexEntry { key, path } = entry?;
if previous.as_deref() != Some(key.as_str()) {
distinct_keys += 1;
previous = Some(key.clone());
}
push(&key, &path)?;
}
Ok(distinct_keys)
}
fn rewrite_the_spine<Sink: SegmentSink>(
store: &WritableRepository,
writer: &mut RecordWriter<Sink>,
head: RecordIdentifier,
rebuilt: &[RebuiltDefinition],
) -> Result<RecordIdentifier> {
let head_node = crate::content::node::NodeState::new(store, head);
let content_root = head_node
.child_node("root")?
.ok_or_else(|| Error::InvalidFormat {
details: "the super-root has no \"root\" child node".to_owned(),
})?;
let oak_index = content_root
.child_node(crate::index::INDEX_DEFINITIONS_NAME)?
.ok_or_else(|| Error::InvalidFormat {
details: "the content root has no /oak:index".to_owned(),
})?;
let mut definition_edits = crate::writer::commit::ChildEdits::new();
for definition in rebuilt {
let name = definition
.path
.rsplit('/')
.next()
.unwrap_or(&definition.path)
.to_owned();
definition_edits.insert(name, Some(definition.rebuilt_record));
}
let new_oak_index = crate::writer::commit::rewrite_node_with_child_edits(
store,
writer,
Some(oak_index.record_identifier()),
&definition_edits,
)?;
let mut root_edits = crate::writer::commit::ChildEdits::new();
root_edits.insert(
crate::index::INDEX_DEFINITIONS_NAME.to_owned(),
Some(new_oak_index),
);
let new_root = crate::writer::commit::rewrite_node_with_child_edits(
store,
writer,
Some(content_root.record_identifier()),
&root_edits,
)?;
let mut super_edits = crate::writer::commit::ChildEdits::new();
super_edits.insert("root".to_owned(), Some(new_root));
crate::writer::commit::rewrite_node_with_child_edits(store, writer, Some(head), &super_edits)
}
fn sanitized_prefix(path: &str) -> String {
let mut prefix = String::with_capacity(path.len());
for character in path.chars() {
if character.is_ascii_alphanumeric() {
prefix.push(character);
} else {
prefix.push('-');
}
}
prefix
}
fn index_error_to_store_error(error: crate::index::IndexError) -> Error {
match error {
crate::index::IndexError::Record(source) => source,
other => Error::InvalidFormat {
details: other.to_string(),
},
}
}
mod verification;
#[cfg(test)]
mod tests;