use crate::metadata_store::{QuarantineScope, RabsMetadataStore, StoreError};
pub trait FilesystemReality {
fn exists(&self, store_path: &str) -> bool;
fn all_paths(&self) -> Vec<String>;
}
#[derive(Debug, Clone, Default)]
pub struct SetFilesystem {
pub paths: std::collections::BTreeSet<String>,
}
impl FilesystemReality for SetFilesystem {
fn exists(&self, store_path: &str) -> bool {
self.paths.contains(store_path)
}
fn all_paths(&self) -> Vec<String> {
self.paths.iter().cloned().collect()
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Drift {
MissingLocationRepaired {
object_key: String,
store_path: String,
},
StaleTombstoneCleared {
object_key: String,
store_path: String,
},
OrphanPathReported {
store_path: String,
},
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum IncompleteState {
PublicationPinMissing {
action_key: String,
},
PublicationPinReleased {
action_key: String,
},
ServingStateMissing {
action_key: String,
},
EvidenceMissing {
action_key: String,
},
AuthorityHistoryMissing,
GenerationHighWaterMissing,
PublicationGenerationMissing {
action_key: String,
},
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ServingDecision {
Allowed,
Refused(Vec<IncompleteState>),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct StartupReport {
pub repaired: Vec<Drift>,
pub reported: Vec<Drift>,
pub serving: ServingDecision,
}
pub fn reconcile_startup(
store: &mut dyn RabsMetadataStore,
filesystem: &dyn FilesystemReality,
) -> Result<StartupReport, StoreError> {
let mut repaired = Vec::new();
let mut reported = Vec::new();
let mut claimed_paths = std::collections::BTreeSet::new();
for row in store.reconciliation_scan()? {
claimed_paths.insert(row.store_path.clone());
if !filesystem.exists(&row.store_path) {
store.remove_location_by_key(&row.object_key, &row.store_path)?;
repaired.push(Drift::MissingLocationRepaired {
object_key: row.object_key,
store_path: row.store_path,
});
}
}
let live_locations: std::collections::BTreeSet<(String, String)> = store
.reconciliation_scan()?
.into_iter()
.map(|r| (r.object_key, r.store_path))
.collect();
for tombstone in store.due_gc_tombstones(u64::MAX)? {
let key = (tombstone.object_key.clone(), tombstone.store_path.clone());
if !live_locations.contains(&key) {
store.remove_gc_tombstone(&tombstone.object_key, &tombstone.store_path)?;
repaired.push(Drift::StaleTombstoneCleared {
object_key: tombstone.object_key,
store_path: tombstone.store_path,
});
}
}
for path in filesystem.all_paths() {
if !claimed_paths.contains(&path) {
reported.push(Drift::OrphanPathReported { store_path: path });
}
}
let incomplete = authoritative_incompleteness(store)?;
Ok(StartupReport {
repaired,
reported,
serving: if incomplete.is_empty() {
ServingDecision::Allowed
} else {
ServingDecision::Refused(incomplete)
},
})
}
fn authoritative_incompleteness(
store: &mut dyn RabsMetadataStore,
) -> Result<Vec<IncompleteState>, StoreError> {
let mut incomplete = Vec::new();
let publications = store.list_publications()?;
for (action_key, pin_hex) in &publications {
if store.serving_disposition_key(action_key)?.as_deref() == Some("quarantined") {
continue;
}
match store.pin_released_by_hex(pin_hex)? {
None => incomplete.push(IncompleteState::PublicationPinMissing {
action_key: action_key.clone(),
}),
Some(true) => incomplete.push(IncompleteState::PublicationPinReleased {
action_key: action_key.clone(),
}),
Some(false) => {}
}
if !store.has_serving_state_key(action_key)? {
incomplete.push(IncompleteState::ServingStateMissing {
action_key: action_key.clone(),
});
}
if !store.has_evidence_key(action_key)? {
incomplete.push(IncompleteState::EvidenceMissing {
action_key: action_key.clone(),
});
}
}
let reset_consumed = store.highest_operator_reset()?.is_some();
let generations = store.generation_count()?;
if (!publications.is_empty() || generations > 0)
&& store.authority_count()? == 0
&& !reset_consumed
{
incomplete.push(IncompleteState::AuthorityHistoryMissing);
}
if generations > 0 && !store.has_generation_high_water()? && !reset_consumed {
incomplete.push(IncompleteState::GenerationHighWaterMissing);
}
for action_key in store.publications_missing_their_generation()? {
if store.serving_disposition_key(&action_key)?.as_deref() == Some("quarantined") {
continue;
}
incomplete.push(IncompleteState::PublicationGenerationMissing { action_key });
}
Ok(incomplete)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ResetOutcome {
pub generation: u64,
pub quarantined_actions: Vec<String>,
}
pub fn apply_operator_reset(
store: &mut dyn RabsMetadataStore,
generation: u64,
seq: u64,
) -> Result<ResetOutcome, StoreError> {
store.record_operator_reset(generation, seq)?;
let incomplete = authoritative_incompleteness(store)?;
let mut quarantined_actions = Vec::new();
for state in incomplete {
let action_key = match state {
IncompleteState::PublicationPinMissing { action_key }
| IncompleteState::PublicationPinReleased { action_key }
| IncompleteState::ServingStateMissing { action_key }
| IncompleteState::EvidenceMissing { action_key }
| IncompleteState::PublicationGenerationMissing { action_key } => action_key,
IncompleteState::AuthorityHistoryMissing
| IncompleteState::GenerationHighWaterMissing => continue,
};
if quarantined_actions.contains(&action_key) {
continue;
}
store.set_serving_disposition_key(&action_key, "quarantined")?;
store.add_quarantine(
QuarantineScope::ActionEntry,
&action_key,
"operator reset: pre-reset publication state incomplete",
)?;
quarantined_actions.push(action_key);
}
Ok(ResetOutcome {
generation,
quarantined_actions,
})
}
#[cfg(test)]
mod tests {
use super::*;
use crate::metadata_store::{
ActionEntryRow, AuthorityRow, FsqliteEngine, PublicationRow, ResultKindTag, RusqliteEngine,
SqlEngine, SqlMetadataStore, digest_key,
};
use rabs_protocol::result_identity::{DigestAlgorithm, TypedDigest};
use std::sync::atomic::{AtomicU64, Ordering};
static DB_COUNTER: AtomicU64 = AtomicU64::new(0);
fn fresh_path(tag: &str) -> std::path::PathBuf {
let n = DB_COUNTER.fetch_add(1, Ordering::SeqCst);
std::env::temp_dir().join(format!("rabs-h013-{}-{}-{}.db", std::process::id(), tag, n))
}
fn digest(domain: &'static str, tag: u8) -> TypedDigest {
TypedDigest {
algorithm: DigestAlgorithm::Sha256V1,
domain,
bytes: [tag; 32],
}
}
fn object(tag: u8) -> TypedDigest {
digest("rabs.object.sha256.v1", tag)
}
fn healthy<E: SqlEngine>(store: &mut SqlMetadataStore<E>) {
let auth = digest("rabs.authority.sha256.v1", 1);
store
.acquire_authority(&AuthorityRow {
digest: auth.clone(),
cluster_id: "c".to_owned(),
incarnation: 1,
term: 1,
acquired_seq: 1,
})
.unwrap();
let action = digest("rabs.action-key.sha256.v1", 7);
store
.upsert_action_entry(&ActionEntryRow {
action_key: action.clone(),
key_epoch: 0,
projection_epoch: 0,
})
.unwrap();
store.create_generation(&auth, 11, &action).unwrap();
store.record_attempt(20, 11, "w", 1).unwrap();
for tag in [50u8, 51] {
store.record_object(&object(tag), 64).unwrap();
store
.add_location(&object(tag), &format!("/cas/{tag}"), Some(1), "raw", true)
.unwrap();
}
store
.commit_publication(
&auth,
None,
&PublicationRow {
action_key: action,
descriptor_digest: digest("rabs.descriptor.sha256.v1", 8),
manifest_digest: object(50),
evidence_digest: object(51),
winner_generation: 11,
winner_attempt: 20,
result_kind: ResultKindTag::Success,
pin_id: 40,
pin_owner: "coordinator".to_owned(),
provisional_ancestors: Vec::new(),
},
)
.unwrap();
}
fn full_filesystem() -> SetFilesystem {
SetFilesystem {
paths: ["/cas/50", "/cas/51"]
.iter()
.map(|s| (*s).to_owned())
.collect(),
}
}
#[test]
fn h013_healthy_store_serves_with_no_drift() {
let mut store = SqlMetadataStore::open(RusqliteEngine::open_in_memory().unwrap()).unwrap();
healthy(&mut store);
let report = reconcile_startup(&mut store, &full_filesystem()).unwrap();
assert!(report.repaired.is_empty());
assert!(report.reported.is_empty());
assert_eq!(report.serving, ServingDecision::Allowed);
}
#[test]
fn h013_seeded_drift_is_repaired_and_reported() {
let mut store = SqlMetadataStore::open(RusqliteEngine::open_in_memory().unwrap()).unwrap();
healthy(&mut store);
store.record_object(&object(60), 64).unwrap();
store
.add_location(&object(60), "/cas/60", None, "raw", true)
.unwrap();
store
.add_gc_tombstone(&digest_key(&object(60)), "/cas/60", 1, 2)
.unwrap();
let mut filesystem = full_filesystem();
filesystem.paths.insert("/cas/99".to_owned());
let report = reconcile_startup(&mut store, &filesystem).unwrap();
assert_eq!(
report.repaired,
vec![
Drift::MissingLocationRepaired {
object_key: digest_key(&object(60)),
store_path: "/cas/60".to_owned(),
},
Drift::StaleTombstoneCleared {
object_key: digest_key(&object(60)),
store_path: "/cas/60".to_owned(),
},
]
);
assert_eq!(
report.reported,
vec![Drift::OrphanPathReported {
store_path: "/cas/99".to_owned(),
}]
);
assert_eq!(report.serving, ServingDecision::Allowed);
let again = reconcile_startup(&mut store, &filesystem).unwrap();
assert!(again.repaired.is_empty());
}
#[test]
fn h013_incomplete_authoritative_state_refuses_serving() {
let mut store = SqlMetadataStore::open(RusqliteEngine::open_in_memory().unwrap()).unwrap();
healthy(&mut store);
let action_key = digest_key(&digest("rabs.action-key.sha256.v1", 7));
for sql in [
"DELETE FROM pins",
"DELETE FROM action_serving_states",
"DELETE FROM action_evidence_index",
"DELETE FROM coordinator_authorities",
"DELETE FROM generation_high_water",
] {
store.engine_mut().execute(sql, &[]).unwrap();
}
let report = reconcile_startup(&mut store, &full_filesystem()).unwrap();
let ServingDecision::Refused(reasons) = report.serving else {
panic!("torn authoritative state must refuse serving");
};
assert_eq!(
reasons,
vec![
IncompleteState::PublicationPinMissing {
action_key: action_key.clone(),
},
IncompleteState::ServingStateMissing {
action_key: action_key.clone(),
},
IncompleteState::EvidenceMissing { action_key },
IncompleteState::AuthorityHistoryMissing,
IncompleteState::GenerationHighWaterMissing,
]
);
}
#[test]
fn h013_released_publication_pin_refuses_serving() {
let mut store = SqlMetadataStore::open(RusqliteEngine::open_in_memory().unwrap()).unwrap();
healthy(&mut store);
store.release_pin(40, "coordinator").unwrap();
let report = reconcile_startup(&mut store, &full_filesystem()).unwrap();
assert_eq!(
report.serving,
ServingDecision::Refused(vec![IncompleteState::PublicationPinReleased {
action_key: digest_key(&digest("rabs.action-key.sha256.v1", 7)),
}])
);
}
#[test]
fn h013_differential_reference_vs_frankensqlite() {
fn scenario<E: SqlEngine>(store: &mut SqlMetadataStore<E>) -> Vec<String> {
healthy(store);
store.record_object(&object(60), 64).unwrap();
store
.add_location(&object(60), "/cas/60", None, "raw", true)
.unwrap();
let report = reconcile_startup(store, &full_filesystem()).unwrap();
assert_eq!(report.repaired.len(), 1); assert_eq!(report.serving, ServingDecision::Allowed);
store.differential_snapshot().unwrap()
}
let mut reference =
SqlMetadataStore::open(RusqliteEngine::open(&fresh_path("ref")).unwrap()).unwrap();
let mut candidate =
SqlMetadataStore::open(FsqliteEngine::open(&fresh_path("fsq")).unwrap()).unwrap();
assert_eq!(scenario(&mut reference), scenario(&mut candidate));
}
fn roll_back<E: SqlEngine>(store: &mut SqlMetadataStore<E>) {
for sql in [
"DELETE FROM pins",
"DELETE FROM action_serving_states",
"DELETE FROM action_evidence_index",
"DELETE FROM coordinator_authorities",
"DELETE FROM generation_high_water",
] {
store.engine_mut().execute(sql, &[]).unwrap();
}
}
#[test]
fn h037_rolled_back_db_refuses_until_reset_then_serves_quarantined() {
let mut store = SqlMetadataStore::open(RusqliteEngine::open_in_memory().unwrap()).unwrap();
healthy(&mut store);
roll_back(&mut store);
let action_key = digest_key(&digest("rabs.action-key.sha256.v1", 7));
let before = reconcile_startup(&mut store, &full_filesystem()).unwrap();
assert!(matches!(before.serving, ServingDecision::Refused(_)));
let outcome = apply_operator_reset(&mut store, 1, 500).unwrap();
assert_eq!(outcome.generation, 1);
assert_eq!(outcome.quarantined_actions, vec![action_key.clone()]);
assert_eq!(
store
.serving_disposition_key(&action_key)
.unwrap()
.as_deref(),
Some("quarantined")
);
let after = reconcile_startup(&mut store, &full_filesystem()).unwrap();
assert_eq!(after.serving, ServingDecision::Allowed);
assert_eq!(
apply_operator_reset(&mut store, 1, 501),
Err(crate::metadata_store::StoreError::StaleOperatorReset)
);
assert_eq!(store.highest_operator_reset().unwrap(), Some(1));
}
#[test]
fn h037_reset_covers_history_gaps_without_publications() {
let mut store = SqlMetadataStore::open(RusqliteEngine::open_in_memory().unwrap()).unwrap();
healthy(&mut store);
for sql in [
"DELETE FROM action_publications",
"DELETE FROM pins",
"DELETE FROM action_serving_states",
"DELETE FROM action_evidence_index",
"DELETE FROM coordinator_authorities",
"DELETE FROM generation_high_water",
] {
store.engine_mut().execute(sql, &[]).unwrap();
}
let before = reconcile_startup(&mut store, &full_filesystem()).unwrap();
assert_eq!(
before.serving,
ServingDecision::Refused(vec![
IncompleteState::AuthorityHistoryMissing,
IncompleteState::GenerationHighWaterMissing,
])
);
let outcome = apply_operator_reset(&mut store, 3, 600).unwrap();
assert!(outcome.quarantined_actions.is_empty());
let after = reconcile_startup(&mut store, &full_filesystem()).unwrap();
assert_eq!(after.serving, ServingDecision::Allowed);
}
#[test]
fn h037_differential_reference_vs_frankensqlite() {
fn scenario<E: SqlEngine>(store: &mut SqlMetadataStore<E>) -> Vec<String> {
healthy(store);
roll_back(store);
assert!(matches!(
reconcile_startup(store, &full_filesystem())
.unwrap()
.serving,
ServingDecision::Refused(_)
));
apply_operator_reset(store, 2, 700).unwrap();
assert_eq!(
reconcile_startup(store, &full_filesystem())
.unwrap()
.serving,
ServingDecision::Allowed
);
store.differential_snapshot().unwrap()
}
let mut reference =
SqlMetadataStore::open(RusqliteEngine::open(&fresh_path("ref37")).unwrap()).unwrap();
let mut candidate =
SqlMetadataStore::open(FsqliteEngine::open(&fresh_path("fsq37")).unwrap()).unwrap();
assert_eq!(scenario(&mut reference), scenario(&mut candidate));
}
#[test]
fn t039_generation_fence_after_the_watermark_is_lost_and_a_reset_consumed() {
let mut store = SqlMetadataStore::open(RusqliteEngine::open_in_memory().unwrap()).unwrap();
healthy(&mut store);
let auth = digest("rabs.authority.sha256.v1", 1);
let action = digest("rabs.action-key.sha256.v1", 7);
store
.engine_mut()
.execute("DELETE FROM generation_high_water", &[])
.unwrap();
assert_eq!(
reconcile_startup(&mut store, &full_filesystem())
.unwrap()
.serving,
ServingDecision::Refused(vec![IncompleteState::GenerationHighWaterMissing]),
"a lost watermark must refuse serving before any reset"
);
let outcome = apply_operator_reset(&mut store, 1, 500).unwrap();
assert!(outcome.quarantined_actions.is_empty());
assert_eq!(
reconcile_startup(&mut store, &full_filesystem())
.unwrap()
.serving,
ServingDecision::Allowed
);
assert!(
store.create_generation(&auth, 11, &action).is_err(),
"an id that was actually minted must stay burned even without the watermark"
);
assert_eq!(
store.create_generation(&auth, 7, &action),
Err(StoreError::GenerationIdNotAboveHighWater),
"a lost watermark row must be recovered from the surviving generation \
rows, not silently read as zero (R108)"
);
assert_eq!(
store.create_generation(&auth, 11, &action),
Err(StoreError::GenerationIdNotAboveHighWater)
);
store
.create_generation(&auth, 12, &action)
.expect("minting resumes above the recovered watermark");
}
#[test]
fn t039_total_generation_state_loss_refuses_until_an_operator_resets() {
let mut store = SqlMetadataStore::open(RusqliteEngine::open_in_memory().unwrap()).unwrap();
healthy(&mut store);
let auth = digest("rabs.authority.sha256.v1", 1);
let action = digest("rabs.action-key.sha256.v1", 7);
for sql in [
"DELETE FROM action_generations",
"DELETE FROM generation_high_water",
] {
store.engine_mut().execute(sql, &[]).unwrap();
}
assert_eq!(store.generation_count().unwrap(), 0);
assert!(!store.list_publications().unwrap().is_empty());
let action_key = digest_key(&action);
assert_eq!(
reconcile_startup(&mut store, &full_filesystem())
.unwrap()
.serving,
ServingDecision::Refused(vec![IncompleteState::PublicationGenerationMissing {
action_key: action_key.clone(),
}]),
"total generation-state loss must refuse serving, never present as an \
ordinary cold cache (R113)"
);
assert_eq!(store.highest_operator_reset().unwrap(), None);
let outcome = apply_operator_reset(&mut store, 1, 500).unwrap();
assert_eq!(outcome.quarantined_actions, vec![action_key.clone()]);
assert_eq!(
store
.serving_disposition_key(&action_key)
.unwrap()
.as_deref(),
Some("quarantined")
);
assert_eq!(
reconcile_startup(&mut store, &full_filesystem())
.unwrap()
.serving,
ServingDecision::Allowed,
"the reset opens a new lineage and serving resumes"
);
store
.create_generation(&auth, 11, &action)
.expect("a post-reset lineage may reuse the id space it declared new");
}
}