#[cfg(test)]
mod tests {
use crate::store::Repository;
use crate::writer::compaction::CompactionKind;
use crate::writer::fault_injection::test_support::{
LUCENE_IMPORT_SCENARIO, TestDirectory, lucene_import_input, run_crash_child,
run_error_child, write_lucene_import_fixture,
};
use crate::writer::index::lucene_import::{LuceneImportOptions, lucene_import};
use crate::writer::maintenance::{CompactionOptions, MaintenanceTask, compact};
use std::collections::BTreeMap;
use std::ffi::OsString;
use std::path::Path;
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";
fn store_files(directory: &Path) -> BTreeMap<OsString, Vec<u8>> {
std::fs::read_dir(directory)
.expect("read the store directory")
.map(|entry| {
let entry = entry.expect("directory entry");
(
entry.file_name(),
std::fs::read(entry.path()).expect("read the file"),
)
})
.collect()
}
fn digest_under(store: &Path, path: &str) -> Vec<String> {
let repository = Repository::open(store).expect("open the repository");
let mut rendered = Vec::new();
crate::tooling::digest::digest_repository_excluding(&repository, &[], &[], &mut rendered)
.expect("digest");
String::from_utf8(rendered)
.expect("UTF-8")
.lines()
.filter(|line| line.starts_with(path))
.map(str::to_owned)
.collect()
}
fn checkpoint_names(store: &Path) -> Vec<String> {
let repository = Repository::open(store).expect("open the repository");
let mut names: Vec<String> = repository
.checkpoints()
.expect("read the checkpoints")
.into_iter()
.map(|(name, _)| name)
.collect();
names.sort();
names
}
fn assert_the_original_store_survived(
store: &Path,
before: &BTreeMap<OsString, Vec<u8>>,
definition_before: &[String],
checkpoints_before: &[String],
) {
let after = store_files(store);
for (name, bytes) in before {
if name == "repo.lock" {
continue;
}
assert_eq!(
after.get(name),
Some(bytes),
"{} must survive a pre-publication failure byte-identical",
Path::new(name).display()
);
}
assert_eq!(
digest_under(store, "/oak:index/lucene"),
definition_before,
"the definition and its :data must be exactly as the run found them"
);
assert_eq!(
checkpoint_names(store),
checkpoints_before,
"an import releases no checkpoint, least of all a failed one"
);
}
fn assert_a_later_compaction_reclaims_the_orphans(
store: &Path,
before: &BTreeMap<OsString, Vec<u8>>,
) {
let orphans: Vec<OsString> = store_files(store)
.into_keys()
.filter(|name| {
!before.contains_key(name)
&& Path::new(name)
.extension()
.is_some_and(|extension| extension == "tar")
})
.collect();
assert!(
!orphans.is_empty(),
"the failed run appended no archive, so there is nothing to reclaim"
);
compact(
store,
CompactionOptions::default()
.with_tasks([MaintenanceTask::Segments])
.with_compaction(CompactionKind::Full),
)
.expect("a later compaction runs over the orphan output");
let after = store_files(store);
for orphan in &orphans {
assert!(
!after.contains_key(orphan),
"{} survived the compaction that should have retired it",
Path::new(orphan).display()
);
}
let repository = Repository::open(store).expect("the store reopens after compaction");
drop(repository);
}
fn assert_whatever_was_appended_is_unreferenced(
store: &Path,
before: &BTreeMap<OsString, Vec<u8>>,
) {
let appended: Vec<OsString> = store_files(store)
.into_keys()
.filter(|name| {
!before.contains_key(name)
&& Path::new(name)
.extension()
.is_some_and(|extension| extension == "tar")
})
.collect();
compact(
store,
CompactionOptions::default()
.with_tasks([MaintenanceTask::Segments])
.with_compaction(CompactionKind::Full),
)
.expect("a later compaction runs over whatever the failed copy left");
let after = store_files(store);
for orphan in &appended {
assert!(
!after.contains_key(orphan),
"{} survived the compaction that should have retired it",
Path::new(orphan).display()
);
}
assert!(
crate::tooling::check_consistency(
store,
&["/".to_owned()],
crate::tooling::BinaryCheck::EveryBlock,
1,
)
.expect("check the store")
.has_good_revision()
);
}
fn assert_the_retry_publishes_the_same_post_state(store: &Path) {
let outcome = lucene_import(store, &LuceneImportOptions::new(lucene_import_input(store)))
.expect("the retry runs from the same input directory");
assert!(outcome.moved_the_head(), "the retry published nothing");
assert_eq!(outcome.indexes.len(), 1);
assert!(
crate::tooling::check_consistency(
store,
&["/".to_owned()],
crate::tooling::BinaryCheck::EveryBlock,
1,
)
.expect("check the store")
.has_good_revision()
);
}
#[test]
fn a_mid_copy_error_leaves_the_store_unchanged_and_names_the_file_and_offset() {
let directory = TestDirectory::new("import-mid-copy-error");
let store = write_lucene_import_fixture(&directory.path);
let before = store_files(&store);
let definition_before = digest_under(&store, "/oak:index/lucene");
let checkpoints_before = checkpoint_names(&store);
run_error_child(&store, LUCENE_IMPORT_SCENARIO, MID_FILE_COPY);
assert_the_original_store_survived(
&store,
&before,
&definition_before,
&checkpoints_before,
);
assert_whatever_was_appended_is_unreferenced(&store, &before);
}
#[test]
fn a_death_mid_copy_leaves_the_head_resolving_the_old_data() {
let directory = TestDirectory::new("import-mid-copy-crash");
let store = write_lucene_import_fixture(&directory.path);
let before = store_files(&store);
let definition_before = digest_under(&store, "/oak:index/lucene");
let checkpoints_before = checkpoint_names(&store);
run_crash_child(&store, LUCENE_IMPORT_SCENARIO, MID_FILE_COPY);
assert_the_original_store_survived(
&store,
&before,
&definition_before,
&checkpoints_before,
);
assert_the_retry_publishes_the_same_post_state(&store);
}
#[test]
fn an_error_before_head_publish_leaves_the_head_and_the_definition_as_they_were() {
let directory = TestDirectory::new("import-publish-error");
let store = write_lucene_import_fixture(&directory.path);
let before = store_files(&store);
let definition_before = digest_under(&store, "/oak:index/lucene");
let checkpoints_before = checkpoint_names(&store);
run_error_child(&store, LUCENE_IMPORT_SCENARIO, BEFORE_HEAD_PUBLISH);
assert_the_original_store_survived(
&store,
&before,
&definition_before,
&checkpoints_before,
);
assert_a_later_compaction_reclaims_the_orphans(&store, &before);
}
#[test]
fn a_death_before_head_publish_leaves_the_head_resolving_the_old_records() {
let directory = TestDirectory::new("import-publish-crash");
let store = write_lucene_import_fixture(&directory.path);
let before = store_files(&store);
let definition_before = digest_under(&store, "/oak:index/lucene");
let checkpoints_before = checkpoint_names(&store);
run_crash_child(&store, LUCENE_IMPORT_SCENARIO, BEFORE_HEAD_PUBLISH);
assert_the_original_store_survived(
&store,
&before,
&definition_before,
&checkpoints_before,
);
assert_the_retry_publishes_the_same_post_state(&store);
}
#[test]
fn an_error_after_head_publish_before_flush_leaves_the_journal_naming_the_old_head() {
let directory = TestDirectory::new("import-flush-error");
let store = write_lucene_import_fixture(&directory.path);
let before = store_files(&store);
let definition_before = digest_under(&store, "/oak:index/lucene");
let checkpoints_before = checkpoint_names(&store);
run_error_child(
&store,
LUCENE_IMPORT_SCENARIO,
AFTER_HEAD_PUBLISH_BEFORE_FLUSH,
);
assert_the_original_store_survived(
&store,
&before,
&definition_before,
&checkpoints_before,
);
assert_a_later_compaction_reclaims_the_orphans(&store, &before);
}
#[test]
fn a_death_between_head_publish_and_flush_leaves_one_resolvable_head() {
let directory = TestDirectory::new("import-flush-crash");
let store = write_lucene_import_fixture(&directory.path);
let before = store_files(&store);
let definition_before = digest_under(&store, "/oak:index/lucene");
let checkpoints_before = checkpoint_names(&store);
run_crash_child(
&store,
LUCENE_IMPORT_SCENARIO,
AFTER_HEAD_PUBLISH_BEFORE_FLUSH,
);
assert_the_original_store_survived(
&store,
&before,
&definition_before,
&checkpoints_before,
);
assert_the_retry_publishes_the_same_post_state(&store);
}
#[test]
fn a_failed_applied_state_verification_reports_rather_than_repairs() {
let directory = TestDirectory::new("import-applied");
let store = write_lucene_import_fixture(&directory.path);
let head_before = Repository::open(&store)
.expect("open")
.head_record_identifier();
run_error_child(&store, LUCENE_IMPORT_SCENARIO, BEFORE_APPLIED_VERIFICATION);
let repository = Repository::open(&store).expect("the store reopens");
assert_ne!(
repository.head_record_identifier(),
head_before,
"this boundary is after publication, so the head has moved"
);
drop(repository);
assert!(
crate::tooling::check_consistency(
&store,
&["/".to_owned()],
crate::tooling::BinaryCheck::EveryBlock,
1,
)
.expect("check the store")
.has_good_revision()
);
}
}