use crate::parallel::worker_count;
use std::sync::Mutex;
use std::sync::atomic::{AtomicUsize, Ordering};
use crate::content::node::{NodeState, PropertyValues};
use crate::content::property::PropertyValue;
use crate::content::provider::SegmentProvider;
use crate::error::Result;
use crate::progress::{DiscardedProgress, ProgressObserver, Step, WorkUnit};
use crate::segment::identifier::SegmentIdentifier;
use crate::segment::record::{RecordIdentifier, RecordType};
use crate::store::{ArchiveSet, open_all_archives_with_progress};
#[derive(Debug, Clone, Default)]
pub struct SearchQuery {
pub has_properties: Vec<String>,
pub has_children: Vec<String>,
pub property_values: Vec<(String, String)>,
}
impl SearchQuery {
#[must_use]
pub fn is_empty(&self) -> bool {
self.has_properties.is_empty()
&& self.has_children.is_empty()
&& self.property_values.is_empty()
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct NodeMatch {
pub record: RecordIdentifier,
pub stable_identifier: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SearchOutcome {
pub matches: Vec<NodeMatch>,
pub unreadable_nodes: u64,
}
pub fn search_nodes(
directory: &std::path::Path,
query: &SearchQuery,
limit: usize,
) -> Result<SearchOutcome> {
search_nodes_with_progress(directory, query, limit, &mut DiscardedProgress)
}
pub fn search_nodes_with_progress(
directory: &std::path::Path,
query: &SearchQuery,
limit: usize,
observer: &mut dyn ProgressObserver,
) -> Result<SearchOutcome> {
let mut matches = Vec::new();
let unreadable_nodes = search_nodes_visiting(directory, query, observer, &mut |found| {
matches.push(found);
if limit != 0 && matches.len() >= limit {
std::ops::ControlFlow::Break(())
} else {
std::ops::ControlFlow::Continue(())
}
})?;
Ok(SearchOutcome {
matches,
unreadable_nodes,
})
}
pub fn search_nodes_visiting(
directory: &std::path::Path,
query: &SearchQuery,
observer: &mut dyn ProgressObserver,
visit: &mut dyn FnMut(NodeMatch) -> std::ops::ControlFlow<()>,
) -> Result<u64> {
let archives = open_all_archives_with_progress(directory, observer)?;
let provider = ArchiveSet::new(archives);
let identifiers: Vec<SegmentIdentifier> = provider.distinct_segment_identifiers().collect();
let mut unreadable_nodes = 0u64;
observer.step_began(
&Step::new("searching segments", WorkUnit::Segments)
.with_total(crate::progress::count(identifiers.len())),
);
let window_size = worker_count(identifiers.len());
let mut searched_segments = 0usize;
for window in identifiers.chunks(window_size.max(1)) {
observer.step_advanced(crate::progress::count(searched_segments));
let window_hits = search_segment_window(&provider, window, query);
searched_segments += window.len();
for hits in window_hits {
unreadable_nodes += hits.unreadable_nodes;
for found in hits.matches {
if visit(found).is_break() {
observer.step_advanced(crate::progress::count(searched_segments));
observer.step_ended();
return Ok(unreadable_nodes);
}
}
}
}
observer.step_advanced(crate::progress::count(searched_segments));
observer.step_ended();
Ok(unreadable_nodes)
}
struct SegmentSearchHits {
matches: Vec<NodeMatch>,
unreadable_nodes: u64,
}
fn search_segment_window(
provider: &ArchiveSet,
identifiers: &[SegmentIdentifier],
query: &SearchQuery,
) -> Vec<SegmentSearchHits> {
let workers = worker_count(identifiers.len());
let next = AtomicUsize::new(0);
let slots: Vec<Mutex<Option<SegmentSearchHits>>> =
identifiers.iter().map(|_| Mutex::new(None)).collect();
std::thread::scope(|scope| {
for _ in 1..workers {
scope.spawn(|| search_window_worker(provider, identifiers, query, &next, &slots));
}
search_window_worker(provider, identifiers, query, &next, &slots);
});
slots
.into_iter()
.map(|slot| {
slot.into_inner()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.expect("every window slot is filled")
})
.collect()
}
fn search_window_worker(
provider: &ArchiveSet,
identifiers: &[SegmentIdentifier],
query: &SearchQuery,
next: &AtomicUsize,
slots: &[Mutex<Option<SegmentSearchHits>>],
) {
loop {
let position = next.fetch_add(1, Ordering::Relaxed);
let Some(identifier) = identifiers.get(position) else {
return;
};
let hits = search_one_segment(provider, *identifier, query);
*slots[position]
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(hits);
}
}
fn search_one_segment(
provider: &ArchiveSet,
segment_identifier: SegmentIdentifier,
query: &SearchQuery,
) -> SegmentSearchHits {
if segment_identifier.is_bulk_segment() {
return SegmentSearchHits {
matches: Vec::new(),
unreadable_nodes: 0,
};
}
let Ok(view) = provider.segment(segment_identifier) else {
return SegmentSearchHits {
matches: Vec::new(),
unreadable_nodes: 0,
};
};
let mut unreadable_nodes = 0u64;
let mut node_records = Vec::new();
for entry in view.structure.record_table() {
match entry.record_type() {
Some(RecordType::Node) => node_records.push(entry.record_number),
Some(_) => {}
None => unreadable_nodes += 1,
}
}
let mut matches = Vec::new();
for record_number in node_records {
let record = RecordIdentifier::new(segment_identifier, record_number);
match node_matches(provider, record, query) {
Ok(false) => {}
Ok(true) => {
let stable_identifier = NodeState::new(provider, record)
.stable_identifier()
.unwrap_or_else(|_| record.to_string());
matches.push(NodeMatch {
record,
stable_identifier,
});
}
Err(_) => unreadable_nodes += 1,
}
}
SegmentSearchHits {
matches,
unreadable_nodes,
}
}
fn node_matches(
provider: &dyn SegmentProvider,
record: RecordIdentifier,
query: &SearchQuery,
) -> Result<bool> {
let node = NodeState::new(provider, record);
for child_name in &query.has_children {
if node.child_node(child_name)?.is_none() {
return Ok(false);
}
}
if query.has_properties.is_empty() && query.property_values.is_empty() {
return Ok(true);
}
let properties = node.properties()?;
for property_name in &query.has_properties {
if !properties
.iter()
.any(|property| property.name == *property_name)
{
return Ok(false);
}
}
for (property_name, expected_value) in &query.property_values {
let Some(property) = properties
.iter()
.find(|property| property.name == *property_name)
else {
return Ok(false);
};
if !property_has_string_value(&property.values, expected_value) {
return Ok(false);
}
}
Ok(true)
}
fn property_has_string_value(values: &PropertyValues, expected: &str) -> bool {
let matches_value =
|value: &PropertyValue| value.as_text().is_some_and(|text| text == expected);
match values {
PropertyValues::Single(value) => matches_value(value),
PropertyValues::Multiple(values) => values.iter().any(matches_value),
}
}
#[cfg(test)]
mod tests {
use super::{SearchQuery, search_nodes};
use crate::writer::record_writer::{ChildNodesToWrite, PropertyToWrite, PropertyValuesToWrite};
use crate::writer::store_writer::WritableRepository;
struct TestDirectory {
path: std::path::PathBuf,
}
impl TestDirectory {
fn new(name: &str) -> Self {
let path =
std::env::temp_dir().join(format!("froe-search-{name}-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&path);
Self { path }
}
}
impl Drop for TestDirectory {
fn drop(&mut self) {
let _ = std::fs::remove_dir_all(&self.path);
}
}
fn populate(directory: &std::path::Path) {
let store = WritableRepository::open(directory).expect("open");
let generation = store.writing_generation().expect("generation");
let mut writer = store.record_writer(generation);
let marked_value = writer.write_string("target").expect("value");
let marked = writer
.write_node(
Some("nt:unstructured"),
&[],
&ChildNodesToWrite::Zero,
&[PropertyToWrite {
name: "marker".to_owned(),
property_type: crate::content::property::PropertyType::String,
values: PropertyValuesToWrite::Single(marked_value),
}],
)
.expect("marked");
let plain = writer
.write_node(Some("nt:unstructured"), &[], &ChildNodesToWrite::Zero, &[])
.expect("plain");
let content = writer
.write_node(
Some("nt:unstructured"),
&[],
&ChildNodesToWrite::Many(vec![
("marked".to_owned(), marked),
("plain".to_owned(), plain),
]),
&[],
)
.expect("content");
let root = writer
.write_node(
None,
&[],
&ChildNodesToWrite::One {
name: "content".to_owned(),
node: content,
},
&[],
)
.expect("root");
let head = writer
.write_node(
None,
&[],
&ChildNodesToWrite::One {
name: "root".to_owned(),
node: root,
},
&[],
)
.expect("super root");
writer.finish().expect("finish");
let previous = store.head();
assert!(store.compare_and_set_head(previous, head));
store.close().expect("close");
}
#[test]
fn finds_nodes_by_property_presence() {
let directory = TestDirectory::new("by-property");
populate(&directory.path);
let query = SearchQuery {
has_properties: vec!["marker".to_owned()],
..SearchQuery::default()
};
let outcome = search_nodes(&directory.path, &query, 0).expect("search");
assert_eq!(outcome.matches.len(), 1, "one node has the marker property");
assert_eq!(outcome.unreadable_nodes, 0);
}
#[test]
fn finds_nodes_by_property_value() {
let directory = TestDirectory::new("by-value");
populate(&directory.path);
let query = SearchQuery {
property_values: vec![("marker".to_owned(), "target".to_owned())],
..SearchQuery::default()
};
assert_eq!(
search_nodes(&directory.path, &query, 0)
.expect("search")
.matches
.len(),
1
);
let wrong_value = SearchQuery {
property_values: vec![("marker".to_owned(), "other".to_owned())],
..SearchQuery::default()
};
assert!(
search_nodes(&directory.path, &wrong_value, 0)
.expect("search")
.matches
.is_empty()
);
}
#[test]
fn finds_nodes_by_child_presence() {
let directory = TestDirectory::new("by-child");
populate(&directory.path);
let query = SearchQuery {
has_children: vec!["marked".to_owned()],
..SearchQuery::default()
};
let outcome = search_nodes(&directory.path, &query, 0).expect("search");
assert_eq!(
outcome.matches.len(),
1,
"the content node has the marked child"
);
}
#[test]
fn the_limit_bounds_the_result_count() {
let directory = TestDirectory::new("limit");
populate(&directory.path);
let query = SearchQuery {
has_properties: vec!["jcr:primaryType".to_owned()],
..SearchQuery::default()
};
let limited = search_nodes(&directory.path, &query, 2).expect("search");
assert_eq!(limited.matches.len(), 2, "the limit caps the results");
}
#[test]
fn matches_are_emitted_in_the_same_order_on_every_run() {
let directory = TestDirectory::new("stable-order");
populate(&directory.path);
let query = SearchQuery {
has_properties: vec!["jcr:primaryType".to_owned()],
..SearchQuery::default()
};
let first = search_nodes(&directory.path, &query, 0).expect("first");
let second = search_nodes(&directory.path, &query, 0).expect("second");
assert_eq!(
first.matches, second.matches,
"windowed parallel search must visit in archive probe order"
);
assert!(!first.matches.is_empty());
}
}