use std::collections::{BTreeSet, HashSet};
use crate::content::node::NodeState;
use crate::content::property::{PropertyType, PropertyValue};
use crate::error::{Error, Result};
use crate::index::definition::IndexDefinition;
use crate::index::path_filter::PathVerdict;
use crate::index::property::key_encoding::keys_for_property;
use crate::index::property::type_predicate::TypePredicate;
use crate::progress::{ProgressObserver, Step, WorkUnit, count, observe};
use crate::segment::record::RecordIdentifier;
use crate::writer::index::{IndexEntry, RunLocation, SortBudget};
const COLLECT_STEP: &str = "collecting index entries";
const SORT_STEP: &str = "sorting index entries";
const VERSION_STORE_PATH: &str = "/jcr:system/jcr:versionStorage";
pub enum EntrySink {
Count,
Runs {
location: RunLocation,
budget: SortBudget,
},
}
pub enum CollectedEntries {
Counted {
entries: u64,
bytes: u64,
},
Sorted(SortedEntries),
}
pub struct SortedEntries {
inner: crate::external_sort::SortedPass<'static, IndexEntry>,
emitted: u64,
}
impl SortedEntries {
#[must_use]
pub fn emitted(&self) -> u64 {
self.emitted
}
}
impl Iterator for SortedEntries {
type Item = Result<IndexEntry>;
fn next(&mut self) -> Option<Self::Item> {
self.inner.next()
}
}
#[derive(Clone, Copy, PartialEq, Eq, Debug, Default)]
pub struct WalkAccounting {
pub nodes_visited: u64,
}
pub struct PropertyCollector;
impl PropertyCollector {
pub fn collect(
state_root: &NodeState<'_>,
definition: &IndexDefinition,
sink: &EntrySink,
observer: &mut dyn ProgressObserver,
) -> Result<(CollectedEntries, WalkAccounting)> {
let predicate = TypePredicate::new(state_root, &definition.property.declaring_node_types)
.map_err(index_error_to_store_error)?;
let mut accounting = WalkAccounting::default();
let mut counted = (0u64, 0u64);
let mut runs = match sink {
EntrySink::Count => None,
EntrySink::Runs { location, budget } => Some(crate::external_sort::SortedRuns::<
IndexEntry,
>::new(
location.clone(), budget.clone()
)),
};
let step = Step::new(COLLECT_STEP, WorkUnit::Nodes);
observe(observer, &step, |observer| {
walk_visible(state_root, |node, path| {
accounting.nodes_visited += 1;
observer.step_advanced(count(accounting.nodes_visited as usize));
let verdict = definition.path_filter.filter(path);
if verdict != PathVerdict::Include {
return Ok(());
}
if !predicate.is_empty()
&& !predicate.test(node).map_err(index_error_to_store_error)?
{
return Ok(());
}
for key in keys_for_node(node, definition)? {
let entry = IndexEntry::new(key, path);
match runs.as_mut() {
None => {
counted.0 += 1;
counted.1 += entry_bytes(&entry);
}
Some(runs) => {
counted.0 += 1;
runs.push(entry)?;
}
}
}
Ok(())
})
})?;
let collected = match runs {
None => CollectedEntries::Counted {
entries: counted.0,
bytes: counted.1,
},
Some(runs) => {
let step = Step::new(SORT_STEP, WorkUnit::IndexEntries).with_total(counted.0);
let sorted = observe(observer, &step, |_| runs.into_sorted())?;
CollectedEntries::Sorted(SortedEntries {
inner: sorted,
emitted: counted.0,
})
}
};
Ok((collected, accounting))
}
}
fn keys_for_node(node: &NodeState<'_>, definition: &IndexDefinition) -> Result<BTreeSet<String>> {
let mut keys = BTreeSet::new();
for name in &definition.property.property_names {
if name.starts_with(':') {
continue;
}
let Some(property) = node.property(name)? else {
continue;
};
keys.extend(
keys_for_property(
&property,
&definition.property.value_pattern,
&definition.path,
)
.map_err(index_error_to_store_error)?,
);
}
Ok(keys)
}
pub enum CollectedReferences {
Counted {
strong: u64,
weak: u64,
bytes: u64,
},
Sorted(Box<SortedReferenceSets>),
}
pub struct SortedReferenceSets {
strong_runs: Option<crate::external_sort::SortedRuns<IndexEntry>>,
weak_runs: Option<crate::external_sort::SortedRuns<IndexEntry>>,
strong_seen: bool,
weak_seen: bool,
}
impl SortedReferenceSets {
#[must_use]
pub fn has_strong(&self) -> bool {
self.strong_seen
}
#[must_use]
pub fn has_weak(&self) -> bool {
self.weak_seen
}
pub fn strong(&mut self) -> Result<Option<SortedEntries>> {
let Some(runs) = self.strong_runs.take() else {
return Ok(None);
};
Ok(Some(SortedEntries {
inner: runs.into_sorted()?,
emitted: 0,
}))
}
pub fn weak(&mut self) -> Result<Option<SortedEntries>> {
assert!(
self.strong_runs.is_none(),
"the weak set is merged only after the strong set, so the open-file \
bound stays at the fan-in plus one"
);
let Some(runs) = self.weak_runs.take() else {
return Ok(None);
};
Ok(Some(SortedEntries {
inner: runs.into_sorted()?,
emitted: 0,
}))
}
}
pub struct ReferenceCollector;
struct ReferenceEmitter {
strong_runs: Option<crate::external_sort::SortedRuns<IndexEntry>>,
weak_runs: Option<crate::external_sort::SortedRuns<IndexEntry>>,
strong: u64,
weak: u64,
bytes: u64,
}
impl ReferenceEmitter {
fn new(sink: &EntrySink) -> Self {
let (strong_runs, weak_runs) = match sink {
EntrySink::Count => (None, None),
EntrySink::Runs { location, budget } => (
Some(crate::external_sort::SortedRuns::new(
RunLocation::new(
location.directory(),
format!("{}-strong", location.name_prefix()),
),
budget.clone(),
)),
Some(crate::external_sort::SortedRuns::new(
RunLocation::new(
location.directory(),
format!("{}-weak", location.name_prefix()),
),
budget.clone(),
)),
),
};
Self {
strong_runs,
weak_runs,
strong: 0,
weak: 0,
bytes: 0,
}
}
fn emit_node(&mut self, node: &NodeState<'_>, path: &str) -> Result<()> {
let in_version_store = path.starts_with(VERSION_STORE_PATH);
for property in node.properties()? {
if property.name.starts_with(':') {
continue;
}
let strong = match property.property_type {
PropertyType::Reference => true,
PropertyType::WeakReference => false,
_ => continue,
};
if strong && in_version_store {
continue;
}
self.emit_property(&property, path, strong)?;
}
Ok(())
}
fn emit_property(
&mut self,
property: &crate::content::node::PropertyState,
path: &str,
strong: bool,
) -> Result<()> {
let identifiers: BTreeSet<String> = crate::index::values_of(property)
.iter()
.filter_map(PropertyValue::as_text)
.collect();
let property_path = relative_property_path(path, &property.name);
for identifier in identifiers {
let entry = IndexEntry::new(identifier, &property_path);
self.bytes += entry_bytes(&entry);
let runs = if strong {
self.strong += 1;
self.strong_runs.as_mut()
} else {
self.weak += 1;
self.weak_runs.as_mut()
};
if let Some(runs) = runs {
runs.push(entry)?;
}
}
Ok(())
}
fn finish(self) -> CollectedReferences {
let (strong_seen, weak_seen) = (self.strong > 0, self.weak > 0);
match (self.strong_runs, self.weak_runs) {
(Some(strong_runs), Some(weak_runs)) => {
CollectedReferences::Sorted(Box::new(SortedReferenceSets {
strong_runs: Some(strong_runs),
weak_runs: Some(weak_runs),
strong_seen,
weak_seen,
}))
}
_ => CollectedReferences::Counted {
strong: self.strong,
weak: self.weak,
bytes: self.bytes,
},
}
}
}
impl ReferenceCollector {
pub fn collect(
state_root: &NodeState<'_>,
sink: &EntrySink,
observer: &mut dyn ProgressObserver,
) -> Result<(CollectedReferences, WalkAccounting)> {
let mut accounting = WalkAccounting::default();
let mut emitter = ReferenceEmitter::new(sink);
let step = Step::new(COLLECT_STEP, WorkUnit::Nodes);
observe(observer, &step, |observer| {
walk_visible(state_root, |node, path| {
accounting.nodes_visited += 1;
observer.step_advanced(count(accounting.nodes_visited as usize));
emitter.emit_node(node, path)
})
})?;
Ok((emitter.finish(), accounting))
}
}
fn relative_property_path(node_path: &str, property_name: &str) -> String {
let absolute = if node_path == "/" {
format!("/{property_name}")
} else {
format!("{node_path}/{property_name}")
};
absolute.strip_prefix('/').unwrap_or(&absolute).to_owned()
}
fn entry_bytes(entry: &IndexEntry) -> u64 {
(entry.key.len() + entry.path.len()) as u64
}
pub(crate) fn walk_visible(
root: &NodeState<'_>,
mut visit: impl FnMut(&NodeState<'_>, &str) -> Result<()>,
) -> Result<()> {
enum Step<'provider> {
Visit {
node: NodeState<'provider>,
path: String,
},
Leave {
record: RecordIdentifier,
},
}
let mut stack = vec![Step::Visit {
node: *root,
path: String::new(),
}];
let mut ancestors: HashSet<RecordIdentifier> = HashSet::new();
while let Some(step) = stack.pop() {
match step {
Step::Leave { record } => {
ancestors.remove(&record);
}
Step::Visit { node, path } => {
let record = node.record_identifier();
if !ancestors.insert(record) {
return Err(Error::InvalidFormat {
details: format!(
"the node at {} is its own ancestor, so the index cannot be \
rebuilt from it; this is corruption, not deep content",
if path.is_empty() { "/" } else { &path }
),
});
}
stack.push(Step::Leave { record });
visit(&node, if path.is_empty() { "/" } else { &path })?;
let mut entries = node.child_node_entries()?;
entries.sort_by(|left, right| left.0.as_bytes().cmp(right.0.as_bytes()));
for (name, child) in entries.into_iter().rev() {
if name.starts_with(':') {
continue;
}
stack.push(Step::Visit {
node: child,
path: format!("{path}/{name}"),
});
}
}
}
}
Ok(())
}
pub(crate) 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(),
},
}
}
#[cfg(test)]
mod tests {
use super::{ReferenceEmitter, SortedReferenceSets};
use crate::external_sort::{
SortedRuns, merge_passes, peak_open_run_files, reset_sort_accounting,
};
use crate::writer::index::{IndexEntry, RunLocation, SortBudget};
struct TestDirectory {
path: std::path::PathBuf,
}
impl TestDirectory {
fn new(name: &str) -> Self {
let path = std::env::temp_dir().join(format!(
"froe-reference-merge-{name}-{}-{:?}",
std::process::id(),
std::thread::current().id()
));
let _ = std::fs::remove_dir_all(&path);
std::fs::create_dir_all(&path).expect("create the run directory");
Self { path }
}
}
impl Drop for TestDirectory {
fn drop(&mut self) {
let _ = std::fs::remove_dir_all(&self.path);
}
}
#[test]
fn the_two_reference_sets_merge_one_at_a_time() {
let directory = TestDirectory::new("one-at-a-time");
let budget = SortBudget::of_bytes(8);
let mut sets = SortedReferenceSets {
strong_runs: Some(SortedRuns::new(
RunLocation::new(&directory.path, "strong"),
budget.clone(),
)),
weak_runs: Some(SortedRuns::new(
RunLocation::new(&directory.path, "weak"),
budget.clone(),
)),
strong_seen: true,
weak_seen: true,
};
for index in 0..200u32 {
let entry = IndexEntry::new(format!("{index:08}"), "/content/a");
sets.strong_runs
.as_mut()
.expect("the strong set")
.push(entry.clone())
.expect("push");
sets.weak_runs
.as_mut()
.expect("the weak set")
.push(entry)
.expect("push");
}
reset_sort_accounting();
let strong: Vec<IndexEntry> = sets
.strong()
.expect("open the strong set")
.expect("it exists")
.collect::<crate::Result<Vec<_>>>()
.expect("read the strong set");
assert_eq!(strong.len(), 200);
let after_strong = peak_open_run_files();
assert!(
after_strong <= crate::external_sort::MAXIMUM_FAN_IN + 1,
"the strong merge held {after_strong} files"
);
let weak: Vec<IndexEntry> = sets
.weak()
.expect("open the weak set")
.expect("it exists")
.collect::<crate::Result<Vec<_>>>()
.expect("read the weak set");
assert_eq!(weak.len(), 200);
assert!(
peak_open_run_files() <= crate::external_sort::MAXIMUM_FAN_IN + 1,
"merging the weak set after the strong one took the peak to {} — the two \
sets must not be open at once",
peak_open_run_files()
);
assert!(merge_passes() >= 2, "both sets needed a reduction pass");
}
#[test]
fn an_emitter_that_saw_nothing_reports_both_sets_empty() {
let emitter = ReferenceEmitter::new(&super::EntrySink::Count);
match emitter.finish() {
super::CollectedReferences::Counted { strong, weak, .. } => {
assert_eq!((strong, weak), (0, 0));
}
super::CollectedReferences::Sorted(_) => panic!("the counting sink counts"),
}
}
}