use std::io::Read;
use crate::content::node::NodeState;
use crate::content::property::{PropertyType, PropertyValue};
use crate::content::value::{BinaryValue, read_binary_stream};
use crate::content::{PropertyValues, SegmentProvider};
use crate::index::{IndexResult, converting_strings};
const STREAM_BUFFER_BYTES: usize = 64 * 1024;
#[derive(Clone, PartialEq, Eq, Debug, Default)]
pub struct LuceneBlobReport {
pub type_mismatch: bool,
pub missing_blobs: Vec<BlobFault>,
pub invalid_blobs: Vec<BlobFault>,
pub blobs_checked: u64,
pub bytes_read: u64,
}
impl LuceneBlobReport {
#[must_use]
pub fn is_consistent(&self) -> bool {
!self.type_mismatch && self.missing_blobs.is_empty() && self.invalid_blobs.is_empty()
}
}
#[derive(Clone, PartialEq, Eq, PartialOrd, Ord, Debug)]
pub struct BlobFault {
pub node_path: String,
pub property_name: String,
pub value_index: usize,
pub declared_length: u64,
pub streamed_length: Option<u64>,
pub reason: Option<String>,
}
pub fn check_blobs(
provider: &dyn SegmentProvider,
definition_node: &NodeState<'_>,
) -> IndexResult<LuceneBlobReport> {
let mut report = LuceneBlobReport::default();
let declared_type = definition_node.property("type")?;
let reads_as_lucene = converting_strings(declared_type.as_ref())
.first()
.is_some_and(|text| text == "lucene");
if !reads_as_lucene {
report.type_mismatch = true;
return Ok(report);
}
let mut buffer = vec![0u8; STREAM_BUFFER_BYTES];
let mut stack = vec![(String::new(), *definition_node)];
while let Some((node_path, node)) = stack.pop() {
for property in node.properties()? {
if property.property_type != PropertyType::Binary {
continue;
}
let values = match &property.values {
PropertyValues::Single(value) => std::slice::from_ref(value),
PropertyValues::Multiple(values) => values.as_slice(),
};
for (value_index, value) in values.iter().enumerate() {
check_one_blob(
provider,
&mut buffer,
&node_path,
&property.name,
value_index,
value,
&mut report,
);
}
}
for (name, child) in node.child_node_entries()?.into_iter().rev() {
stack.push((format!("{node_path}/{name}"), child));
}
}
Ok(report)
}
fn check_one_blob(
provider: &dyn SegmentProvider,
buffer: &mut [u8],
node_path: &str,
property_name: &str,
value_index: usize,
value: &PropertyValue,
report: &mut LuceneBlobReport,
) {
let fault = |declared_length, streamed_length, reason| BlobFault {
node_path: node_path.to_owned(),
property_name: property_name.to_owned(),
value_index,
declared_length,
streamed_length,
reason,
};
let (declared_length, record_identifier) = match value {
PropertyValue::Binary(BinaryValue::Inline {
length,
record_identifier,
}) => (*length, *record_identifier),
PropertyValue::Binary(BinaryValue::External { blob_identifier }) => {
report.missing_blobs.push(fault(
0,
None,
Some(format!(
"the value names the external blob {blob_identifier:?}, which lives \
outside the segment store"
)),
));
return;
}
_ => return,
};
let mut stream = match read_binary_stream(provider, record_identifier) {
Ok(stream) => stream,
Err(error) => {
report
.missing_blobs
.push(fault(declared_length, None, Some(error.to_string())));
return;
}
};
let mut streamed_length = 0u64;
loop {
match stream.read(buffer) {
Ok(0) => break,
Ok(count) => streamed_length += count as u64,
Err(error) => {
report.invalid_blobs.push(fault(
declared_length,
Some(streamed_length),
Some(error.to_string()),
));
return;
}
}
}
report.bytes_read += streamed_length;
if streamed_length == declared_length {
report.blobs_checked += 1;
} else {
report
.invalid_blobs
.push(fault(declared_length, Some(streamed_length), None));
}
}
#[derive(Clone, PartialEq, Eq, Debug, Default)]
pub struct LuceneStructuralReport {
pub generation: Option<i64>,
pub commit_file: Option<String>,
pub codec_name: Option<String>,
pub unregistered_codecs: Vec<String>,
pub segment_count: usize,
pub live_document_count: i64,
pub missing_files: Vec<String>,
pub unreferenced_files: Vec<String>,
pub unreadable_files: Vec<(String, String)>,
}
impl LuceneStructuralReport {
#[must_use]
pub fn is_coherent(&self) -> bool {
self.missing_files.is_empty()
&& self.unreferenced_files.is_empty()
&& self.unreadable_files.is_empty()
&& self.commit_file.is_some()
}
}
pub const REGISTERED_CODEC_NAMES: [&str; 3] = ["oakCodec", "Lucene46", "compressingCodec"];
fn is_never_referenced(name: &str) -> bool {
name == crate::index::lucene::segments::SEGMENTS_GEN_FILE_NAME
}
pub trait LuceneFileSource {
fn file_names(&self) -> Vec<String>;
fn read_file(&self, name: &str) -> std::result::Result<Vec<u8>, String>;
}
impl LuceneFileSource for crate::index::lucene::OakDirectory<'_> {
fn file_names(&self) -> Vec<String> {
crate::index::lucene::OakDirectory::file_names(self).to_vec()
}
fn read_file(&self, name: &str) -> std::result::Result<Vec<u8>, String> {
let file = self.file(name).map_err(|error| error.to_string())?;
let mut bytes = Vec::with_capacity(file.length() as usize);
std::io::Read::read_to_end(&mut file.reader(), &mut bytes)
.map_err(|error| error.to_string())?;
Ok(bytes)
}
}
pub struct LocalIndexDirectory {
path: std::path::PathBuf,
}
impl LocalIndexDirectory {
#[must_use]
pub fn new(path: impl Into<std::path::PathBuf>) -> Self {
Self { path: path.into() }
}
}
impl LuceneFileSource for LocalIndexDirectory {
fn file_names(&self) -> Vec<String> {
let Ok(entries) = std::fs::read_dir(&self.path) else {
return Vec::new();
};
let mut names: Vec<String> = entries
.filter_map(Result::ok)
.filter(|entry| entry.path().is_file())
.filter_map(|entry| entry.file_name().to_str().map(str::to_owned))
.collect();
names.sort();
names
}
fn read_file(&self, name: &str) -> std::result::Result<Vec<u8>, String> {
std::fs::read(self.path.join(name)).map_err(|error| error.to_string())
}
}
pub fn check_structure<Source: LuceneFileSource + ?Sized>(
directory: &Source,
) -> IndexResult<LuceneStructuralReport> {
let mut report = LuceneStructuralReport::default();
let listing: Vec<String> = LuceneFileSource::file_names(directory);
let hint = listing
.iter()
.any(|name| name == crate::index::lucene::segments::SEGMENTS_GEN_FILE_NAME)
.then(|| {
let bytes = directory
.read_file(crate::index::lucene::segments::SEGMENTS_GEN_FILE_NAME)
.ok()?;
let length = bytes.len() as u64;
let mut reader = crate::index::lucene::read::Reader::new(
std::io::Cursor::new(bytes),
crate::index::lucene::segments::SEGMENTS_GEN_FILE_NAME,
length,
);
crate::index::lucene::segments::read_segments_gen(&mut reader)
})
.flatten();
let Some(generation) = crate::index::lucene::segments::commit_generation(&listing, hint) else {
report.unreferenced_files = listing
.iter()
.filter(|name| !is_never_referenced(name))
.cloned()
.collect();
report.unreferenced_files.sort();
return Ok(report);
};
report.generation = Some(generation);
let commit_name: Option<String> = listing
.iter()
.filter(|name| crate::index::lucene::segments::is_commit_file_name(name))
.filter(|name| {
crate::index::lucene::segments::generation_from_commit_file_name(name)
== Some(generation)
})
.max()
.cloned();
let Some(commit_name) = commit_name else {
report.unreadable_files.push((
format!("segments_{generation:x}"),
"the generation hint names a commit file the directory does not hold".to_owned(),
));
return Ok(report);
};
let commit = match read_commit(directory, &commit_name) {
Ok(commit) => commit,
Err(details) => {
report.unreadable_files.push((commit_name, details));
return Ok(report);
}
};
report.commit_file = Some(commit.file_name.clone());
report.segment_count = commit.segments.len();
report.live_document_count = commit.live_document_count();
let mut codec_names: Vec<String> = commit
.segments
.iter()
.map(|segment| segment.codec_name.clone())
.collect();
codec_names.sort();
codec_names.dedup();
report.unregistered_codecs = codec_names
.iter()
.filter(|name| !REGISTERED_CODEC_NAMES.contains(&name.as_str()))
.cloned()
.collect();
if codec_names.len() == 1 {
report.codec_name = codec_names.into_iter().next();
}
let referenced = commit.referenced_files();
report.missing_files = referenced
.iter()
.filter(|name| !listing.contains(name))
.cloned()
.collect();
report.unreferenced_files = listing
.iter()
.filter(|name| !referenced.contains(name) && !is_never_referenced(name))
.cloned()
.collect();
report.missing_files.sort();
report.unreferenced_files.sort();
Ok(report)
}
fn read_commit<Source: LuceneFileSource + ?Sized>(
directory: &Source,
commit_name: &str,
) -> std::result::Result<crate::index::lucene::segments::CommitFile, String> {
let bytes = directory.read_file(commit_name)?;
let length = bytes.len() as u64;
let mut reader =
crate::index::lucene::read::Reader::new(std::io::Cursor::new(bytes), commit_name, length);
crate::index::lucene::segments::read_commit_file(&mut reader, commit_name, |info_name| {
let bytes = directory.read_file(info_name).map_err(|details| {
crate::index::lucene::read::LuceneReadError::Source {
file: info_name.to_owned(),
details,
}
})?;
let info_length = bytes.len() as u64;
Ok(crate::index::lucene::read::Reader::new(
std::io::Cursor::new(bytes),
info_name,
info_length,
))
})
.map_err(|error| error.to_string())
}