use sha2::{Digest, Sha256};
use rabs_protocol::result_identity::TypedDigest;
use crate::metadata_store::{
DivergenceIncidentRow, QuarantineScope, RabsMetadataStore, StoreError, digest_key,
};
use crate::serving_state::{
RevalidationError, SERVABLE_DISPOSITION, action_quarantine_present,
divergence_quarantine_disposition,
};
use crate::trust_evidence::{
DISPOSITION_QUARANTINED, require_active_authority, verification_evidence,
};
pub const SERVING_SAMPLE_QUARANTINE_REASON: &str = "k008-served-divergence";
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ActionClassRisk {
LowRiskRegistry,
Elevated,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct SamplingPolicy {
pub min_samples: u32,
pub min_pass_rate_basis_points: u32,
pub sample_rate_basis_points: u32,
}
impl SamplingPolicy {
#[must_use]
pub const fn sample_all(min_samples: u32, min_pass_rate_basis_points: u32) -> Self {
Self {
min_samples,
min_pass_rate_basis_points,
sample_rate_basis_points: 10_000,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum PrivateExecutionReason {
ElevatedClassRisk,
InvalidPolicy,
NoPublishedResult,
Quarantined,
ServingNotEligible {
disposition: String,
},
InsufficientVerificationSamples {
observed: u32,
required: u32,
},
VerificationRateBelowPolicy {
observed_basis_points: u32,
required_basis_points: u32,
},
AdverseVerificationSamples {
observed: u64,
},
NotSampledThisEpoch {
key_bucket_basis_points: u32,
},
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum SampleGateDecision {
ServeFromCache,
ExecutePrivately(PrivateExecutionReason),
}
#[must_use]
pub fn key_bucket_basis_points(action: &TypedDigest) -> u32 {
let digest = Sha256::digest(digest_key(action).as_bytes());
let window = u32::from(digest[0]) << 24
| u32::from(digest[1]) << 16
| u32::from(digest[2]) << 8
| u32::from(digest[3]);
((window >> 16).saturating_mul(10_000)) >> 16
}
pub fn serving_sample_decision(
store: &mut dyn RabsMetadataStore,
action: &TypedDigest,
risk: ActionClassRisk,
policy: &SamplingPolicy,
) -> Result<SampleGateDecision, StoreError> {
if risk == ActionClassRisk::Elevated {
return Ok(SampleGateDecision::ExecutePrivately(
PrivateExecutionReason::ElevatedClassRisk,
));
}
if policy.min_pass_rate_basis_points > 10_000 || policy.sample_rate_basis_points > 10_000 {
return Ok(SampleGateDecision::ExecutePrivately(
PrivateExecutionReason::InvalidPolicy,
));
}
if !store.has_publication(action)? {
return Ok(SampleGateDecision::ExecutePrivately(
PrivateExecutionReason::NoPublishedResult,
));
}
let action_key = digest_key(action);
let Some(record) = store.serving_record(&action_key)? else {
return Ok(SampleGateDecision::ExecutePrivately(
PrivateExecutionReason::NoPublishedResult,
));
};
if !record.blocking.is_empty()
|| action_quarantine_present(store, &action_key)?
|| divergence_quarantine_disposition(store, &action_key)?.is_some()
{
return Ok(SampleGateDecision::ExecutePrivately(
PrivateExecutionReason::Quarantined,
));
}
if record.disposition != SERVABLE_DISPOSITION {
return Ok(SampleGateDecision::ExecutePrivately(
PrivateExecutionReason::ServingNotEligible {
disposition: record.disposition,
},
));
}
let evidence = verification_evidence(store, action)?;
let observed = u32::try_from(evidence.attempts).unwrap_or(u32::MAX);
let required = policy.min_samples.max(1);
if observed < required {
return Ok(SampleGateDecision::ExecutePrivately(
PrivateExecutionReason::InsufficientVerificationSamples { observed, required },
));
}
let pass_rate_basis_points = u32::try_from(
u128::from(evidence.passed_attempts) * 10_000 / u128::from(evidence.attempts),
)
.unwrap_or(0);
if pass_rate_basis_points < policy.min_pass_rate_basis_points {
return Ok(SampleGateDecision::ExecutePrivately(
PrivateExecutionReason::VerificationRateBelowPolicy {
observed_basis_points: pass_rate_basis_points,
required_basis_points: policy.min_pass_rate_basis_points,
},
));
}
if evidence.adverse_samples > 0 {
return Ok(SampleGateDecision::ExecutePrivately(
PrivateExecutionReason::AdverseVerificationSamples {
observed: evidence.adverse_samples,
},
));
}
let bucket = key_bucket_basis_points(action);
if bucket >= policy.sample_rate_basis_points {
return Ok(SampleGateDecision::ExecutePrivately(
PrivateExecutionReason::NotSampledThisEpoch {
key_bucket_basis_points: bucket,
},
));
}
Ok(SampleGateDecision::ServeFromCache)
}
pub fn quarantine_served_divergence(
store: &mut dyn RabsMetadataStore,
authority: &TypedDigest,
action_key_str: &str,
expected_revision: u64,
generation: u128,
attempt: u128,
detail: &str,
) -> Result<u64, 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,
});
}
require_active_authority(store, authority)
.map_err(|e| RevalidationError::Store(format!("{e:?}")))?;
let new_revision = expected_revision
.checked_add(1)
.ok_or(RevalidationError::RevisionExhausted)?;
if !action_quarantine_present(store, action_key_str)
.map_err(|e| RevalidationError::Store(format!("{e:?}")))?
{
store
.add_quarantine(
QuarantineScope::ActionEntry,
action_key_str,
SERVING_SAMPLE_QUARANTINE_REASON,
)
.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(0), |seq| seq.checked_add(1))
.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: "serving-sample-divergence".to_owned(),
committed_manifest_key: String::new(),
candidate_manifest_key: String::new(),
candidate_evidence_key: String::new(),
candidate_pin_hex: String::new(),
generation_hex: format!("{generation:x}"),
attempt_hex: format!("{attempt:x}"),
detail: detail.to_owned(),
},
)
.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,
DISPOSITION_QUARANTINED,
new_revision,
&record.validity,
&blocking,
)
.map_err(|e| RevalidationError::Store(format!("{e:?}")))?;
Ok(new_revision)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::metadata_store::{
ActionEntryRow, AuthorityRow, CommitOutcome, FsqliteEngine, PublicationRow, ResultKindTag,
RusqliteEngine, SqlMetadataStore,
};
use crate::serving_state::{ServeDecision, serving_gate};
use rabs_protocol::result_identity::DigestAlgorithm;
use rabs_protocol::serving::ServingValidity;
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-k008-{tag}-{}-{n}.db", std::process::id()))
}
fn digest(domain: &'static str, tag: u8) -> TypedDigest {
TypedDigest {
algorithm: DigestAlgorithm::Sha256V1,
domain,
bytes: [tag; 32],
}
}
fn action(tag: u8) -> TypedDigest {
digest("rabs.action-key.sha256.v1", tag)
}
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 published(store: &mut dyn RabsMetadataStore, tag: u8, generation: u128) -> String {
let active = digest("rabs.authority.sha256.v1", 1);
store.acquire_authority(&authority_row(1)).unwrap();
let entry = ActionEntryRow {
action_key: action(tag),
key_epoch: 0,
projection_epoch: 0,
};
store.upsert_action_entry(&entry).unwrap();
store
.create_generation(&active, generation, &entry.action_key)
.unwrap();
let attempt = generation * 10 + 1;
store
.record_attempt(attempt, generation, "worker-a", 5)
.unwrap();
let row = PublicationRow {
action_key: entry.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: generation,
winner_attempt: attempt,
result_kind: ResultKindTag::Success,
pin_id: generation * 10 + 2,
pin_owner: "coordinator".to_owned(),
provisional_ancestors: Vec::new(),
};
assert_eq!(
store.commit_publication(&active, None, &row).unwrap(),
CommitOutcome::Committed
);
digest_key(&entry.action_key)
}
fn samples(
store: &mut dyn RabsMetadataStore,
tag: u8,
generation: u128,
passes: u32,
fails: u32,
) {
let first = store.list_verification_samples(&action(tag)).unwrap().len() as u64 + 1;
for (seq, passed) in (first..).zip(
std::iter::repeat_n(true, passes as usize)
.chain(std::iter::repeat_n(false, fails as usize)),
) {
let attempt = generation * 1_000 + u128::from(seq);
store
.record_attempt(attempt, generation, "worker-sampler", seq)
.unwrap();
store
.record_verification_sample(&action(tag), attempt, passed, seq)
.unwrap();
}
}
fn k008_scenarios(store: &mut dyn RabsMetadataStore) -> Vec<String> {
let active = digest("rabs.authority.sha256.v1", 1);
let strict = SamplingPolicy {
min_samples: 4,
min_pass_rate_basis_points: 9_900,
sample_rate_basis_points: 10_000,
};
assert_eq!(
serving_sample_decision(
store,
&action(1),
ActionClassRisk::Elevated,
&SamplingPolicy::sample_all(0, 0),
)
.unwrap(),
SampleGateDecision::ExecutePrivately(PrivateExecutionReason::ElevatedClassRisk)
);
assert_eq!(
serving_sample_decision(
store,
&action(99),
ActionClassRisk::LowRiskRegistry,
&strict
)
.unwrap(),
SampleGateDecision::ExecutePrivately(PrivateExecutionReason::NoPublishedResult)
);
let key_one = published(store, 1, 10);
assert_eq!(
serving_sample_decision(
store,
&action(1),
ActionClassRisk::LowRiskRegistry,
&SamplingPolicy::sample_all(0, 0),
)
.unwrap(),
SampleGateDecision::ExecutePrivately(
PrivateExecutionReason::InsufficientVerificationSamples {
observed: 0,
required: 1,
}
)
);
for invalid in [
SamplingPolicy {
min_pass_rate_basis_points: 10_001,
..strict
},
SamplingPolicy {
sample_rate_basis_points: 10_001,
..strict
},
] {
assert_eq!(
serving_sample_decision(
store,
&action(1),
ActionClassRisk::LowRiskRegistry,
&invalid
)
.unwrap(),
SampleGateDecision::ExecutePrivately(PrivateExecutionReason::InvalidPolicy)
);
}
assert_eq!(
serving_sample_decision(store, &action(1), ActionClassRisk::LowRiskRegistry, &strict)
.unwrap(),
SampleGateDecision::ExecutePrivately(
PrivateExecutionReason::InsufficientVerificationSamples {
observed: 0,
required: 4,
}
)
);
samples(store, 1, 10, 3, 1);
assert_eq!(
serving_sample_decision(store, &action(1), ActionClassRisk::LowRiskRegistry, &strict)
.unwrap(),
SampleGateDecision::ExecutePrivately(
PrivateExecutionReason::VerificationRateBelowPolicy {
observed_basis_points: 7_500,
required_basis_points: 9_900,
}
)
);
samples(store, 1, 10, 5, 0);
assert_eq!(
serving_sample_decision(store, &action(1), ActionClassRisk::LowRiskRegistry, &strict)
.unwrap(),
SampleGateDecision::ExecutePrivately(
PrivateExecutionReason::VerificationRateBelowPolicy {
observed_basis_points: 8_888,
required_basis_points: 9_900,
}
)
);
let key_two = published(store, 2, 11);
samples(store, 2, 11, 4, 0);
let none = SamplingPolicy {
sample_rate_basis_points: 0,
..strict
};
match serving_sample_decision(store, &action(2), ActionClassRisk::LowRiskRegistry, &none)
.unwrap()
{
SampleGateDecision::ExecutePrivately(PrivateExecutionReason::NotSampledThisEpoch {
key_bucket_basis_points: observed_bucket,
}) => {
assert_eq!(observed_bucket, key_bucket_basis_points(&action(2)));
assert!(observed_bucket < 10_000);
}
other => panic!("zero-rate policy served anyway: {other:?}"),
}
assert_eq!(
serving_sample_decision(
store,
&action(2),
ActionClassRisk::LowRiskRegistry,
&SamplingPolicy::sample_all(4, 9_900),
)
.unwrap(),
SampleGateDecision::ServeFromCache
);
assert_eq!(
key_bucket_basis_points(&action(2)),
key_bucket_basis_points(&action(2))
);
store
.record_verification_sample(&action(2), 11_001, true, 500)
.unwrap();
store
.record_verification_sample(&action(2), 99_999, true, 501)
.unwrap();
assert_eq!(
serving_sample_decision(
store,
&action(2),
ActionClassRisk::LowRiskRegistry,
&SamplingPolicy::sample_all(5, 9_900),
)
.unwrap(),
SampleGateDecision::ExecutePrivately(
PrivateExecutionReason::InsufficientVerificationSamples {
observed: 4,
required: 5,
}
)
);
store
.set_serving_disposition_key(&key_two, "evidence-pending")
.unwrap();
assert_eq!(
serving_sample_decision(store, &action(2), ActionClassRisk::LowRiskRegistry, &strict)
.unwrap(),
SampleGateDecision::ExecutePrivately(PrivateExecutionReason::ServingNotEligible {
disposition: "evidence-pending".to_owned(),
})
);
store
.set_serving_disposition_key(&key_two, SERVABLE_DISPOSITION)
.unwrap();
store
.add_quarantine(QuarantineScope::ActionEntry, &key_two, "corrupt closure")
.unwrap();
assert!(
store
.serving_record(&key_two)
.unwrap()
.unwrap()
.blocking
.is_empty()
);
assert_eq!(
serving_sample_decision(store, &action(2), ActionClassRisk::LowRiskRegistry, &strict)
.unwrap(),
SampleGateDecision::ExecutePrivately(PrivateExecutionReason::Quarantined)
);
store
.add_quarantine(
QuarantineScope::LogicalObject,
"object:damaged",
"bad bytes",
)
.unwrap();
store
.put_serving_record(
&active,
&key_one,
SERVABLE_DISPOSITION,
7,
&ServingValidity {
evaluated_at_unix_micros: 1_000,
maximum_age_micros: None,
clock_uncertainty_micros: 0,
coordinator_clock_epoch: 1,
},
&[(QuarantineScope::LogicalObject, "object:damaged".to_owned())],
)
.unwrap();
assert_eq!(
serving_sample_decision(store, &action(1), ActionClassRisk::LowRiskRegistry, &strict)
.unwrap(),
SampleGateDecision::ExecutePrivately(PrivateExecutionReason::Quarantined)
);
let before = store.differential_snapshot().unwrap();
assert_eq!(
quarantine_served_divergence(store, &active, &key_one, 6, 11, 22, "mismatch"),
Err(RevalidationError::StaleRevision { stored: 7 })
);
let wrong = digest("rabs.authority.sha256.v1", 2);
assert_eq!(
quarantine_served_divergence(store, &wrong, &key_one, 7, 11, 22, "mismatch"),
Err(RevalidationError::Store(format!(
"{:?}",
StoreError::NotActiveAuthority
)))
);
assert_eq!(store.differential_snapshot().unwrap(), before);
store
.record_divergence_incident(
&active,
&DivergenceIncidentRow {
action_key: key_one.clone(),
seq: 40,
class: "prior-incident".to_owned(),
committed_manifest_key: String::new(),
candidate_manifest_key: String::new(),
candidate_evidence_key: String::new(),
candidate_pin_hex: String::new(),
generation_hex: "a".to_owned(),
attempt_hex: "14".to_owned(),
detail: "preexisting sparse sequence".to_owned(),
},
)
.unwrap();
let new_revision = quarantine_served_divergence(
store,
&active,
&key_one,
7,
11,
22,
"stdout digest mismatch",
)
.unwrap();
assert_eq!(new_revision, 8);
let incidents = store.list_divergence_incidents(&key_one).unwrap();
assert_eq!(incidents.len(), 2);
assert_eq!(incidents[0].seq, 40);
assert_eq!(incidents[1].seq, 41);
assert_eq!(incidents[1].class, "serving-sample-divergence");
assert_eq!(incidents[1].detail, "stdout digest mismatch");
let record = store.serving_record(&key_one).unwrap().unwrap();
assert_eq!(record.disposition, DISPOSITION_QUARANTINED);
assert_eq!(
record.blocking,
vec![
("action-entry".to_owned(), key_one.clone()),
("logical-object".to_owned(), "object:damaged".to_owned()),
]
);
assert_eq!(
serving_gate(store, &key_one, 2_000, 1).unwrap(),
ServeDecision::NotServable {
disposition: DISPOSITION_QUARANTINED.to_owned(),
}
);
assert_eq!(
quarantine_served_divergence(store, &active, &key_one, 7, 11, 23, "again"),
Err(RevalidationError::StaleRevision { stored: 8 })
);
assert_eq!(
quarantine_served_divergence(store, &active, &key_one, 8, 11, 23, "next"),
Ok(9)
);
let incidents = store.list_divergence_incidents(&key_one).unwrap();
assert_eq!(
incidents
.iter()
.map(|incident| incident.seq)
.collect::<Vec<_>>(),
vec![40, 41, 42]
);
assert_eq!(
store.serving_record(&key_one).unwrap().unwrap().blocking,
record.blocking
);
assert_eq!(
quarantine_served_divergence(store, &active, "missing:key", 1, 1, 1, "x"),
Err(RevalidationError::NoServingRecord)
);
store.differential_snapshot().unwrap()
}
#[test]
fn k008_reference_backend() {
let engine = RusqliteEngine::open_in_memory().unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
k008_scenarios(&mut store);
}
#[test]
fn k008_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!(
k008_scenarios(&mut reference),
k008_scenarios(&mut candidate)
);
}
fn adverse_verification_scenarios(store: &mut dyn RabsMetadataStore) -> Vec<String> {
let permissive = SamplingPolicy::sample_all(3, 7_500);
for (tag, generation) in [(1_u8, 10_u128), (2, 11), (3, 12)] {
let key = published(store, tag, generation);
samples(store, tag, generation, 3, 0);
assert_eq!(
serving_sample_decision(
store,
&action(tag),
ActionClassRisk::LowRiskRegistry,
&permissive
)
.unwrap(),
SampleGateDecision::ServeFromCache
);
let failed_attempt = match tag {
1 => {
let attempt = generation * 1_000 + 4;
store
.record_attempt(attempt, generation, "worker-sampler", 100)
.unwrap();
attempt
}
2 => 99_999, _ => {
published(store, 9, 13);
131 }
};
store
.record_verification_sample(&action(tag), failed_attempt, false, 100)
.unwrap();
for policy in [permissive, SamplingPolicy::sample_all(0, 0)] {
let before = store.differential_snapshot().unwrap();
assert_eq!(
serving_sample_decision(
store,
&action(tag),
ActionClassRisk::LowRiskRegistry,
&policy
)
.unwrap(),
SampleGateDecision::ExecutePrivately(
PrivateExecutionReason::AdverseVerificationSamples { observed: 1 }
)
);
assert_eq!(store.differential_snapshot().unwrap(), before);
}
store
.record_verification_sample(&action(tag), failed_attempt, true, 101)
.unwrap();
samples(store, tag, generation, 1, 0);
assert_eq!(
serving_sample_decision(
store,
&action(tag),
ActionClassRisk::LowRiskRegistry,
&permissive
)
.unwrap(),
SampleGateDecision::ExecutePrivately(
PrivateExecutionReason::AdverseVerificationSamples { observed: 1 }
)
);
assert_eq!(
store.serving_disposition_key(&key).unwrap().as_deref(),
Some(SERVABLE_DISPOSITION)
);
assert!(!action_quarantine_present(store, &key).unwrap());
}
store.differential_snapshot().unwrap()
}
#[test]
fn adverse_verification_is_a_safety_veto_reference() {
let engine = RusqliteEngine::open_in_memory().unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
adverse_verification_scenarios(&mut store);
}
#[test]
fn adverse_verification_is_a_safety_veto_differential() {
let reference_engine = RusqliteEngine::open(&fresh_path("adverse-ref")).unwrap();
let candidate_engine = FsqliteEngine::open(&fresh_path("adverse-fsq")).unwrap();
let mut reference = SqlMetadataStore::open(reference_engine).unwrap();
let mut candidate = SqlMetadataStore::open(candidate_engine).unwrap();
assert_eq!(
adverse_verification_scenarios(&mut reference),
adverse_verification_scenarios(&mut candidate)
);
}
}