use std::collections::BTreeMap;
use std::io::Read as _;
use std::path::{Path, PathBuf};
use crate::error::{Error, Result};
use crate::index::IndexDefinition;
use crate::progress::{ProgressObserver, Step, WorkUnit, count, observe};
use crate::segment::record::RecordIdentifier;
use crate::writer::index::definition_update::{
DefinitionEdits, INDEX_VERSION_PROPERTY, ReindexCount, STORED_DEFINITION_CHILD,
clone_updated_definition_state, disabler_verdict, fresh_index_format_version,
rewrite_definition,
};
use crate::writer::index::lucene_directory::{DirectoryListing, OakDirectoryWriter};
use crate::writer::index::lucene_import::plan::PlannedImport;
use crate::writer::index::lucene_import::prepared::PreparedLuceneImport;
use crate::writer::record_writer::{
PropertyToWrite, PropertyValuesToWrite, RecordWriter, SegmentSink,
};
use crate::writer::store_writer::WritableRepository;
const COPY_STEP: &str = "copying index files";
const VERIFY_STEP: &str = "verifying the imported index";
#[cfg(test)]
const MID_FILE_COPY: &str = "lucene-import.mid-file-copy";
const BEFORE_HEAD_PUBLISH: &str = "lucene-import.before-head-publish";
const AFTER_HEAD_PUBLISH_BEFORE_FLUSH: &str = "lucene-import.after-head-publish-before-flush";
const BEFORE_APPLIED_VERIFICATION: &str = "lucene-import.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(())
}
#[derive(Clone, PartialEq, Eq, Debug)]
#[non_exhaustive]
pub struct ImportedIndex {
pub path: String,
pub files: Vec<(String, u64)>,
pub unique_identifier: String,
pub reindex_count: u64,
pub dropped_hidden_children: Vec<String>,
pub skipped_mappings: Vec<(String, String)>,
}
#[derive(Clone, PartialEq, Eq, Debug)]
#[non_exhaustive]
pub struct LuceneImportOutcome {
pub indexes: Vec<ImportedIndex>,
pub checkpoint: String,
pub head_before: RecordIdentifier,
pub head_after: RecordIdentifier,
}
impl LuceneImportOutcome {
#[must_use]
pub fn moved_the_head(&self) -> bool {
self.head_before != self.head_after
}
#[must_use]
pub fn file_count(&self) -> usize {
self.indexes.iter().map(|index| index.files.len()).sum()
}
}
pub(crate) fn apply_prepared(
prepared: &PreparedLuceneImport,
observer: &mut dyn ProgressObserver,
) -> Result<LuceneImportOutcome> {
prepared.recheck_before_mutation()?;
let store = WritableRepository::open_prepared(
&prepared.directory,
prepared.repository_lock.clone(),
prepared.certified_archive_number,
)?;
let head_before = store.head();
let written = write_every_index(&store, prepared, observer);
if written.is_err() {
let current = store.head();
if current != head_before {
store.compare_and_set_head(current, head_before);
}
}
let closed = store.close();
let written = written?;
closed?;
probe(BEFORE_APPLIED_VERIFICATION)?;
verify_after_reopen(&prepared.directory, written.head_after, &written.indexes)?;
Ok(LuceneImportOutcome {
indexes: written.indexes,
checkpoint: prepared.plan.checkpoint.clone(),
head_before,
head_after: written.head_after,
})
}
struct WrittenImport {
indexes: Vec<ImportedIndex>,
head_after: RecordIdentifier,
}
fn write_every_index(
store: &WritableRepository,
prepared: &PreparedLuceneImport,
observer: &mut dyn ProgressObserver,
) -> Result<WrittenImport> {
let head_before = store.head();
let total = prepared.plan.file_count();
let mut rewritten: BTreeMap<String, RecordIdentifier> = BTreeMap::new();
let mut indexes = Vec::new();
{
let mut writer = store.record_writer(store.writing_generation()?);
let step = Step::new(COPY_STEP, WorkUnit::Files).with_total(count(total));
let mut copied = 0usize;
observe(observer, &step, |observer| {
for import in &prepared.plan.imports {
let (record, imported) = import_one(store, &mut writer, import, &mut |written| {
observer.step_advanced(count(copied + written));
})?;
copied += imported.files.len();
rewritten.insert(import.path.clone(), record);
indexes.push(imported);
}
observer.step_advanced(count(copied));
Ok::<_, Error>(())
})?;
writer.finish()?;
}
let mut writer = store.record_writer(store.writing_generation()?);
let head_after = rewrite_the_spine(store, &mut writer, head_before, &rewritten)?;
let step = Step::new(VERIFY_STEP, WorkUnit::Files).with_total(count(total));
observe(observer, &step, |observer| {
verify_before_publication(store, prepared, &rewritten, observer)
})?;
writer.finish()?;
probe(BEFORE_HEAD_PUBLISH)?;
if !store.compare_and_set_head(head_before, head_after) {
return Err(Error::InvalidFormat {
details: "the head moved while the import held the lock, which cannot happen \
and means the lock did not hold"
.to_owned(),
});
}
probe(AFTER_HEAD_PUBLISH_BEFORE_FLUSH)?;
store.flush()?;
Ok(WrittenImport {
indexes,
head_after,
})
}
fn import_one<Sink: SegmentSink>(
store: &WritableRepository,
writer: &mut RecordWriter<Sink>,
import: &PlannedImport,
report: &mut dyn FnMut(usize),
) -> Result<(RecordIdentifier, ImportedIndex)> {
let content_root =
store
.head_node()
.child_node("root")?
.ok_or_else(|| Error::InvalidFormat {
details: "the super-root has no \"root\" child node".to_owned(),
})?;
let definition_node =
descend(&content_root, &import.path)?.ok_or_else(|| Error::InvalidFormat {
details: format!("{} vanished between the plan and the apply", import.path),
})?;
let definition = IndexDefinition::read(&definition_node, &import.path).map_err(index_error)?;
let mut hidden_children = Vec::new();
let mut files = Vec::new();
for (jcr_name, source) in &import.mappings {
let (record, written) = write_directory(writer, &definition, source, &mut |count| {
report(files.len() + count);
})?;
files.extend(written);
hidden_children.push((jcr_name.clone(), record));
}
let unique_identifier = epoch_milliseconds().to_string();
let status = write_status_node(writer, &unique_identifier)?;
hidden_children.push((":status".to_owned(), status));
let dropped: Vec<String> = definition_node
.child_node_entries()?
.into_iter()
.map(|(name, _)| name)
.filter(|name| {
name.starts_with(':') && !hidden_children.iter().any(|(kept, _)| kept == name)
})
.collect();
let verdict = disabler_verdict(&content_root, &definition_node)?;
let mut edits = DefinitionEdits::reindexed(verdict, hidden_children);
edits.reindex_count = ReindexCount::Set(import.reindex_count);
edits.property_removals.push("corrupt".to_owned());
edits.property_removals.push("indexImportState".to_owned());
let version = fresh_index_format_version(&definition_node)?;
let version_value = writer.write_string(&version.to_string())?;
edits.property_replacements.push(PropertyToWrite {
name: INDEX_VERSION_PROPERTY.to_owned(),
property_type: crate::PropertyType::Long,
values: PropertyValuesToWrite::Single(version_value),
});
let stored = clone_updated_definition_state(store, writer, &definition_node, &edits)?;
edits
.hidden_children
.push((STORED_DEFINITION_CHILD.to_owned(), stored));
let record = rewrite_definition(store, writer, &definition_node, &edits)?;
Ok((
record,
ImportedIndex {
path: import.path.clone(),
files,
unique_identifier,
reindex_count: import.reindex_count,
dropped_hidden_children: dropped,
skipped_mappings: import.skipped_mappings.clone(),
},
))
}
fn write_directory<Sink: SegmentSink>(
writer: &mut RecordWriter<Sink>,
definition: &IndexDefinition,
source: &Path,
report: &mut dyn FnMut(usize),
) -> Result<(RecordIdentifier, Vec<(String, u64)>)> {
let mut names: Vec<PathBuf> = std::fs::read_dir(source)?
.filter_map(|entry| entry.ok().map(|entry| entry.path()))
.filter(|path| path.is_file())
.collect();
names.sort();
let blob_size = definition.lucene.blob_size;
let listing = if definition.lucene.save_directory_listing {
DirectoryListing::Saved
} else {
DirectoryListing::Omitted
};
let mut builder = OakDirectoryWriter::new(writer, blob_size, listing);
#[cfg(test)]
let largest = names
.iter()
.max_by_key(|path| std::fs::metadata(path).map_or(0, |data| data.len()))
.cloned();
let mut written = Vec::new();
for path in names {
let name = path
.file_name()
.and_then(|name| name.to_str())
.ok_or_else(|| Error::InvalidFormat {
details: format!("{} has a name that is not UTF-8", path.display()),
})?
.to_owned();
let length = std::fs::metadata(&path)?.len();
let handle = std::fs::File::open(&path)?;
#[cfg(test)]
{
if Some(&path) == largest.as_ref() {
builder.add_file(&name, ProbingReader::halfway_through(handle, &name, length))?;
} else {
builder.add_file(&name, handle)?;
}
}
#[cfg(not(test))]
{
builder.add_file(&name, handle)?;
}
written.push((name, length));
report(written.len());
}
Ok((builder.finish()?, written))
}
#[cfg(test)]
struct ProbingReader<Source> {
source: Source,
name: String,
read: u64,
fire_at: u64,
fired: bool,
}
#[cfg(test)]
impl<Source: std::io::Read> ProbingReader<Source> {
fn halfway_through(source: Source, name: &str, length: u64) -> Self {
Self {
source,
name: name.to_owned(),
read: 0,
fire_at: length / 2,
fired: false,
}
}
}
#[cfg(test)]
impl<Source: std::io::Read> std::io::Read for ProbingReader<Source> {
fn read(&mut self, buffer: &mut [u8]) -> std::io::Result<usize> {
let count = self.source.read(buffer)?;
self.read += count as u64;
if !self.fired && self.read >= self.fire_at {
self.fired = true;
probe(MID_FILE_COPY).map_err(|error| {
std::io::Error::other(format!(
"the copy of {} stopped at byte {}: {error}",
self.name, self.read
))
})?;
}
Ok(count)
}
}
fn write_status_node<Sink: SegmentSink>(
writer: &mut RecordWriter<Sink>,
unique_identifier: &str,
) -> Result<RecordIdentifier> {
let value = writer.write_string(unique_identifier)?;
writer.write_node(
None,
&[],
&crate::writer::record_writer::ChildNodesToWrite::Zero,
&[PropertyToWrite {
name: "uid".to_owned(),
property_type: crate::PropertyType::String,
values: PropertyValuesToWrite::Single(value),
}],
)
}
fn rewrite_the_spine<Sink: SegmentSink>(
store: &WritableRepository,
writer: &mut RecordWriter<Sink>,
head: RecordIdentifier,
rewritten: &BTreeMap<String, RecordIdentifier>,
) -> Result<RecordIdentifier> {
let super_root = store.head_node();
let content_root = super_root
.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 store has no /oak:index".to_owned(),
})?;
let mut definition_edits = crate::writer::commit::ChildEdits::new();
for (path, record) in rewritten {
let name = path.rsplit('/').next().unwrap_or(path).to_owned();
definition_edits.insert(name, Some(*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 verify_before_publication(
store: &WritableRepository,
prepared: &PreparedLuceneImport,
rewritten: &BTreeMap<String, RecordIdentifier>,
observer: &mut dyn ProgressObserver,
) -> Result<()> {
let mut verified = 0usize;
for import in &prepared.plan.imports {
let record = rewritten
.get(&import.path)
.copied()
.ok_or_else(|| Error::InvalidFormat {
details: format!("{} was planned but not written", import.path),
})?;
let node = crate::content::node::NodeState::new(store, record);
let definition = IndexDefinition::read(&node, &import.path).map_err(index_error)?;
for (jcr_name, source) in &import.mappings {
let directory =
crate::index::lucene::OakDirectory::open(store, &node, &definition, jcr_name)
.map_err(index_error)?
.ok_or_else(|| Error::InvalidFormat {
details: format!("{} has no {jcr_name} after the import", import.path),
})?;
for file_name in directory.file_names() {
let file = directory.file(file_name).map_err(index_error)?;
let mut stored = Vec::new();
file.reader().read_to_end(&mut stored)?;
let original = std::fs::read(source.join(file_name))?;
if stored != original {
return Err(Error::InvalidFormat {
details: format!(
"{}'s {jcr_name}/{file_name} reads back as {} bytes where the \
file on disk holds {}; the import is refused before anything \
is published",
import.path,
stored.len(),
original.len()
),
});
}
verified += 1;
observer.step_advanced(count(verified));
}
}
}
Ok(())
}
fn verify_after_reopen(
directory: &Path,
published: RecordIdentifier,
indexes: &[ImportedIndex],
) -> Result<()> {
let repository = crate::store::Repository::open(directory)?;
if repository.head_record_identifier() != published {
return Err(Error::InvalidFormat {
details: format!(
"the reopened store's head is {} rather than the {published} just published",
repository.head_record_identifier()
),
});
}
for index in indexes {
let node = repository
.node_at_path(&index.path)?
.ok_or_else(|| Error::InvalidFormat {
details: format!("the published head does not reach {}", index.path),
})?;
crate::tooling::check::verify_node_tree(&repository, node.record_identifier())?;
}
Ok(())
}
fn descend<'store>(
root: &crate::content::node::NodeState<'store>,
path: &str,
) -> Result<Option<crate::content::node::NodeState<'store>>> {
let mut node = *root;
for element in path.split('/').filter(|element| !element.is_empty()) {
let Some(child) = node.child_node(element)? else {
return Ok(None);
};
node = child;
}
Ok(Some(node))
}
fn index_error(error: crate::index::IndexError) -> Error {
match error {
crate::index::IndexError::Record(source) => source,
other => Error::InvalidFormat {
details: other.to_string(),
},
}
}
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))
}