use rabs_protocol::result_identity::TypedDigest;
use rabs_protocol::serving::ServingValidity;
use crate::metadata_store::{
DivergenceIncidentRow, QuarantineScope, RabsMetadataStore, SqlValue, StoreError, digest_key,
};
use crate::publication::DISPOSITION_PRESENTATION_QUARANTINED;
use crate::trust_evidence::{DISPOSITION_QUARANTINED, require_active_authority};
pub const SERVABLE_DISPOSITION: &str = "servable";
pub(crate) fn action_quarantine_present(
store: &mut dyn RabsMetadataStore,
action_key: &str,
) -> Result<bool, StoreError> {
Ok(!store
.query(
"SELECT 1 FROM quarantines WHERE scope = 'action-entry' AND subject = ?1 LIMIT 1",
&[SqlValue::Text(action_key.to_owned())],
)?
.is_empty())
}
pub(crate) fn divergence_quarantine_disposition(
store: &mut dyn RabsMetadataStore,
action_key: &str,
) -> Result<Option<&'static str>, StoreError> {
let key = [SqlValue::Text(action_key.to_owned())];
if !store
.query(
"SELECT 1 FROM divergence_incidents \
WHERE action_key = ?1 AND class != 'observable-only' LIMIT 1",
&key,
)?
.is_empty()
{
return Ok(Some(DISPOSITION_QUARANTINED));
}
Ok((!store
.query(
"SELECT 1 FROM divergence_incidents WHERE action_key = ?1 LIMIT 1",
&key,
)?
.is_empty())
.then_some(DISPOSITION_PRESENTATION_QUARANTINED))
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ServeDecision {
NoRecord,
NotServable {
disposition: String,
},
Blocked {
references: Vec<(String, String)>,
},
ExpiredClockEpoch,
ExpiredClockRollback,
ExpiredTtl,
Servable,
}
pub fn serving_gate(
store: &mut dyn RabsMetadataStore,
action_key: &str,
now_unix_micros: i64,
now_epoch: u64,
) -> Result<ServeDecision, StoreError> {
let Some(record) = store.serving_record(action_key)? else {
return Ok(ServeDecision::NoRecord);
};
if record.disposition != SERVABLE_DISPOSITION {
return Ok(ServeDecision::NotServable {
disposition: record.disposition,
});
}
if !record.blocking.is_empty() {
return Ok(ServeDecision::Blocked {
references: record.blocking,
});
}
if action_quarantine_present(store, action_key)? {
return Ok(ServeDecision::Blocked {
references: vec![("action-entry".to_owned(), action_key.to_owned())],
});
}
if let Some(disposition) = divergence_quarantine_disposition(store, action_key)? {
return Ok(ServeDecision::NotServable {
disposition: disposition.to_owned(),
});
}
if !record.validity.still_valid(now_unix_micros, now_epoch) {
let validity: &ServingValidity = &record.validity;
if now_epoch != validity.coordinator_clock_epoch {
return Ok(ServeDecision::ExpiredClockEpoch);
}
if now_unix_micros < validity.evaluated_at_unix_micros {
return Ok(ServeDecision::ExpiredClockRollback);
}
return Ok(ServeDecision::ExpiredTtl);
}
Ok(ServeDecision::Servable)
}
pub const DEFAULT_REVALIDATION_TTL_MICROS: u64 = 60_000_000;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum TerminalOutcome {
Exit(i32),
Oom,
Signal(i32),
Cancelled,
Timeout,
WorkerLost,
TransportFailed,
Panicked,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum FailureRefusal {
ZeroExitIsNotAFailure,
NotDeterministic(&'static str),
CaptureIncomplete,
InputsNotClosed,
UndeclaredSideEffects,
ClassPolicyRefused,
TrustPolicyRefused,
}
impl std::fmt::Display for FailureRefusal {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::ZeroExitIsNotAFailure => write!(f, "exit 0 is a success candidate"),
Self::NotDeterministic(class) => write!(f, "nondeterministic outcome: {class}"),
Self::CaptureIncomplete => write!(f, "canonical capture incomplete"),
Self::InputsNotClosed => write!(f, "input closure not closed"),
Self::UndeclaredSideEffects => write!(f, "undeclared side effects"),
Self::ClassPolicyRefused => write!(f, "class policy refuses failure caching"),
Self::TrustPolicyRefused => write!(f, "trust policy refuses failure serving"),
}
}
}
impl std::error::Error for FailureRefusal {}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct FailureAdmission {
pub normalized_exit: i32,
}
pub fn classify_deterministic_failure(
outcome: TerminalOutcome,
capture_complete: bool,
inputs_closed: bool,
no_undeclared_side_effects: bool,
class_policy_permits: bool,
trust_policy_permits: bool,
) -> Result<FailureAdmission, FailureRefusal> {
let normalized_exit = match outcome {
TerminalOutcome::Exit(0) => {
return Err(FailureRefusal::ZeroExitIsNotAFailure);
}
TerminalOutcome::Exit(code) => code,
TerminalOutcome::Oom => return Err(FailureRefusal::NotDeterministic("oom")),
TerminalOutcome::Signal(_) => {
return Err(FailureRefusal::NotDeterministic("signal"));
}
TerminalOutcome::Cancelled => return Err(FailureRefusal::NotDeterministic("cancelled")),
TerminalOutcome::Timeout => return Err(FailureRefusal::NotDeterministic("timeout")),
TerminalOutcome::WorkerLost => return Err(FailureRefusal::NotDeterministic("worker-loss")),
TerminalOutcome::TransportFailed => {
return Err(FailureRefusal::NotDeterministic("transport-failure"));
}
TerminalOutcome::Panicked => return Err(FailureRefusal::NotDeterministic("panic")),
};
if !capture_complete {
return Err(FailureRefusal::CaptureIncomplete);
}
if !inputs_closed {
return Err(FailureRefusal::InputsNotClosed);
}
if !no_undeclared_side_effects {
return Err(FailureRefusal::UndeclaredSideEffects);
}
if !class_policy_permits {
return Err(FailureRefusal::ClassPolicyRefused);
}
if !trust_policy_permits {
return Err(FailureRefusal::TrustPolicyRefused);
}
Ok(FailureAdmission { normalized_exit })
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum RevalidationVerdict {
IdenticalEvidenceAppended {
new_revision: u64,
},
SoundnessIncidentQuarantined {
incident_seq: u64,
new_revision: u64,
},
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum RevalidationError {
NoServingRecord,
ActionKeyMismatch,
ServingBlocked,
RevisionExhausted,
StaleRevision {
stored: u64,
},
Store(String),
}
impl std::fmt::Display for RevalidationError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::NoServingRecord => write!(f, "no serving record to revalidate"),
Self::ActionKeyMismatch => write!(f, "revalidation action keys do not match"),
Self::ServingBlocked => write!(f, "serving is blocked; revalidation is not a repair"),
Self::RevisionExhausted => write!(f, "serving revision exhausted"),
Self::StaleRevision { stored } => {
write!(f, "serving revision moved to {stored}; retry")
}
Self::Store(error) => write!(f, "store: {error}"),
}
}
}
impl std::error::Error for RevalidationError {}
#[allow(clippy::too_many_arguments)]
pub fn apply_revalidation(
store: &mut dyn RabsMetadataStore,
authority: &TypedDigest,
action_key_str: &str,
action_key_typed: &TypedDigest,
expected_revision: u64,
attempt: u128,
generation: u128,
published_signature: &str,
revalidated_signature: &str,
committed_manifest_key: &str,
candidate_manifest_key: &str,
candidate_evidence_key: &str,
now_unix_micros: i64,
now_epoch: u64,
ttl_micros: Option<u64>,
) -> Result<RevalidationVerdict, RevalidationError> {
let record = store
.serving_record(action_key_str)
.map_err(|e| RevalidationError::Store(format!("{e:?}")))?
.ok_or(RevalidationError::NoServingRecord)?;
if record.state_revision != expected_revision {
return Err(RevalidationError::StaleRevision {
stored: record.state_revision,
});
}
if digest_key(action_key_typed) != action_key_str {
return Err(RevalidationError::ActionKeyMismatch);
}
require_active_authority(store, authority)
.map_err(|e| RevalidationError::Store(format!("{e:?}")))?;
let new_revision = expected_revision
.checked_add(1)
.ok_or(RevalidationError::RevisionExhausted)?;
let ttl = ttl_micros.unwrap_or(DEFAULT_REVALIDATION_TTL_MICROS);
let validity = ServingValidity {
evaluated_at_unix_micros: now_unix_micros,
maximum_age_micros: Some(ttl),
clock_uncertainty_micros: 0,
coordinator_clock_epoch: now_epoch,
};
if published_signature == revalidated_signature {
if record.disposition != SERVABLE_DISPOSITION
|| !record.blocking.is_empty()
|| action_quarantine_present(store, action_key_str)
.map_err(|e| RevalidationError::Store(format!("{e:?}")))?
|| divergence_quarantine_disposition(store, action_key_str)
.map_err(|e| RevalidationError::Store(format!("{e:?}")))?
.is_some()
{
return Err(RevalidationError::ServingBlocked);
}
store
.record_verification_sample(action_key_typed, attempt, true, new_revision)
.map_err(|e| RevalidationError::Store(format!("{e:?}")))?;
store
.put_serving_record(
authority,
action_key_str,
SERVABLE_DISPOSITION,
new_revision,
&validity,
&[],
)
.map_err(|e| RevalidationError::Store(format!("{e:?}")))?;
Ok(RevalidationVerdict::IdenticalEvidenceAppended { new_revision })
} else {
if !action_quarantine_present(store, action_key_str)
.map_err(|e| RevalidationError::Store(format!("{e:?}")))?
{
store
.add_quarantine(
QuarantineScope::ActionEntry,
action_key_str,
"k007-soundness-incident",
)
.map_err(|e| RevalidationError::Store(format!("{e:?}")))?;
}
let incidents = store
.list_divergence_incidents(action_key_str)
.map_err(|e| RevalidationError::Store(format!("{e:?}")))?;
let seq = incidents
.iter()
.map(|incident| incident.seq)
.max()
.map_or(Some(new_revision), |seq| {
seq.checked_add(1).map(|next| next.max(new_revision))
})
.ok_or_else(|| {
RevalidationError::Store("divergence incident sequence exhausted".to_owned())
})?;
store
.record_divergence_incident(
authority,
&DivergenceIncidentRow {
action_key: action_key_str.to_owned(),
seq,
class: "soundness".to_owned(),
committed_manifest_key: committed_manifest_key.to_owned(),
candidate_manifest_key: candidate_manifest_key.to_owned(),
candidate_evidence_key: candidate_evidence_key.to_owned(),
candidate_pin_hex: String::new(),
generation_hex: format!("{generation:x}"),
attempt_hex: format!("{attempt:x}"),
detail: format!(
"revalidation diverged: published {published_signature}, \
revalidated {revalidated_signature}"
),
},
)
.map_err(|e| RevalidationError::Store(format!("{e:?}")))?;
let mut blocking = Vec::new();
for (scope, subject) in &record.blocking {
let scope = match scope.as_str() {
"location" => QuarantineScope::Location,
"logical-object" => QuarantineScope::LogicalObject,
"action-entry" => QuarantineScope::ActionEntry,
_ => {
return Err(RevalidationError::Store(format!(
"unknown blocking quarantine scope: {scope}"
)));
}
};
blocking.push((scope, subject.clone()));
}
if !blocking.iter().any(|(scope, subject)| {
scope == &QuarantineScope::ActionEntry && subject == action_key_str
}) {
blocking.push((QuarantineScope::ActionEntry, action_key_str.to_owned()));
}
store
.put_serving_record(
authority,
action_key_str,
"quarantined",
new_revision,
&validity,
&blocking,
)
.map_err(|e| RevalidationError::Store(format!("{e:?}")))?;
Ok(RevalidationVerdict::SoundnessIncidentQuarantined {
incident_seq: seq,
new_revision,
})
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::metadata_store::{
ActionEntryRow, AuthorityRow, CommitOutcome, FsqliteEngine, PublicationRow,
QuarantineScope, ResultKindTag, RusqliteEngine, 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-h040-{}-{}-{}.db", std::process::id(), tag, n))
}
fn digest(domain: &'static str, tag: u8) -> TypedDigest {
TypedDigest {
algorithm: DigestAlgorithm::Sha256V1,
domain,
bytes: [tag; 32],
}
}
fn authority_row(tag: u8) -> AuthorityRow {
AuthorityRow {
digest: digest("rabs.authority.sha256.v1", tag),
cluster_id: "cluster-a".to_owned(),
incarnation: u128::from(tag),
term: u64::from(tag),
acquired_seq: 1,
}
}
fn validity(
evaluated_at: i64,
max_age: Option<u64>,
uncertainty: u64,
epoch: u64,
) -> ServingValidity {
ServingValidity {
evaluated_at_unix_micros: evaluated_at,
maximum_age_micros: max_age,
clock_uncertainty_micros: uncertainty,
coordinator_clock_epoch: epoch,
}
}
fn published_fixture(store: &mut dyn RabsMetadataStore) -> (TypedDigest, String) {
store.acquire_authority(&authority_row(1)).unwrap();
let active = digest("rabs.authority.sha256.v1", 1);
let action = ActionEntryRow {
action_key: digest("rabs.action-key.sha256.v1", 7),
key_epoch: 0,
projection_epoch: 0,
};
store.upsert_action_entry(&action).unwrap();
store
.create_generation(&active, 10, &action.action_key)
.unwrap();
store.record_attempt(20, 10, "worker-a", 5).unwrap();
let row = PublicationRow {
action_key: action.action_key.clone(),
descriptor_digest: digest("rabs.descriptor.sha256.v1", 1),
manifest_digest: digest("rabs.result-manifest.sha256.v1", 1),
evidence_digest: digest("rabs.evidence-bundle.sha256.v1", 1),
winner_generation: 10,
winner_attempt: 20,
result_kind: ResultKindTag::Success,
pin_id: 40,
pin_owner: "coordinator".to_owned(),
provisional_ancestors: Vec::new(),
};
assert_eq!(
store.commit_publication(&active, None, &row).unwrap(),
CommitOutcome::Committed
);
(active, digest_key(&action.action_key))
}
fn t048_scenarios(store: &mut dyn RabsMetadataStore) -> Vec<String> {
let (active, action_key) = published_fixture(store);
let legacy = store.serving_record(&action_key).unwrap().unwrap();
assert_eq!(legacy.state_revision, 0);
assert_eq!(legacy.disposition, "servable");
assert_eq!(
serving_gate(store, &action_key, 1_000, 0).unwrap(),
ServeDecision::Servable
);
assert_eq!(
serving_gate(store, "missing:key", 1_000, 0).unwrap(),
ServeDecision::NoRecord
);
store
.add_quarantine(
QuarantineScope::ActionEntry,
"other:key",
"unrelated incident",
)
.unwrap();
assert_eq!(
serving_gate(store, &action_key, 1_000, 0).unwrap(),
ServeDecision::Servable
);
store
.put_serving_record(
&active,
&action_key,
"servable",
1,
&validity(1_000, Some(500), 100, 1),
&[],
)
.unwrap();
let record = store.serving_record(&action_key).unwrap().unwrap();
assert_eq!(record.state_revision, 1);
assert_eq!(record.authority_key, digest_key(&active));
assert_eq!(
serving_gate(store, &action_key, 1_300, 1).unwrap(),
ServeDecision::Servable
);
assert_eq!(
serving_gate(store, &action_key, 1_450, 1).unwrap(),
ServeDecision::ExpiredTtl
);
assert_eq!(
serving_gate(store, &action_key, 900, 1).unwrap(),
ServeDecision::ExpiredClockRollback
);
assert_eq!(
serving_gate(store, &action_key, 1_300, 2).unwrap(),
ServeDecision::ExpiredClockEpoch
);
assert_eq!(
store.put_serving_record(
&active,
&action_key,
"servable",
1,
&validity(2_000, None, 0, 1),
&[],
),
Err(StoreError::StaleServingRevision)
);
assert_eq!(
store.put_serving_record(
&active,
&action_key,
"servable",
0,
&validity(2_000, None, 0, 1),
&[],
),
Err(StoreError::StaleServingRevision)
);
assert_eq!(
store.serving_record(&action_key).unwrap().unwrap(),
record,
"refused replays must not move the stored record"
);
let wrong = digest("rabs.authority.sha256.v1", 2);
assert_eq!(
store.put_serving_record(
&wrong,
&action_key,
"servable",
2,
&validity(2_000, None, 0, 1),
&[],
),
Err(StoreError::NotActiveAuthority)
);
let reference = (QuarantineScope::ActionEntry, action_key.clone());
assert_eq!(
store.put_serving_record(
&active,
&action_key,
"servable",
2,
&validity(2_000, None, 0, 1),
std::slice::from_ref(&reference),
),
Err(StoreError::UnknownQuarantineReference)
);
store
.add_quarantine(
QuarantineScope::ActionEntry,
&action_key,
"divergent recompute",
)
.unwrap();
assert_eq!(
serving_gate(store, &action_key, 1_300, 1).unwrap(),
ServeDecision::Blocked {
references: vec![("action-entry".to_owned(), action_key.clone())],
}
);
store
.put_serving_record(
&active,
&action_key,
"servable",
2,
&validity(2_000, None, 0, 1),
std::slice::from_ref(&reference),
)
.unwrap();
assert_eq!(
serving_gate(store, &action_key, 2_100, 1).unwrap(),
ServeDecision::Blocked {
references: vec![("action-entry".to_owned(), action_key.clone())],
}
);
let before = store.differential_snapshot().unwrap();
assert_eq!(
store.put_serving_record(
&active,
&action_key,
"servable",
3,
&validity(3_000, None, 0, 1),
&[],
),
Err(StoreError::QuarantineRequiresRepair)
);
assert_eq!(store.differential_snapshot().unwrap(), before);
assert_eq!(
serving_gate(store, &action_key, 3_100, 1).unwrap(),
ServeDecision::Blocked {
references: vec![("action-entry".to_owned(), action_key.clone())],
}
);
store
.put_serving_record(
&active,
&action_key,
"servable",
3,
&validity(3_000, None, 0, 1),
std::slice::from_ref(&reference),
)
.unwrap();
assert_eq!(
store.serving_record(&action_key).unwrap().unwrap().blocking,
vec![("action-entry".to_owned(), action_key.clone())]
);
store
.put_serving_record(
&active,
&action_key,
"evidence-pending",
4,
&validity(3_000, None, 0, 1),
std::slice::from_ref(&reference),
)
.unwrap();
assert_eq!(
serving_gate(store, &action_key, 3_100, 1).unwrap(),
ServeDecision::NotServable {
disposition: "evidence-pending".to_owned(),
}
);
let before = store.serving_record(&action_key).unwrap().unwrap();
store
.set_serving_disposition_key(&action_key, "servable")
.unwrap();
let after = store.serving_record(&action_key).unwrap().unwrap();
assert_eq!(
after.state_revision, 5,
"disposition-only write failed to fence an older serving evaluation"
);
assert_eq!(after.validity, before.validity);
assert_eq!(after.authority_key, before.authority_key);
assert_eq!(after.blocking, before.blocking);
assert_eq!(after.disposition, "servable");
assert!(matches!(
serving_gate(store, &action_key, 3_100, 1).unwrap(),
ServeDecision::Blocked { .. }
));
store.differential_snapshot().unwrap()
}
#[test]
fn t048_reference_backend() {
let engine = RusqliteEngine::open_in_memory().unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
t048_scenarios(&mut store);
}
#[test]
fn disposition_revision_exhaustion_never_partially_updates_quarantine() {
fn scenario(store: &mut dyn RabsMetadataStore) -> Vec<String> {
let (active, action_key) = published_fixture(store);
store
.put_serving_record(
&active,
&action_key,
DISPOSITION_PRESENTATION_QUARANTINED,
u64::try_from(i64::MAX).unwrap(),
&validity(1_000, Some(5_000), 10, 2),
&[],
)
.unwrap();
let before = store.differential_snapshot().unwrap();
assert_eq!(
store.set_serving_disposition_key(&action_key, DISPOSITION_QUARANTINED),
Err(StoreError::Corruption(
"state_revision out of range".to_owned()
))
);
assert_eq!(store.differential_snapshot().unwrap(), before);
store
.set_serving_disposition_key(&action_key, DISPOSITION_PRESENTATION_QUARANTINED)
.unwrap();
assert_eq!(store.differential_snapshot().unwrap(), before);
assert_eq!(
serving_gate(store, &action_key, 1_100, 2).unwrap(),
ServeDecision::NotServable {
disposition: DISPOSITION_PRESENTATION_QUARANTINED.to_owned(),
}
);
before
}
let mut reference =
SqlMetadataStore::open(RusqliteEngine::open_in_memory().unwrap()).unwrap();
let mut candidate =
SqlMetadataStore::open(FsqliteEngine::open(&fresh_path("revision-exhausted")).unwrap())
.unwrap();
assert_eq!(scenario(&mut reference), scenario(&mut candidate));
}
#[test]
fn t048_differential_reference_vs_frankensqlite() {
let reference_engine = RusqliteEngine::open(&fresh_path("ref")).unwrap();
let mut reference = SqlMetadataStore::open(reference_engine).unwrap();
let candidate_engine = FsqliteEngine::open(&fresh_path("fsq")).unwrap();
let mut candidate = SqlMetadataStore::open(candidate_engine).unwrap();
assert_eq!(
t048_scenarios(&mut reference),
t048_scenarios(&mut candidate)
);
}
#[test]
fn h040_record_survives_reopen() {
let path = fresh_path("reopen");
let action_key;
{
let engine = RusqliteEngine::open(&path).unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
let (active, key) = published_fixture(&mut store);
action_key = key;
store
.put_serving_record(
&active,
&action_key,
"servable",
5,
&validity(1_000, Some(500), 100, 3),
&[],
)
.unwrap();
}
let engine = RusqliteEngine::open(&path).unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
assert_eq!(
serving_gate(&mut store, &action_key, 1_200, 3).unwrap(),
ServeDecision::Servable
);
assert_eq!(
serving_gate(&mut store, &action_key, 1_450, 3).unwrap(),
ServeDecision::ExpiredTtl
);
assert_eq!(
serving_gate(&mut store, &action_key, 1_200, 4).unwrap(),
ServeDecision::ExpiredClockEpoch
);
assert_eq!(
serving_gate(&mut store, &action_key, 900, 3).unwrap(),
ServeDecision::ExpiredClockRollback
);
store.acquire_authority(&authority_row(1)).unwrap();
let active = digest("rabs.authority.sha256.v1", 1);
assert_eq!(
store.put_serving_record(
&active,
&action_key,
"servable",
5,
&validity(2_000, None, 0, 3),
&[],
),
Err(StoreError::StaleServingRevision)
);
}
#[test]
fn unreferenced_quarantine_survives_reopen() {
let path = fresh_path("quarantine-reopen");
let action_key;
{
let engine = RusqliteEngine::open(&path).unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
let (_, key) = published_fixture(&mut store);
action_key = key;
store
.add_quarantine(
QuarantineScope::ActionEntry,
&action_key,
"interrupted demotion",
)
.unwrap();
}
let engine = RusqliteEngine::open(&path).unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
let record = store.serving_record(&action_key).unwrap().unwrap();
assert_eq!(record.disposition, "servable");
assert!(record.blocking.is_empty());
assert_eq!(
serving_gate(&mut store, &action_key, 1_000, 0).unwrap(),
ServeDecision::Blocked {
references: vec![("action-entry".to_owned(), action_key)],
}
);
}
#[test]
fn k007_r28_class_never_publishes() {
let r28 = [
(TerminalOutcome::Oom, "oom"),
(TerminalOutcome::Signal(9), "signal"),
(TerminalOutcome::Signal(137), "signal"),
(TerminalOutcome::Cancelled, "cancelled"),
(TerminalOutcome::Timeout, "timeout"),
(TerminalOutcome::WorkerLost, "worker-loss"),
(TerminalOutcome::TransportFailed, "transport-failure"),
(TerminalOutcome::Panicked, "panic"),
];
for (outcome, class) in r28 {
assert_eq!(
classify_deterministic_failure(outcome, true, true, true, true, true),
Err(FailureRefusal::NotDeterministic(class)),
"{class} must never publish as a deterministic failure"
);
}
assert_eq!(
classify_deterministic_failure(TerminalOutcome::Exit(0), true, true, true, true, true),
Err(FailureRefusal::ZeroExitIsNotAFailure),
);
assert_eq!(
classify_deterministic_failure(
TerminalOutcome::Exit(101),
true,
true,
true,
true,
true
),
Ok(FailureAdmission {
normalized_exit: 101
}),
);
assert_eq!(
classify_deterministic_failure(
TerminalOutcome::Exit(1),
false,
false,
false,
false,
false
),
Err(FailureRefusal::CaptureIncomplete),
);
assert_eq!(
classify_deterministic_failure(
TerminalOutcome::Exit(1),
true,
false,
false,
false,
false
),
Err(FailureRefusal::InputsNotClosed),
);
assert_eq!(
classify_deterministic_failure(
TerminalOutcome::Exit(1),
true,
true,
false,
false,
false
),
Err(FailureRefusal::UndeclaredSideEffects),
);
assert_eq!(
classify_deterministic_failure(
TerminalOutcome::Exit(1),
true,
true,
true,
false,
false
),
Err(FailureRefusal::ClassPolicyRefused),
);
assert_eq!(
classify_deterministic_failure(TerminalOutcome::Exit(1), true, true, true, true, false),
Err(FailureRefusal::TrustPolicyRefused),
);
}
fn k007_lifecycle_scenarios(store: &mut dyn RabsMetadataStore) -> Vec<String> {
let (active, action_key) = published_fixture(store);
let action_typed = digest("rabs.action-key.sha256.v1", 7);
assert_eq!(
apply_revalidation(
store,
&active,
"missing:key",
&digest("rabs.action-key.sha256.v1", 99),
1,
21,
10,
"exit=1|diag=d1",
"exit=1|diag=d1",
"manifest-a",
"manifest-a",
"ev-a",
1_000,
1,
None,
),
Err(RevalidationError::NoServingRecord)
);
store
.put_serving_record(
&active,
&action_key,
"servable",
1,
&validity(1_000_000, Some(DEFAULT_REVALIDATION_TTL_MICROS), 0, 1),
&[],
)
.unwrap();
assert_eq!(
serving_gate(
store,
&action_key,
1_000_000 + i64::try_from(DEFAULT_REVALIDATION_TTL_MICROS / 2).unwrap(),
1,
)
.unwrap(),
ServeDecision::Servable
);
assert_eq!(
serving_gate(
store,
&action_key,
1_000_000 + i64::try_from(DEFAULT_REVALIDATION_TTL_MICROS + 1).unwrap(),
1,
)
.unwrap(),
ServeDecision::ExpiredTtl
);
assert_eq!(
apply_revalidation(
store,
&active,
&action_key,
&action_typed,
7,
21,
10,
"exit=1|diag=d1",
"exit=1|diag=d1",
"manifest-a",
"manifest-a",
"ev-a",
1_040_000,
1,
None,
),
Err(RevalidationError::StaleRevision { stored: 1 })
);
assert_eq!(
apply_revalidation(
store,
&active,
&action_key,
&action_typed,
1,
21,
10,
"exit=1|diag=d1",
"exit=1|diag=d1",
"manifest-a",
"manifest-a",
"ev-a",
1_040_000,
1,
None,
),
Ok(RevalidationVerdict::IdenticalEvidenceAppended { new_revision: 2 })
);
assert_eq!(
store
.serving_record(&action_key)
.unwrap()
.unwrap()
.state_revision,
2
);
assert_eq!(
serving_gate(
store,
&action_key,
1_040_000 + i64::try_from(DEFAULT_REVALIDATION_TTL_MICROS / 2).unwrap(),
1,
)
.unwrap(),
ServeDecision::Servable
);
assert_eq!(
serving_gate(
store,
&action_key,
1_040_000 + i64::try_from(DEFAULT_REVALIDATION_TTL_MICROS + 1).unwrap(),
1,
)
.unwrap(),
ServeDecision::ExpiredTtl
);
assert_eq!(
apply_revalidation(
store,
&active,
&action_key,
&action_typed,
2,
22,
11,
"exit=1|diag=d1",
"success",
"manifest-a",
"manifest-b",
"ev-b",
1_060_000,
1,
None,
),
Ok(RevalidationVerdict::SoundnessIncidentQuarantined {
incident_seq: 3,
new_revision: 3,
})
);
let incidents = store.list_divergence_incidents(&action_key).unwrap();
assert_eq!(incidents.len(), 1);
assert_eq!(incidents[0].class, "soundness");
assert_eq!(incidents[0].seq, 3);
let record = store.serving_record(&action_key).unwrap().unwrap();
assert_eq!(record.disposition, "quarantined");
assert_eq!(
record.blocking,
vec![("action-entry".to_owned(), action_key.clone())]
);
assert_eq!(
serving_gate(store, &action_key, 1_060_000, 1).unwrap(),
ServeDecision::NotServable {
disposition: "quarantined".to_owned(),
}
);
let before = store.differential_snapshot().unwrap();
assert_eq!(
apply_revalidation(
store,
&active,
&action_key,
&action_typed,
3,
23,
11,
"exit=1|diag=d1",
"exit=1|diag=d1",
"manifest-a",
"manifest-a",
"ev-c",
1_080_000,
1,
None,
),
Err(RevalidationError::ServingBlocked)
);
assert_eq!(store.differential_snapshot().unwrap(), before);
store
.put_serving_record(
&active,
&action_key,
"servable",
4,
&validity(1_080_000, Some(DEFAULT_REVALIDATION_TTL_MICROS), 0, 1),
&[],
)
.unwrap();
let before = store.differential_snapshot().unwrap();
assert_eq!(
apply_revalidation(
store,
&active,
&action_key,
&action_typed,
4,
24,
11,
"exit=1|diag=d1",
"exit=1|diag=d1",
"manifest-a",
"manifest-a",
"ev-d",
1_100_000,
1,
None,
),
Err(RevalidationError::ServingBlocked)
);
assert_eq!(store.differential_snapshot().unwrap(), before);
store.differential_snapshot().unwrap()
}
#[test]
fn k007_failure_ttl_lifecycle_reference_backend() {
let engine = RusqliteEngine::open_in_memory().unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
k007_lifecycle_scenarios(&mut store);
}
#[test]
fn k007_failure_ttl_lifecycle_differential_reference_vs_frankensqlite() {
let reference_engine = RusqliteEngine::open(&fresh_path("k007ref")).unwrap();
let mut reference = SqlMetadataStore::open(reference_engine).unwrap();
let candidate_engine = FsqliteEngine::open(&fresh_path("k007fsq")).unwrap();
let mut candidate = SqlMetadataStore::open(candidate_engine).unwrap();
assert_eq!(
k007_lifecycle_scenarios(&mut reference),
k007_lifecycle_scenarios(&mut candidate)
);
}
fn k007_fencing_and_incident_history(store: &mut dyn RabsMetadataStore) -> Vec<String> {
let (active, action_key) = published_fixture(store);
let action_typed = digest("rabs.action-key.sha256.v1", 7);
let wrong_action = digest("rabs.action-key.sha256.v1", 8);
let wrong_authority = digest("rabs.authority.sha256.v1", 2);
store
.put_serving_record(
&active,
&action_key,
SERVABLE_DISPOSITION,
1,
&validity(1_000, Some(500), 0, 1),
&[],
)
.unwrap();
let before = store.differential_snapshot().unwrap();
for signature in ["published", "diverged"] {
assert_eq!(
apply_revalidation(
store,
&wrong_authority,
&action_key,
&action_typed,
1,
20,
10,
"published",
signature,
"manifest-a",
"manifest-b",
"ev-b",
1_100,
1,
None,
),
Err(RevalidationError::Store(format!(
"{:?}",
StoreError::NotActiveAuthority
)))
);
assert_eq!(store.differential_snapshot().unwrap(), before);
assert_eq!(
apply_revalidation(
store,
&active,
&action_key,
&wrong_action,
1,
20,
10,
"published",
signature,
"manifest-a",
"manifest-b",
"ev-b",
1_100,
1,
None,
),
Err(RevalidationError::ActionKeyMismatch)
);
assert_eq!(store.differential_snapshot().unwrap(), before);
}
store
.add_quarantine(
QuarantineScope::LogicalObject,
"object:damaged",
"bad bytes",
)
.unwrap();
store
.add_quarantine(QuarantineScope::ActionEntry, &action_key, "operator hold")
.unwrap();
store
.put_serving_record(
&active,
&action_key,
SERVABLE_DISPOSITION,
2,
&validity(1_000, Some(500), 0, 1),
&[(QuarantineScope::LogicalObject, "object:damaged".to_owned())],
)
.unwrap();
store
.record_divergence_incident(
&active,
&DivergenceIncidentRow {
action_key: action_key.clone(),
seq: 40,
class: "prior-incident".to_owned(),
committed_manifest_key: "manifest-a".to_owned(),
candidate_manifest_key: "manifest-old".to_owned(),
candidate_evidence_key: "ev-old".to_owned(),
candidate_pin_hex: String::new(),
generation_hex: "a".to_owned(),
attempt_hex: "14".to_owned(),
detail: "prior sparse incident".to_owned(),
},
)
.unwrap();
let frozen_publication: Vec<_> = store
.differential_snapshot()
.unwrap()
.into_iter()
.filter(|line| line.starts_with("action_publications|"))
.collect();
assert_eq!(
apply_revalidation(
store,
&active,
&action_key,
&action_typed,
2,
20,
10,
"published",
"diverged",
"manifest-a",
"manifest-b",
"ev-b",
1_200,
1,
None,
),
Ok(RevalidationVerdict::SoundnessIncidentQuarantined {
incident_seq: 41,
new_revision: 3,
})
);
let incidents = store.list_divergence_incidents(&action_key).unwrap();
assert_eq!(incidents.len(), 2);
assert_eq!(incidents[0].seq, 40);
assert_eq!(incidents[1].seq, 41);
let record = store.serving_record(&action_key).unwrap().unwrap();
assert_eq!(
record.blocking,
vec![
("action-entry".to_owned(), action_key.clone()),
("logical-object".to_owned(), "object:damaged".to_owned()),
]
);
assert_eq!(
store
.query(
"SELECT reason FROM quarantines WHERE scope = 'action-entry' AND subject = ?1",
&[SqlValue::Text(action_key.clone())],
)
.unwrap(),
vec![vec![SqlValue::Text("operator hold".to_owned())]]
);
assert!(matches!(
serving_gate(store, &action_key, 1_200, 1).unwrap(),
ServeDecision::NotServable { .. }
));
let snapshot = store.differential_snapshot().unwrap();
assert_eq!(
snapshot
.iter()
.filter(|line| line.starts_with("action_publications|"))
.cloned()
.collect::<Vec<_>>(),
frozen_publication
);
snapshot
}
#[test]
fn k007_revalidation_fencing_reference() {
let engine = RusqliteEngine::open_in_memory().unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
k007_fencing_and_incident_history(&mut store);
}
#[test]
fn k007_revalidation_fencing_differential() {
let reference_engine = RusqliteEngine::open(&fresh_path("k007fence-ref")).unwrap();
let candidate_engine = FsqliteEngine::open(&fresh_path("k007fence-fsq")).unwrap();
let mut reference = SqlMetadataStore::open(reference_engine).unwrap();
let mut candidate = SqlMetadataStore::open(candidate_engine).unwrap();
assert_eq!(
k007_fencing_and_incident_history(&mut reference),
k007_fencing_and_incident_history(&mut candidate)
);
}
}