use std::path::{Path, PathBuf};
use crate::content::node::NodeState;
use crate::error::{Error, Result};
use crate::index::lanes::AsyncLanes;
use crate::index::lucene::check::LocalIndexDirectory;
use crate::index::lucene::layout::{
INDEX_DETAILS_FILE_NAME, INDEXER_INFO_FILE_NAME, IndexDetails, IndexerInfo,
};
use crate::index::{IndexDefinition, IndexType};
use crate::segment::record::RecordIdentifier;
use crate::store::Repository;
pub const INDEX_DEFINITIONS_FILE_NAME: &str = "index-definitions.json";
#[derive(Clone, Debug)]
pub struct LuceneImportOptions {
input: PathBuf,
indexes: Vec<String>,
}
impl LuceneImportOptions {
#[must_use]
pub fn new(input: PathBuf) -> Self {
Self {
input,
indexes: Vec::new(),
}
}
#[must_use]
pub fn with_indexes(mut self, indexes: impl IntoIterator<Item = String>) -> Self {
self.indexes = indexes.into_iter().collect();
self
}
#[must_use]
pub fn input(&self) -> &Path {
&self.input
}
#[must_use]
pub fn indexes(&self) -> &[String] {
&self.indexes
}
}
#[derive(Clone, PartialEq, Eq, Debug)]
pub struct PlannedImport {
pub path: String,
pub directory: PathBuf,
pub mappings: Vec<(String, PathBuf)>,
pub skipped_mappings: Vec<(String, String)>,
pub file_count: usize,
pub byte_count: u64,
pub reindex_count: u64,
}
#[derive(Clone, PartialEq, Eq, Debug)]
#[non_exhaustive]
pub struct LuceneImportPlan {
pub directory: PathBuf,
pub input: PathBuf,
pub checkpoint: String,
pub imports: Vec<PlannedImport>,
pub definitions_without_directories: Vec<String>,
}
impl LuceneImportPlan {
#[must_use]
pub fn is_empty(&self) -> bool {
self.imports.is_empty()
}
#[must_use]
pub fn file_count(&self) -> usize {
self.imports.iter().map(|import| import.file_count).sum()
}
}
pub const SUGGEST_SKIP_REASON: &str = "froe never imports suggester data; Oak's own writer rebuilds it when its \
lastUpdated is missing";
pub fn plan_lucene_import(
directory: &Path,
options: &LuceneImportOptions,
) -> Result<LuceneImportPlan> {
let directory = crate::writer::maintenance::canonical_repository_directory(directory)?;
let repository = Repository::open(&directory)?;
let plan = build_plan(&directory, &repository, options)?;
drop(repository);
Ok(plan)
}
pub(crate) fn build_plan(
directory: &Path,
repository: &Repository,
options: &LuceneImportOptions,
) -> Result<LuceneImportPlan> {
let input = options.input();
let info = read_indexer_info(input)?;
let local = read_local_directories(input)?;
let definitions_file = read_definitions_file(input)?;
let content_root = repository.content_root()?;
let lanes = AsyncLanes::read(&content_root).map_err(index_error)?;
let checkpoint_root = resolve_checkpoint_root(repository, &info.checkpoint)?;
let mut imports = Vec::new();
let mut refusals: Vec<String> = Vec::new();
for (index_path, local_directory) in &local {
if !options.indexes().is_empty()
&& !options.indexes().iter().any(|wanted| wanted == index_path)
{
continue;
}
match plan_one(
repository,
&content_root,
&lanes,
checkpoint_root,
&info.checkpoint,
index_path,
local_directory,
&definitions_file,
) {
Ok(import) => imports.push(import),
Err(refusal) => refusals.push(refusal.to_string()),
}
}
if !refusals.is_empty() {
return Err(Error::InvalidFormat {
details: format!(
"{} of the {} definitions in {} cannot be imported:\n {}",
refusals.len(),
local.len(),
input.display(),
refusals.join("\n ")
),
});
}
for wanted in options.indexes() {
if !imports.iter().any(|import| &import.path == wanted) {
return Err(Error::InvalidFormat {
details: format!("{wanted} has no index directory in {}", input.display()),
});
}
}
let definitions_without_directories = definitions_file
.definitions
.keys()
.filter(|path| !local.iter().any(|(index_path, _)| &index_path == path))
.cloned()
.collect();
imports.sort_by(|left, right| left.path.cmp(&right.path));
Ok(LuceneImportPlan {
directory: directory.to_owned(),
input: input.to_owned(),
checkpoint: info.checkpoint,
imports,
definitions_without_directories,
})
}
#[allow(
clippy::too_many_arguments,
reason = "every argument is a distinct fact the rule needs; bundling them would hide \
which check reads which"
)]
fn plan_one(
repository: &Repository,
content_root: &NodeState<'_>,
lanes: &AsyncLanes,
checkpoint_root: RecordIdentifier,
checkpoint: &str,
index_path: &str,
local_directory: &Path,
definitions_file: &crate::index::definitions_json_reader::ParsedDefinitions,
) -> Result<PlannedImport> {
let _ = content_root;
let Some(node) = repository.node_at_path(index_path)? else {
return Err(Error::InvalidFormat {
details: format!(
"{index_path} names no node in the store; adding a definition is oak-run's \
job, not an import's"
),
});
};
let definition = IndexDefinition::read(&node, index_path).map_err(index_error)?;
if definition.index_type != Some(IndexType::Lucene) {
return Err(Error::InvalidFormat {
details: format!("{index_path} is not a lucene definition"),
});
}
let Some(lane) = definition.lane.clone() else {
return Err(Error::InvalidFormat {
details: format!(
"{index_path} is synchronous, and oak-run's own importer never completes \
that case; there is no Oak behaviour for froe to match"
),
});
};
if definition.indexing_mode.synchronous_synonym {
return Err(Error::InvalidFormat {
details: format!(
"{index_path} is hybrid — it lists sync beside lane {lane} — and Oak keeps \
its synchronous property index in the hidden :property-index child, which \
froe does not build"
),
});
}
let lane_checkpoint = lanes
.lane(&lane)
.and_then(|state| state.checkpoint.clone())
.ok_or_else(|| Error::InvalidFormat {
details: format!("{index_path} indexes on lane {lane}, which has no state on /:async"),
})?;
let lane_root = resolve_checkpoint_root(repository, &lane_checkpoint)?;
if lane_root != checkpoint_root {
return Err(Error::InvalidFormat {
details: format!(
"{index_path} was built at checkpoint {checkpoint} (root {checkpoint_root}), \
but lane {lane} resumes from {lane_checkpoint} (root {lane_root}); rebuild \
at the lane's own checkpoint"
),
});
}
let Some(parsed) = definitions_file.definitions.get(index_path) else {
return Err(Error::InvalidFormat {
details: format!(
"{index_path} has an index directory but no entry in \
{INDEX_DEFINITIONS_FILE_NAME}; oak-run's importer requires the file to \
describe every directory"
),
});
};
let materialized = crate::writer::index::lucene_import::materialize::materialize(parsed)?;
let verdict = super::drift::compare(&materialized.node(), &node).map_err(index_error)?;
if !verdict.is_clean() {
return Err(index_error(super::drift::refusal(index_path, &verdict)));
}
let Mappings {
written: mappings,
skipped,
file_count,
byte_count,
} = read_mappings(index_path, local_directory)?;
for (jcr_name, source) in &mappings {
refuse_an_incoherent_directory(index_path, jcr_name, source)?;
}
let file_count_property = definitions_file
.definitions
.get(index_path)
.and_then(|parsed| parsed.properties.get("reindexCount"))
.and_then(|property| property.values.first())
.and_then(|text| text.parse::<u64>().ok())
.unwrap_or(0);
Ok(PlannedImport {
path: index_path.to_owned(),
directory: local_directory.to_owned(),
mappings,
skipped_mappings: skipped,
file_count,
byte_count,
reindex_count: file_count_property + 1,
})
}
fn refuse_an_incoherent_directory(index_path: &str, jcr_name: &str, source: &Path) -> Result<()> {
let report =
crate::index::lucene::check::check_structure(&LocalIndexDirectory::new(source.to_owned()))
.map_err(index_error)?;
if report.is_coherent() {
return Ok(());
}
let reason = if let Some((file, details)) = report.unreadable_files.first() {
format!("{file} does not read: {details}")
} else if let Some(missing) = report.missing_files.first() {
format!("the commit names {missing}, which the directory does not hold")
} else if let Some(extra) = report.unreferenced_files.first() {
format!("{extra} is in the directory and no segment names it")
} else {
"it holds no commit file".to_owned()
};
Err(Error::InvalidFormat {
details: format!(
"{index_path}'s {jcr_name} in {} is not a coherent Lucene index: {reason}; \
froe refuses it before copying a byte",
source.display()
),
})
}
struct Mappings {
written: Vec<(String, PathBuf)>,
skipped: Vec<(String, String)>,
file_count: usize,
byte_count: u64,
}
fn read_mappings(index_path: &str, local_directory: &Path) -> Result<Mappings> {
let details = read_index_details(local_directory)?;
let mut resolved = Mappings {
written: Vec::new(),
skipped: Vec::new(),
file_count: 0,
byte_count: 0,
};
for (filesystem_name, jcr_name) in &details.directory_mappings {
let source = local_directory.join(filesystem_name);
if crate::index::lucene::is_suggest_directory_name(jcr_name) {
resolved
.skipped
.push((jcr_name.clone(), SUGGEST_SKIP_REASON.to_owned()));
continue;
}
if !crate::index::lucene::is_index_directory_name(jcr_name) {
return Err(Error::InvalidFormat {
details: format!(
"{index_path}'s {INDEX_DETAILS_FILE_NAME} maps {filesystem_name} to \
{jcr_name}, which is neither an index nor a suggester directory"
),
});
}
for entry in std::fs::read_dir(&source)? {
let entry = entry?;
if entry.file_type()?.is_file() {
resolved.file_count += 1;
resolved.byte_count += entry.metadata()?.len();
}
}
resolved.written.push((jcr_name.clone(), source));
}
Ok(resolved)
}
fn read_indexer_info(input: &Path) -> Result<IndexerInfo> {
let path = input.join(INDEXER_INFO_FILE_NAME);
let content = std::fs::read_to_string(&path).map_err(|error| Error::InvalidFormat {
details: format!(
"{} could not be read ({error}); an import needs the checkpoint the index was \
built at, and a dump writes none when it could not name one",
path.display()
),
})?;
IndexerInfo::parse(&content)
}
fn read_index_details(local: &Path) -> Result<IndexDetails> {
let path = local.join(INDEX_DETAILS_FILE_NAME);
let content = std::fs::read_to_string(&path)?;
IndexDetails::parse(&content)
}
fn read_local_directories(input: &Path) -> Result<Vec<(String, PathBuf)>> {
let mut found = Vec::new();
for entry in std::fs::read_dir(input)? {
let entry = entry?;
if !entry.file_type()?.is_dir() {
continue;
}
let path = entry.path();
if !path.join(INDEX_DETAILS_FILE_NAME).is_file() {
continue;
}
let details = read_index_details(&path)?;
if details.index_path.is_empty() {
return Err(Error::InvalidFormat {
details: format!(
"{}/{INDEX_DETAILS_FILE_NAME} names no indexPath",
path.display()
),
});
}
found.push((details.index_path, path));
}
found.sort();
Ok(found)
}
fn read_definitions_file(
input: &Path,
) -> Result<crate::index::definitions_json_reader::ParsedDefinitions> {
let path = input.join(INDEX_DEFINITIONS_FILE_NAME);
let content = std::fs::read_to_string(&path).map_err(|error| Error::InvalidFormat {
details: format!(
"{} could not be read ({error}); oak-run's importer requires it and froe reads \
it back to check the definition has not drifted",
path.display()
),
})?;
crate::index::definitions_json_reader::parse(&content).map_err(index_error)
}
fn resolve_checkpoint_root(repository: &Repository, checkpoint: &str) -> Result<RecordIdentifier> {
let root = repository
.checkpoints()?
.into_iter()
.find(|(name, _)| name == checkpoint)
.map(|(_, node)| node.child_node("root"))
.transpose()?
.flatten()
.ok_or_else(|| Error::InvalidFormat {
details: format!("checkpoint {checkpoint} does not resolve in this store"),
})?;
Ok(root.record_identifier())
}
fn index_error(error: crate::index::IndexError) -> Error {
match error {
crate::index::IndexError::Record(source) => source,
other => Error::InvalidFormat {
details: other.to_string(),
},
}
}