use rabs_protocol::generation::AttemptAuthority;
use rabs_protocol::raw_bytes::RawBytes;
use rabs_protocol::result_identity::{
ActionPublicationRecord, AttemptEvidenceBundle, CanonicalActionResultManifest, DigestAlgorithm,
DivergenceClass, ObjectId, OutputRole, ResultKind, TypedDigest,
};
use sha2::{Digest, Sha256};
use crate::metadata_store::{
CommitOutcome, DivergenceIncidentRow, ProvisionalAncestorRow, PublicationRow, QuarantineScope,
RabsMetadataStore, ResultKindTag, StoreError, digest_key,
};
use crate::trust_evidence::DISPOSITION_QUARANTINED;
pub const AUTHORITY_DIGEST_DOMAIN: &str = rabs_key::typed_digest::DOMAIN_COORDINATOR_AUTHORITY;
pub const SEMANTIC_PROJECTION_DOMAIN: &str = "rabs.semantic-result-projection.sha256.v1";
pub const OBSERVABLE_PROJECTION_DOMAIN: &str = "rabs.observable-result-projection.sha256.v1";
pub const DISPOSITION_PRESENTATION_QUARANTINED: &str = "presentation-quarantined";
pub const DIVERGENCE_EVIDENCE_PIN_CLASS: &str = "divergence-evidence";
pub(crate) struct Framing(Sha256);
impl Framing {
pub(crate) fn new(domain: &str) -> Self {
let mut hasher = Sha256::new();
hasher.update((domain.len() as u64).to_be_bytes());
hasher.update(domain.as_bytes());
Self(hasher)
}
pub(crate) fn field(&mut self, bytes: &[u8]) -> &mut Self {
self.0.update((bytes.len() as u64).to_be_bytes());
self.0.update(bytes);
self
}
pub(crate) fn u64(&mut self, v: u64) -> &mut Self {
self.field(&v.to_be_bytes())
}
pub(crate) fn digest_field(&mut self, d: &TypedDigest) -> &mut Self {
self.field(d.domain.as_bytes());
self.field(&d.bytes)
}
pub(crate) fn finish(self, domain: &'static str) -> TypedDigest {
TypedDigest {
algorithm: DigestAlgorithm::Sha256V1,
domain,
bytes: self.0.finalize().into(),
}
}
}
#[must_use]
pub fn authority_digest(authority: &rabs_protocol::authority::CoordinatorAuthority) -> TypedDigest {
rabs_key::authority_binding::coordinator_authority_digest(authority)
}
const fn result_kind_tag(kind: ResultKind) -> u64 {
match kind {
ResultKind::Success => 0,
ResultKind::DeterministicFailure => 1,
}
}
pub(crate) const fn output_role_tag(role: OutputRole) -> u64 {
match role {
OutputRole::Materializable => 0,
OutputRole::DepInfo => 1,
OutputRole::ProvisionalMetadata => 2,
OutputRole::BuildScriptMetadata => 3,
OutputRole::TestSideEffect => 4,
}
}
#[must_use]
pub fn semantic_result_digest_v1(manifest: &CanonicalActionResultManifest) -> TypedDigest {
let mut framing = Framing::new(SEMANTIC_PROJECTION_DOMAIN);
framing
.digest_field(&manifest.action_key)
.digest_field(&manifest.canonical_descriptor_digest)
.u64(u64::from(manifest.key_epoch))
.u64(u64::from(manifest.projection_epoch))
.u64(result_kind_tag(manifest.result_kind));
match &manifest.artifact_bundle_root {
None => framing.u64(0),
Some(root) => framing.u64(1).digest_field(&root.0),
};
let mut outputs: Vec<&rabs_protocol::result_identity::LogicalOutput> =
manifest.logical_outputs.iter().collect();
outputs.sort_by(|a, b| {
(output_role_tag(a.role), a.virtual_path.as_bytes())
.cmp(&(output_role_tag(b.role), b.virtual_path.as_bytes()))
});
framing.u64(outputs.len() as u64);
for output in outputs {
framing
.u64(output_role_tag(output.role))
.field(output.virtual_path.as_bytes())
.digest_field(&output.object.0);
}
framing.finish(SEMANTIC_PROJECTION_DOMAIN)
}
#[must_use]
pub fn observable_result_digest_v1(
manifest: &CanonicalActionResultManifest,
canonical_observations: &TypedDigest,
) -> TypedDigest {
let mut framing = Framing::new(OBSERVABLE_PROJECTION_DOMAIN);
framing
.digest_field(&semantic_result_digest_v1(manifest))
.digest_field(canonical_observations);
framing.finish(OBSERVABLE_PROJECTION_DOMAIN)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum OfferBuildError {
ManifestInvalid(&'static str),
EvidenceManifestMismatch,
ActionKeyMismatch,
UndeclaredOutput {
path: String,
},
DuplicateAncestorRef {
path: String,
},
SelfAncestorRef,
BundleRoot(rabs_key::logical_output_map::BundleRootError),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ProvisionalAncestorRef {
pub producer_action_key: TypedDigest,
pub role: OutputRole,
pub virtual_path: RawBytes,
pub consumed_object: ObjectId,
}
const fn output_role_name(role: OutputRole) -> &'static str {
match role {
OutputRole::Materializable => "materializable",
OutputRole::DepInfo => "dep-info",
OutputRole::ProvisionalMetadata => "provisional-metadata",
OutputRole::BuildScriptMetadata => "build-script-metadata",
OutputRole::TestSideEffect => "test-side-effect",
}
}
pub(crate) const fn output_role_name_for_tag(tag: i64) -> &'static str {
match tag {
1 => "dep-info",
2 => "provisional-metadata",
3 => "build-script-metadata",
4 => "test-side-effect",
_ => "materializable",
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct OfferPreparedActionResult {
pub authority: AttemptAuthority,
pub manifest: CanonicalActionResultManifest,
pub manifest_id: ObjectId,
pub evidence: AttemptEvidenceBundle,
pub evidence_id: ObjectId,
pub canonical_observations: TypedDigest,
pub provisional_ancestors: Vec<ProvisionalAncestorRef>,
}
impl OfferPreparedActionResult {
#[allow(clippy::too_many_arguments)]
pub fn build(
authority: AttemptAuthority,
mut manifest: CanonicalActionResultManifest,
manifest_id: ObjectId,
evidence: AttemptEvidenceBundle,
evidence_id: ObjectId,
canonical_observations: TypedDigest,
declared_outputs: &[(OutputRole, RawBytes)],
mut provisional_ancestors: Vec<ProvisionalAncestorRef>,
) -> Result<Self, OfferBuildError> {
manifest
.validate()
.map_err(OfferBuildError::ManifestInvalid)?;
provisional_ancestors.sort_by(|a, b| {
(
digest_key(&a.producer_action_key),
output_role_name(a.role),
a.virtual_path.as_bytes(),
)
.cmp(&(
digest_key(&b.producer_action_key),
output_role_name(b.role),
b.virtual_path.as_bytes(),
))
});
for pair in provisional_ancestors.windows(2) {
if pair[0].producer_action_key == pair[1].producer_action_key
&& pair[0].role == pair[1].role
&& pair[0].virtual_path == pair[1].virtual_path
{
return Err(OfferBuildError::DuplicateAncestorRef {
path: pair[1].virtual_path.escaped(),
});
}
}
if provisional_ancestors
.iter()
.any(|a| a.producer_action_key == manifest.action_key)
{
return Err(OfferBuildError::SelfAncestorRef);
}
if evidence.canonical_result_manifest_id != manifest_id {
return Err(OfferBuildError::EvidenceManifestMismatch);
}
if evidence.action_key != manifest.action_key || authority.action_key != manifest.action_key
{
return Err(OfferBuildError::ActionKeyMismatch);
}
for output in &manifest.logical_outputs {
let declared = declared_outputs
.iter()
.any(|(role, path)| *role == output.role && *path == output.virtual_path);
if !declared {
return Err(OfferBuildError::UndeclaredOutput {
path: output.virtual_path.escaped(),
});
}
}
manifest.artifact_bundle_root =
rabs_key::logical_output_map::compute_bundle_root(&manifest.logical_outputs);
rabs_key::logical_output_map::verify_manifest_bundle_root(&manifest)
.map_err(OfferBuildError::BundleRoot)?;
manifest.semantic_result_digest = semantic_result_digest_v1(&manifest);
manifest.observable_result_digest =
observable_result_digest_v1(&manifest, &canonical_observations);
Ok(Self {
authority,
manifest,
manifest_id,
evidence,
evidence_id,
canonical_observations,
provisional_ancestors,
})
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum OfferRefusal {
NotActiveAuthority,
GenerationAuthorityMismatch,
UnknownGeneration,
GenerationTombstoned,
UnknownAttempt,
UnknownLease,
LeaseReleased,
LeaseExpired,
DescriptorMismatch,
ActionKeyMismatch,
UnknownActionEntry,
EpochMismatch,
BundleRootMismatch {
rule: String,
},
SemanticDigestMismatch,
ObservableDigestMismatch,
IncompleteObjectClosure {
missing: String,
},
ObjectNotDurable {
missing: String,
},
CommittedManifestUnavailable,
ProvisionalProducerNotCommitted {
producer: String,
},
AncestorManifestUnavailable {
producer: String,
},
AncestorOutputMissing {
producer: String,
path: String,
},
DivergentProvisionalAncestor {
producer: String,
path: String,
},
UndeclaredProvisionalConsumption {
pin_key: String,
},
ConsumptionLineageCancelled {
pin_key: String,
},
NativeChildUnresolved {
child_action_key: String,
},
Store(StoreError),
}
impl From<StoreError> for OfferRefusal {
fn from(e: StoreError) -> Self {
Self::Store(e)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ConsumerEscalation {
pub consumer: String,
pub trust_state: String,
pub decision: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct DivergenceQuarantine {
pub class: DivergenceClass,
pub incident_seq: u64,
pub candidate_pin_id: u128,
pub escalations: Vec<ConsumerEscalation>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum PublicationOutcome {
Committed(ActionPublicationRecord),
IdempotentEvidenceAppended,
Quarantined(DivergenceQuarantine),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CommitDurabilityProfile {
RequireDurableClosure,
AcceptVolatileLocations,
}
const fn result_kind_to_tag(kind: ResultKind) -> ResultKindTag {
match kind {
ResultKind::Success => ResultKindTag::Success,
ResultKind::DeterministicFailure => ResultKindTag::DeterministicFailure,
}
}
#[allow(clippy::too_many_arguments)]
pub fn process_offer(
store: &mut dyn RabsMetadataStore,
offer: &OfferPreparedActionResult,
expected_descriptor: &TypedDigest,
manifest_resolver: impl Fn(&str) -> Option<CanonicalActionResultManifest>,
pin_id: u128,
seq: u64,
durability: CommitDurabilityProfile,
own_monotonic_now_ms: impl Fn() -> u64,
) -> Result<PublicationOutcome, OfferRefusal> {
let offered_authority = authority_digest(&offer.authority.coordinator);
let active = store.active_authority()?;
match active {
Some(row) if row.digest == offered_authority => {}
_ => return Err(OfferRefusal::NotActiveAuthority),
}
if offer
.authority
.action_generation
.created_under_authority_digest
!= offered_authority
{
return Err(OfferRefusal::GenerationAuthorityMismatch);
}
let generation_id = offer.authority.action_generation.generation_id.0;
match store.generation_state(generation_id)? {
None => return Err(OfferRefusal::UnknownGeneration),
Some(state) if state.tombstoned => return Err(OfferRefusal::GenerationTombstoned),
Some(_) => {}
}
let attempt_id = offer.authority.attempt_id.0;
if !store.attempt_exists(attempt_id, generation_id)? {
return Err(OfferRefusal::UnknownAttempt);
}
store
.validate_attempt_lease(&offer.authority, own_monotonic_now_ms())
.map_err(lease_refusal)?;
if offer.manifest.canonical_descriptor_digest != *expected_descriptor {
return Err(OfferRefusal::DescriptorMismatch);
}
if offer.manifest.action_key != offer.authority.action_key {
return Err(OfferRefusal::ActionKeyMismatch);
}
let entry = store
.lookup_action(&offer.manifest.action_key)?
.ok_or(OfferRefusal::UnknownActionEntry)?;
if entry.key_epoch != offer.manifest.key_epoch
|| entry.projection_epoch != offer.manifest.projection_epoch
{
return Err(OfferRefusal::EpochMismatch);
}
if let Err(rule) = rabs_key::logical_output_map::verify_manifest_bundle_root(&offer.manifest) {
return Err(OfferRefusal::BundleRootMismatch {
rule: format!("{rule:?}"),
});
}
if semantic_result_digest_v1(&offer.manifest) != offer.manifest.semantic_result_digest {
return Err(OfferRefusal::SemanticDigestMismatch);
}
if observable_result_digest_v1(&offer.manifest, &offer.canonical_observations)
!= offer.manifest.observable_result_digest
{
return Err(OfferRefusal::ObservableDigestMismatch);
}
let mut closure: Vec<&TypedDigest> = vec![&offer.manifest_id.0, &offer.evidence_id.0];
if let Some(root) = &offer.manifest.artifact_bundle_root {
closure.push(&root.0);
}
for output in &offer.manifest.logical_outputs {
closure.push(&output.object.0);
}
closure.push(&offer.evidence.execution_snapshot_root.0);
closure.push(&offer.evidence.observed_input_report.0);
closure.push(&offer.evidence.raw_process_and_event_evidence.0);
closure.push(&offer.evidence.provenance_receipt.0);
if let Some(snapshot) = &offer.evidence.incremental_snapshot {
closure.push(&snapshot.0);
}
for object in closure {
if !store.object_located(object)? {
return Err(OfferRefusal::IncompleteObjectClosure {
missing: digest_key(object),
});
}
if durability == CommitDurabilityProfile::RequireDurableClosure
&& !store.object_durably_located(object)?
{
return Err(OfferRefusal::ObjectNotDurable {
missing: digest_key(object),
});
}
}
let ancestor_rows = verify_provisional_ancestry(store, offer, &manifest_resolver)?;
verify_consumption_obligations(store, offer, &manifest_resolver)?;
crate::native_children::enforce_native_children_resolved(store, &offer.manifest.action_key)
.map_err(|err| match err {
crate::native_children::NativeChildError::ChildUnresolved { child_action_key } => {
OfferRefusal::NativeChildUnresolved { child_action_key }
}
crate::native_children::NativeChildError::Store(store_err) => {
OfferRefusal::Store(store_err)
}
})?;
if let Some(committed_key) = store.published_manifest_key(&offer.manifest.action_key)? {
store
.validate_attempt_lease(&offer.authority, own_monotonic_now_ms())
.map_err(lease_refusal)?;
let candidate_key = digest_key(&offer.manifest_id.0);
if committed_key == candidate_key {
store
.append_evidence_for_attempt(
&offer.authority,
&committed_key,
&offer.evidence_id.0,
&own_monotonic_now_ms,
)
.map_err(lease_refusal)?;
return Ok(PublicationOutcome::IdempotentEvidenceAppended);
}
let committed =
manifest_resolver(&committed_key).ok_or(OfferRefusal::CommittedManifestUnavailable)?;
store
.validate_attempt_lease(&offer.authority, own_monotonic_now_ms())
.map_err(lease_refusal)?;
let class = if committed.semantic_result_digest != offer.manifest.semantic_result_digest {
DivergenceClass::SemanticDivergence
} else if committed.observable_result_digest != offer.manifest.observable_result_digest {
DivergenceClass::ObservableOnlyDivergence
} else {
DivergenceClass::ProjectionCompletenessIncident
};
let quarantine = quarantine_divergence(
store,
&offered_authority,
offer,
class,
&committed_key,
generation_id,
attempt_id,
pin_id,
seq,
&own_monotonic_now_ms,
)?;
return Ok(PublicationOutcome::Quarantined(quarantine));
}
let row = PublicationRow {
action_key: offer.manifest.action_key.clone(),
descriptor_digest: offer.manifest.canonical_descriptor_digest.clone(),
manifest_digest: offer.manifest_id.0.clone(),
evidence_digest: offer.evidence_id.0.clone(),
winner_generation: generation_id,
winner_attempt: attempt_id,
result_kind: result_kind_to_tag(offer.manifest.result_kind),
pin_id,
pin_owner: "coordinator".to_owned(),
provisional_ancestors: ancestor_rows,
};
match store
.commit_publication(
&offered_authority,
Some((&offer.authority, &own_monotonic_now_ms)),
&row,
)
.map_err(lease_refusal)?
{
CommitOutcome::Committed => {}
CommitOutcome::IdempotentDuplicate | CommitOutcome::ConflictQuarantined => {
return Err(OfferRefusal::Store(StoreError::Corruption(
"publication row appeared during admission".into(),
)));
}
}
Ok(PublicationOutcome::Committed(ActionPublicationRecord {
action_key: offer.manifest.action_key.clone(),
canonical_result_manifest_id: offer.manifest_id.clone(),
winner_evidence_bundle_id: offer.evidence_id.clone(),
committed_causal_sequence: seq,
}))
}
fn lease_refusal(error: StoreError) -> OfferRefusal {
match error {
StoreError::UnknownLease => OfferRefusal::UnknownLease,
StoreError::LeaseReleased => OfferRefusal::LeaseReleased,
StoreError::LeaseExpired => OfferRefusal::LeaseExpired,
error => OfferRefusal::Store(error),
}
}
struct PendingAncestorCheck {
producer_key: String,
role_tag: String,
path: Vec<u8>,
consumed_key: String,
direct: bool,
}
fn verify_provisional_ancestry(
store: &mut dyn RabsMetadataStore,
offer: &OfferPreparedActionResult,
manifest_resolver: &impl Fn(&str) -> Option<CanonicalActionResultManifest>,
) -> Result<Vec<ProvisionalAncestorRow>, OfferRefusal> {
let mut direct_rows = Vec::new();
let mut queue: std::collections::VecDeque<PendingAncestorCheck> = offer
.provisional_ancestors
.iter()
.map(|a| PendingAncestorCheck {
producer_key: digest_key(&a.producer_action_key),
role_tag: output_role_name(a.role).to_owned(),
path: a.virtual_path.as_bytes().to_vec(),
consumed_key: digest_key(&a.consumed_object.0),
direct: true,
})
.collect();
let mut walked_producers = std::collections::BTreeSet::new();
while let Some(check) = queue.pop_front() {
let manifest_key = store
.published_manifest_key_str(&check.producer_key)?
.ok_or_else(|| OfferRefusal::ProvisionalProducerNotCommitted {
producer: check.producer_key.clone(),
})?;
let manifest = manifest_resolver(&manifest_key).ok_or_else(|| {
OfferRefusal::AncestorManifestUnavailable {
producer: check.producer_key.clone(),
}
})?;
let output = manifest
.logical_outputs
.iter()
.find(|o| {
output_role_name(o.role) == check.role_tag
&& o.virtual_path.as_bytes() == check.path.as_slice()
})
.ok_or_else(|| OfferRefusal::AncestorOutputMissing {
producer: check.producer_key.clone(),
path: RawBytes::new(check.path.clone()).escaped(),
})?;
let committed_key = digest_key(&output.object.0);
let adopted = if committed_key == check.consumed_key {
false
} else if store.has_adoption_edge(
&check.producer_key,
&check.role_tag,
&check.path,
&check.consumed_key,
&committed_key,
)? {
true
} else {
return Err(OfferRefusal::DivergentProvisionalAncestor {
producer: check.producer_key.clone(),
path: RawBytes::new(check.path.clone()).escaped(),
});
};
if check.direct {
direct_rows.push(ProvisionalAncestorRow {
producer_action_key: check.producer_key.clone(),
role: check.role_tag.clone(),
virtual_path: check.path.clone(),
object_key: check.consumed_key.clone(),
adopted,
});
}
if walked_producers.insert(check.producer_key.clone()) {
for row in store.list_provisional_ancestors(&check.producer_key)? {
if walked_producers.contains(&row.producer_action_key) {
continue;
}
queue.push_back(PendingAncestorCheck {
producer_key: row.producer_action_key,
role_tag: row.role,
path: row.virtual_path,
consumed_key: row.object_key,
direct: false,
});
}
}
}
Ok(direct_rows)
}
fn verify_consumption_obligations(
store: &mut dyn RabsMetadataStore,
offer: &OfferPreparedActionResult,
manifest_resolver: &impl Fn(&str) -> Option<CanonicalActionResultManifest>,
) -> Result<(), OfferRefusal> {
let rows = store.list_open_provisional_obligations(
offer.authority.worker_peer_id.0.as_str(),
&format!("{:032x}", offer.authority.attempt_id.0),
)?;
for obligation in rows {
if obligation.status == "cancelled" {
return Err(OfferRefusal::ConsumptionLineageCancelled {
pin_key: obligation.pin_key,
});
}
let declared = offer.provisional_ancestors.iter().any(|a| {
digest_key(&a.producer_action_key) == obligation.producer_action_key
&& output_role_tag(a.role) == u64::try_from(obligation.role_tag).unwrap_or(u64::MAX)
&& a.virtual_path.as_bytes() == obligation.virtual_path.as_slice()
});
if !declared {
return Err(OfferRefusal::UndeclaredProvisionalConsumption {
pin_key: obligation.pin_key,
});
}
let manifest_key = store
.published_manifest_key_str(&obligation.producer_action_key)?
.ok_or_else(|| OfferRefusal::ProvisionalProducerNotCommitted {
producer: obligation.producer_action_key.clone(),
})?;
let manifest = manifest_resolver(&manifest_key).ok_or_else(|| {
OfferRefusal::AncestorManifestUnavailable {
producer: obligation.producer_action_key.clone(),
}
})?;
let output = manifest
.logical_outputs
.iter()
.find(|o| {
output_role_tag(o.role) == u64::try_from(obligation.role_tag).unwrap_or(u64::MAX)
&& o.virtual_path.as_bytes() == obligation.virtual_path.as_slice()
})
.ok_or_else(|| OfferRefusal::AncestorOutputMissing {
producer: obligation.producer_action_key.clone(),
path: RawBytes::new(obligation.virtual_path.clone()).escaped(),
})?;
let committed_key = digest_key(&output.object.0);
if committed_key != obligation.object_key
&& !store.has_adoption_edge(
&obligation.producer_action_key,
output_role_name_for_tag(obligation.role_tag),
&obligation.virtual_path,
&obligation.object_key,
&committed_key,
)?
{
return Err(OfferRefusal::DivergentProvisionalAncestor {
producer: obligation.producer_action_key.clone(),
path: RawBytes::new(obligation.virtual_path.clone()).escaped(),
});
}
store.resolve_provisional_obligations(&obligation.pin_key, &committed_key)?;
}
Ok(())
}
const fn divergence_class_tag(class: DivergenceClass) -> &'static str {
match class {
DivergenceClass::IdempotentSameResult => "idempotent",
DivergenceClass::SemanticDivergence => "semantic",
DivergenceClass::ObservableOnlyDivergence => "observable-only",
DivergenceClass::ProjectionCompletenessIncident => "projection-completeness",
}
}
const fn divergence_quarantine_reason(class: DivergenceClass) -> &'static str {
match class {
DivergenceClass::IdempotentSameResult => "same-key re-offer",
DivergenceClass::SemanticDivergence => {
"semantic divergence: determinism/key-soundness incident"
}
DivergenceClass::ObservableOnlyDivergence => {
"observable-only divergence: presentation quarantine"
}
DivergenceClass::ProjectionCompletenessIncident => {
"projection completeness incident: equal digests, different manifests"
}
}
}
fn escalation_decision(trust_state: &str) -> &'static str {
match trust_state {
"ci-policy-approved" | "project-release-eligible" => "recall-and-reverify",
_ => "notify-and-reverify",
}
}
fn set_divergence_serving_disposition(
store: &mut dyn RabsMetadataStore,
action_key: &str,
class: DivergenceClass,
attempt_lease: Option<(&AttemptAuthority, &dyn Fn() -> u64)>,
) -> Result<(), StoreError> {
let mut disposition = match class {
DivergenceClass::SemanticDivergence | DivergenceClass::ProjectionCompletenessIncident => {
DISPOSITION_QUARANTINED
}
DivergenceClass::ObservableOnlyDivergence => DISPOSITION_PRESENTATION_QUARANTINED,
DivergenceClass::IdempotentSameResult => {
return Err(StoreError::Corruption(
"idempotent result entered divergence quarantine".into(),
));
}
};
if disposition == DISPOSITION_PRESENTATION_QUARANTINED
&& store.serving_disposition_key(action_key)?.as_deref() == Some(DISPOSITION_QUARANTINED)
{
disposition = DISPOSITION_QUARANTINED;
}
match attempt_lease {
Some((authority, clock)) => {
if digest_key(&authority.action_key) != action_key {
return Err(StoreError::AttemptAuthorityMismatch);
}
store.set_serving_disposition_for_attempt(authority, disposition, clock)
}
None => store.set_serving_disposition_key(action_key, disposition),
}
}
#[allow(clippy::too_many_arguments)]
fn quarantine_divergence(
store: &mut dyn RabsMetadataStore,
authority: &TypedDigest,
offer: &OfferPreparedActionResult,
class: DivergenceClass,
committed_key: &str,
generation_id: u128,
attempt_id: u128,
pin_id: u128,
seq: u64,
own_monotonic_now_ms: &dyn Fn() -> u64,
) -> Result<DivergenceQuarantine, OfferRefusal> {
let action_key = digest_key(&offer.manifest.action_key);
let candidate_key = digest_key(&offer.manifest_id.0);
set_divergence_serving_disposition(
store,
&action_key,
class,
Some((&offer.authority, own_monotonic_now_ms)),
)
.map_err(lease_refusal)?;
if class != DivergenceClass::ObservableOnlyDivergence {
store.add_quarantine(
QuarantineScope::ActionEntry,
&action_key,
divergence_quarantine_reason(class),
)?;
}
store.add_object_edge(
&offer.manifest_id.0,
&offer.evidence_id.0,
"divergence-candidate",
)?;
if let Some(root) = &offer.manifest.artifact_bundle_root {
store.add_object_edge(&offer.manifest_id.0, &root.0, "divergence-candidate")?;
}
for output in &offer.manifest.logical_outputs {
store.add_object_edge(
&offer.manifest_id.0,
&output.object.0,
"divergence-candidate",
)?;
}
store.create_pin(
pin_id,
&offer.manifest_id.0,
"coordinator",
DIVERGENCE_EVIDENCE_PIN_CLASS,
None,
Some(&format!("divergence-incident:{action_key}:{seq}")),
true,
divergence_quarantine_reason(class),
)?;
store.append_evidence(
&offer.manifest.action_key,
&candidate_key,
&offer.evidence_id.0,
generation_id,
attempt_id,
)?;
store.record_divergence_incident(
authority,
&DivergenceIncidentRow {
action_key: action_key.clone(),
seq,
class: divergence_class_tag(class).to_owned(),
committed_manifest_key: committed_key.to_owned(),
candidate_manifest_key: candidate_key,
candidate_evidence_key: digest_key(&offer.evidence_id.0),
candidate_pin_hex: format!("{pin_id:032x}"),
generation_hex: format!("{generation_id:032x}"),
attempt_hex: format!("{attempt_id:032x}"),
detail: divergence_quarantine_reason(class).to_owned(),
},
)?;
let mut escalations = Vec::new();
if class == DivergenceClass::SemanticDivergence {
let trust_state = store
.latest_trust_evaluation(&offer.manifest.action_key)?
.map_or_else(|| "unevaluated".to_owned(), |row| row.state);
for consumer in store.list_served_consumers(&action_key)? {
let decision = escalation_decision(&trust_state);
store.record_decision_receipt(
"divergence-escalation",
&consumer,
seq,
decision,
&format!("semantic divergence on {action_key}; trust {trust_state}"),
)?;
escalations.push(ConsumerEscalation {
consumer,
trust_state: trust_state.clone(),
decision: decision.to_owned(),
});
}
}
Ok(DivergenceQuarantine {
class,
incident_seq: seq,
candidate_pin_id: pin_id,
escalations,
})
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum RecomputationCheck {
NoTombstone,
Reproduced,
Divergence(DivergenceClass),
}
pub fn retain_eviction_tombstone(
store: &mut dyn RabsMetadataStore,
manifest: &CanonicalActionResultManifest,
evicted_seq: u64,
) -> Result<(), StoreError> {
store.record_eviction_tombstone(
&manifest.action_key,
&manifest.semantic_result_digest,
&manifest.observable_result_digest,
evicted_seq,
)
}
pub fn check_recomputation_against_tombstone(
store: &mut dyn RabsMetadataStore,
action: &TypedDigest,
new_semantic: &TypedDigest,
new_observable: &TypedDigest,
) -> Result<RecomputationCheck, StoreError> {
let Some((retained_semantic, retained_observable)) = store.eviction_tombstone(action)? else {
return Ok(RecomputationCheck::NoTombstone);
};
if retained_semantic == *new_semantic && retained_observable == *new_observable {
store.consume_eviction_tombstone(action)?;
return Ok(RecomputationCheck::Reproduced);
}
let class = if retained_semantic == *new_semantic {
DivergenceClass::ObservableOnlyDivergence
} else {
DivergenceClass::SemanticDivergence
};
let action_key = digest_key(action);
set_divergence_serving_disposition(store, &action_key, class, None)?;
if class == DivergenceClass::SemanticDivergence {
store.add_quarantine(
QuarantineScope::ActionEntry,
&action_key,
"post-eviction recomputation divergence",
)?;
}
Ok(RecomputationCheck::Divergence(class))
}
#[cfg(test)]
mod tests {
use super::*;
use rabs_protocol::authority::{ClusterId, CoordinatorAuthority, CoordinatorIncarnationId};
use rabs_protocol::generation::{
ActionGeneration, ActionGenerationId, AttemptId, ExecutionLeaseId, LeaseRenewal,
LeaseRenewalSeq, WorkerBootGeneration, WorkerIncarnationId,
};
use rabs_protocol::result_identity::LogicalOutput;
use rabs_protocol::wire_time::PeerId;
use rabs_protocol::worker_fence::WorkerSessionOffer;
use crate::metadata_store::{
ActionEntryRow, AuthorityRow, FsqliteEngine, ProvisionalObligationInsert, RusqliteEngine,
SqlMetadataStore,
};
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-h011-{}-{}-{}.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) -> ObjectId {
ObjectId(digest("rabs.object.sha256.v1", tag))
}
fn content_object(bytes: &[u8]) -> ObjectId {
ObjectId(
crate::digest_set::digest_set(bytes, crate::digest_set::DigestRequest::default(), None)
.unwrap()
.atp_content_id,
)
}
fn coordinator_authority() -> CoordinatorAuthority {
CoordinatorAuthority {
cluster_id: ClusterId("cluster-a".to_owned()),
credential_generation: 1,
term: 3,
incarnation_id: CoordinatorIncarnationId(77),
}
}
fn attempt_authority() -> AttemptAuthority {
let coordinator = coordinator_authority();
let created_under = authority_digest(&coordinator);
AttemptAuthority {
coordinator,
action_key: digest("rabs.action-key.sha256.v1", 7),
action_generation: ActionGeneration {
generation_id: ActionGenerationId(11),
per_key_ordinal: 1,
created_under_authority_digest: created_under,
},
attempt_id: AttemptId(20),
execution_lease_id: ExecutionLeaseId(30),
lease_renewal_seq: LeaseRenewalSeq(1),
worker_peer_id: PeerId("worker-a".to_owned()),
worker_boot_generation: WorkerBootGeneration(1),
worker_incarnation_id: WorkerIncarnationId(5),
}
}
fn distinct_attempt_authority() -> AttemptAuthority {
let mut authority = attempt_authority();
authority.attempt_id = AttemptId(21);
authority.execution_lease_id = ExecutionLeaseId(31);
authority.worker_peer_id = PeerId("worker-b".to_owned());
authority.worker_boot_generation = WorkerBootGeneration(2);
authority.worker_incarnation_id = WorkerIncarnationId(6);
authority
}
fn manifest() -> CanonicalActionResultManifest {
CanonicalActionResultManifest {
action_key: digest("rabs.action-key.sha256.v1", 7),
canonical_descriptor_digest: digest("rabs.descriptor.sha256.v1", 8),
key_epoch: 1,
projection_epoch: 1,
result_kind: ResultKind::Success,
artifact_bundle_root: Some(object(40)),
logical_outputs: vec![LogicalOutput {
role: OutputRole::Materializable,
virtual_path: RawBytes::new(b"out/lib.rlib".to_vec()),
object: object(41),
}],
semantic_result_digest: digest(SEMANTIC_PROJECTION_DOMAIN, 0),
observable_result_digest: digest(OBSERVABLE_PROJECTION_DOMAIN, 0),
}
}
fn evidence(manifest_id: &ObjectId) -> AttemptEvidenceBundle {
AttemptEvidenceBundle {
action_key: digest("rabs.action-key.sha256.v1", 7),
canonical_result_manifest_id: manifest_id.clone(),
execution_snapshot_root: object(60),
observed_input_report: object(61),
raw_process_and_event_evidence: object(62),
provenance_receipt: object(63),
incremental_snapshot: None,
}
}
fn declared() -> Vec<(OutputRole, RawBytes)> {
vec![(
OutputRole::Materializable,
RawBytes::new(b"out/lib.rlib".to_vec()),
)]
}
fn offer() -> OfferPreparedActionResult {
let manifest_id = object(50);
OfferPreparedActionResult::build(
attempt_authority(),
manifest(),
manifest_id.clone(),
evidence(&manifest_id),
object(51),
digest("rabs.observation-stream.sha256.v1", 9),
&declared(),
Vec::new(),
)
.unwrap()
}
fn ready_store(store: &mut dyn RabsMetadataStore) {
let attempt_authority = attempt_authority();
let auth = authority_digest(&attempt_authority.coordinator);
store
.acquire_authority(&AuthorityRow {
digest: auth.clone(),
cluster_id: "cluster-a".to_owned(),
incarnation: 77,
term: 3,
acquired_seq: 1,
})
.unwrap();
store
.upsert_action_entry(&ActionEntryRow {
action_key: digest("rabs.action-key.sha256.v1", 7),
key_epoch: 1,
projection_epoch: 1,
})
.unwrap();
store
.create_bound_generation(
&auth,
&attempt_authority.action_generation,
&digest("rabs.action-key.sha256.v1", 7),
)
.unwrap();
store
.admit_worker_session(
&auth,
&WorkerSessionOffer {
worker_peer_id: attempt_authority.worker_peer_id.clone(),
boot_generation: attempt_authority.worker_boot_generation,
incarnation: attempt_authority.worker_incarnation_id,
reenrollment_proof: None,
},
1,
)
.unwrap();
store
.admit_attempt_lease(&attempt_authority, 1, 100)
.unwrap();
for tag in [40, 41, 50, 51, 60, 61, 62, 63] {
let id = object(tag);
store.record_object(&id.0, 64).unwrap();
store
.add_location(&id.0, &format!("/cas/{tag}"), Some(1), "raw", true)
.unwrap();
}
locate_bundle_root(store, &offer());
}
fn locate_bundle_root(store: &mut dyn RabsMetadataStore, offer: &OfferPreparedActionResult) {
if let Some(root) = &offer.manifest.artifact_bundle_root {
store.record_object(&root.0, 0).unwrap();
store
.add_location(&root.0, "/cas/bundle-root", Some(1), "raw", true)
.unwrap();
}
}
fn locate_test_object(
store: &mut dyn RabsMetadataStore,
object: &ObjectId,
path: &str,
logical_size: usize,
) {
store
.record_object(&object.0, u64::try_from(logical_size).unwrap())
.unwrap();
store
.add_location(&object.0, path, Some(1), "raw", true)
.unwrap();
}
fn expected_descriptor() -> TypedDigest {
digest("rabs.descriptor.sha256.v1", 8)
}
fn no_committed(_: &str) -> Option<CanonicalActionResultManifest> {
None
}
#[allow(clippy::too_many_arguments)]
fn plant_producer(
store: &mut dyn RabsMetadataStore,
resolver_map: &mut std::collections::BTreeMap<String, CanonicalActionResultManifest>,
key_tag: u8,
manifest_tag: u8,
path: &str,
output_tag: u8,
pin: u128,
ancestors: Vec<ProvisionalAncestorRow>,
) {
let action = digest("rabs.action-key.sha256.v1", key_tag);
let mut producer_manifest = manifest();
producer_manifest.action_key = action.clone();
producer_manifest.logical_outputs = vec![LogicalOutput {
role: OutputRole::ProvisionalMetadata,
virtual_path: RawBytes::from(path),
object: object(output_tag),
}];
resolver_map.insert(digest_key(&object(manifest_tag).0), producer_manifest);
let auth = authority_digest(&coordinator_authority());
assert_eq!(
store
.commit_publication(
&auth,
None,
&PublicationRow {
action_key: action,
descriptor_digest: expected_descriptor(),
manifest_digest: object(manifest_tag).0,
evidence_digest: object(manifest_tag).0.clone(),
winner_generation: 11,
winner_attempt: 20,
result_kind: ResultKindTag::Success,
pin_id: pin,
pin_owner: "coordinator".to_owned(),
provisional_ancestors: ancestors,
},
)
.unwrap(),
CommitOutcome::Committed
);
}
fn ancestor_ref(producer_tag: u8, path: &str, consumed_tag: u8) -> ProvisionalAncestorRef {
ProvisionalAncestorRef {
producer_action_key: digest("rabs.action-key.sha256.v1", producer_tag),
role: OutputRole::ProvisionalMetadata,
virtual_path: RawBytes::from(path),
consumed_object: object(consumed_tag),
}
}
fn offer_with_ancestors(ancestors: Vec<ProvisionalAncestorRef>) -> OfferPreparedActionResult {
let manifest_id = object(50);
OfferPreparedActionResult::build(
attempt_authority(),
manifest(),
manifest_id.clone(),
evidence(&manifest_id),
object(51),
digest("rabs.observation-stream.sha256.v1", 9),
&declared(),
ancestors,
)
.unwrap()
}
#[test]
fn h028_transitive_ancestor_closure_adoption_and_divergence() {
let engine = RusqliteEngine::open_in_memory().unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
ready_store(&mut store);
let mut manifests = std::collections::BTreeMap::new();
plant_producer(
&mut store,
&mut manifests,
100,
120,
"out/a.rmeta",
110,
800,
vec![],
);
plant_producer(
&mut store,
&mut manifests,
101,
121,
"out/b.rmeta",
111,
801,
vec![ProvisionalAncestorRow {
producer_action_key: digest_key(&digest("rabs.action-key.sha256.v1", 100)),
role: "provisional-metadata".to_owned(),
virtual_path: b"out/a.rmeta".to_vec(),
object_key: digest_key(&object(110).0),
adopted: false,
}],
);
let resolver = {
let manifests = manifests.clone();
move |key: &str| manifests.get(key).cloned()
};
let c_offer = offer_with_ancestors(vec![ancestor_ref(101, "out/b.rmeta", 111)]);
assert!(matches!(
process_offer(
&mut store,
&c_offer,
&expected_descriptor(),
resolver.clone(),
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
)
.unwrap(),
PublicationOutcome::Committed(_)
));
let c_key = digest_key(&digest("rabs.action-key.sha256.v1", 7));
let recorded = store.list_provisional_ancestors(&c_key).unwrap();
assert_eq!(recorded.len(), 1);
assert_eq!(
recorded[0].producer_action_key,
digest_key(&digest("rabs.action-key.sha256.v1", 101))
);
assert_eq!(recorded[0].object_key, digest_key(&object(111).0));
assert!(!recorded[0].adopted);
}
#[test]
fn m008_transitive_debt_blocks_commit_until_producers_finalize() {
let engine = RusqliteEngine::open_in_memory().unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
ready_store(&mut store);
let mut manifests = std::collections::BTreeMap::new();
plant_producer(
&mut store,
&mut manifests,
100,
120,
"out/a.rmeta",
110,
800,
vec![],
);
store
.record_provisional_consumption(&ProvisionalObligationInsert {
consumer_worker: "worker-a".to_owned(),
consumer_attempt: 20,
pin_key: "test-pin-b".to_owned(),
producer_action_key: digest_key(&digest("rabs.action-key.sha256.v1", 101)),
producer_generation: 11,
producer_attempt: 21,
role_tag: 2,
virtual_path: b"out/b.rmeta".to_vec(),
object_key: digest_key(&object(111).0),
created_seq: 1,
})
.unwrap();
let resolver = {
let manifests = manifests.clone();
move |key: &str| manifests.get(key).cloned()
};
let c_offer = offer_with_ancestors(vec![ancestor_ref(101, "out/b.rmeta", 111)]);
assert!(matches!(
process_offer(
&mut store,
&c_offer,
&expected_descriptor(),
resolver,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
),
Err(OfferRefusal::ProvisionalProducerNotCommitted { .. })
));
assert_eq!(
store
.count_open_provisional_obligations("test-pin-b")
.unwrap(),
1
);
plant_producer(
&mut store,
&mut manifests,
101,
121,
"out/b.rmeta",
111,
801,
vec![ProvisionalAncestorRow {
producer_action_key: digest_key(&digest("rabs.action-key.sha256.v1", 100)),
role: "provisional-metadata".to_owned(),
virtual_path: b"out/a.rmeta".to_vec(),
object_key: digest_key(&object(110).0),
adopted: false,
}],
);
let resolver = {
let manifests = manifests.clone();
move |key: &str| manifests.get(key).cloned()
};
assert!(matches!(
process_offer(
&mut store,
&c_offer,
&expected_descriptor(),
resolver,
901,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
),
Ok(PublicationOutcome::Committed(_))
));
assert_eq!(
store
.count_open_provisional_obligations("test-pin-b")
.unwrap(),
0
);
}
#[test]
fn m008_concealed_consumption_refuses_commit() {
let engine = RusqliteEngine::open_in_memory().unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
ready_store(&mut store);
let mut manifests = std::collections::BTreeMap::new();
plant_producer(
&mut store,
&mut manifests,
100,
120,
"out/a.rmeta",
110,
800,
vec![],
);
let pin_key = "test-pin-concealed".to_owned();
store
.record_provisional_consumption(&ProvisionalObligationInsert {
consumer_worker: "worker-a".to_owned(),
consumer_attempt: 20,
pin_key: pin_key.clone(),
producer_action_key: digest_key(&digest("rabs.action-key.sha256.v1", 100)),
producer_generation: 11,
producer_attempt: 20,
role_tag: 2,
virtual_path: b"out/a.rmeta".to_vec(),
object_key: digest_key(&object(110).0),
created_seq: 1,
})
.unwrap();
let resolver = {
let manifests = manifests.clone();
move |key: &str| manifests.get(key).cloned()
};
let c_offer = offer_with_ancestors(vec![]);
assert_eq!(
process_offer(
&mut store,
&c_offer,
&expected_descriptor(),
resolver,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
)
.unwrap_err(),
OfferRefusal::UndeclaredProvisionalConsumption { pin_key }
);
}
#[test]
fn m008_cancelled_lineage_permanently_refuses_commit() {
let engine = RusqliteEngine::open_in_memory().unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
ready_store(&mut store);
let mut manifests = std::collections::BTreeMap::new();
plant_producer(
&mut store,
&mut manifests,
100,
120,
"out/a.rmeta",
110,
800,
vec![],
);
let pin_key = "test-pin-cancelled".to_owned();
store
.record_provisional_consumption(&ProvisionalObligationInsert {
consumer_worker: "worker-a".to_owned(),
consumer_attempt: 20,
pin_key: pin_key.clone(),
producer_action_key: digest_key(&digest("rabs.action-key.sha256.v1", 100)),
producer_generation: 11,
producer_attempt: 20,
role_tag: 2,
virtual_path: b"out/a.rmeta".to_vec(),
object_key: digest_key(&object(110).0),
created_seq: 1,
})
.unwrap();
store.cancel_provisional_obligations(&pin_key).unwrap();
let resolver = {
let manifests = manifests.clone();
move |key: &str| manifests.get(key).cloned()
};
let c_offer = offer_with_ancestors(vec![ancestor_ref(100, "out/a.rmeta", 110)]);
assert_eq!(
process_offer(
&mut store,
&c_offer,
&expected_descriptor(),
resolver,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
)
.unwrap_err(),
OfferRefusal::ConsumptionLineageCancelled { pin_key }
);
}
#[test]
fn h028_transitive_hole_refuses_where_direct_only_would_commit() {
let engine = RusqliteEngine::open_in_memory().unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
ready_store(&mut store);
let mut manifests = std::collections::BTreeMap::new();
plant_producer(
&mut store,
&mut manifests,
101,
121,
"out/b.rmeta",
111,
801,
vec![ProvisionalAncestorRow {
producer_action_key: digest_key(&digest("rabs.action-key.sha256.v1", 102)),
role: "provisional-metadata".to_owned(),
virtual_path: b"out/x.rmeta".to_vec(),
object_key: digest_key(&object(112).0),
adopted: false,
}],
);
let resolver = move |key: &str| manifests.get(key).cloned();
let c_offer = offer_with_ancestors(vec![ancestor_ref(101, "out/b.rmeta", 111)]);
assert_eq!(
process_offer(
&mut store,
&c_offer,
&expected_descriptor(),
resolver,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
),
Err(OfferRefusal::ProvisionalProducerNotCommitted {
producer: digest_key(&digest("rabs.action-key.sha256.v1", 102)),
})
);
assert!(
!store
.has_publication(&digest("rabs.action-key.sha256.v1", 7))
.unwrap()
);
}
#[test]
fn h028_divergent_consumed_object_refuses_until_adoption_edge() {
let engine = RusqliteEngine::open_in_memory().unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
ready_store(&mut store);
let mut manifests = std::collections::BTreeMap::new();
plant_producer(
&mut store,
&mut manifests,
101,
121,
"out/b.rmeta",
111,
801,
vec![],
);
let resolver = {
let manifests = manifests.clone();
move |key: &str| manifests.get(key).cloned()
};
let b_key = digest_key(&digest("rabs.action-key.sha256.v1", 101));
let c_offer = offer_with_ancestors(vec![ancestor_ref(101, "out/b.rmeta", 119)]);
assert_eq!(
process_offer(
&mut store,
&c_offer,
&expected_descriptor(),
resolver.clone(),
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
),
Err(OfferRefusal::DivergentProvisionalAncestor {
producer: b_key.clone(),
path: "out/b.rmeta".to_owned(),
})
);
let auth = authority_digest(&coordinator_authority());
store
.record_adoption_edge(
&auth,
&b_key,
"provisional-metadata",
b"out/b.rmeta",
&digest_key(&object(119).0),
&digest_key(&object(111).0),
)
.unwrap();
assert_eq!(
store.record_adoption_edge(
&auth,
&b_key,
"provisional-metadata",
b"out/b.rmeta",
&digest_key(&object(119).0),
&digest_key(&object(112).0),
),
Err(StoreError::AdoptionEdgeConflict)
);
assert!(matches!(
process_offer(
&mut store,
&c_offer,
&expected_descriptor(),
resolver,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
)
.unwrap(),
PublicationOutcome::Committed(_)
));
let c_key = digest_key(&digest("rabs.action-key.sha256.v1", 7));
let recorded = store.list_provisional_ancestors(&c_key).unwrap();
assert_eq!(recorded.len(), 1);
assert!(
recorded[0].adopted,
"adoption edge must be recorded as such"
);
}
#[test]
fn h028_canonical_ancestor_set_refuses_duplicates_and_self() {
let manifest_id = object(50);
assert_eq!(
OfferPreparedActionResult::build(
attempt_authority(),
manifest(),
manifest_id.clone(),
evidence(&manifest_id),
object(51),
digest("rabs.observation-stream.sha256.v1", 9),
&declared(),
vec![
ancestor_ref(101, "out/b.rmeta", 111),
ancestor_ref(101, "out/b.rmeta", 119),
],
),
Err(OfferBuildError::DuplicateAncestorRef {
path: "out/b.rmeta".to_owned()
})
);
assert_eq!(
OfferPreparedActionResult::build(
attempt_authority(),
manifest(),
manifest_id.clone(),
evidence(&manifest_id),
object(51),
digest("rabs.observation-stream.sha256.v1", 9),
&declared(),
vec![ancestor_ref(7, "out/self.rmeta", 111)],
),
Err(OfferBuildError::SelfAncestorRef)
);
let shuffled = OfferPreparedActionResult::build(
attempt_authority(),
manifest(),
manifest_id.clone(),
evidence(&manifest_id),
object(51),
digest("rabs.observation-stream.sha256.v1", 9),
&declared(),
vec![
ancestor_ref(102, "out/x.rmeta", 112),
ancestor_ref(101, "out/b.rmeta", 111),
],
)
.unwrap();
let ordered = OfferPreparedActionResult::build(
attempt_authority(),
manifest(),
manifest_id.clone(),
evidence(&manifest_id),
object(51),
digest("rabs.observation-stream.sha256.v1", 9),
&declared(),
vec![
ancestor_ref(101, "out/b.rmeta", 111),
ancestor_ref(102, "out/x.rmeta", 112),
],
)
.unwrap();
assert_eq!(shuffled, ordered, "ancestor order never changes identity");
}
#[test]
fn h036_publication_pin_is_coordinator_owned_and_worker_release_refuses() {
let engine = RusqliteEngine::open_in_memory().unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
ready_store(&mut store);
assert!(matches!(
process_offer(
&mut store,
&offer(),
&expected_descriptor(),
no_committed,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
)
.unwrap(),
PublicationOutcome::Committed(_)
));
let pin = store.pin_row(900).unwrap().unwrap();
assert_eq!(pin.owner, "coordinator");
assert_eq!(pin.class, "action-publication");
assert_eq!(pin.root_key, digest_key(&object(50).0));
assert!(!pin.released);
assert_eq!(pin.expires_at_seq, None, "publication pins never expire");
for intruder in ["worker-a", "worker-b", "rabs-wkr:22", ""] {
assert_eq!(
store.release_pin(900, intruder),
Err(StoreError::PinOwnerMismatch),
"worker identity {intruder:?} must be refused"
);
}
let after = store.pin_row(900).unwrap().unwrap();
assert!(!after.released, "refused releases must not release");
assert_eq!(after, pin, "refused releases must change nothing");
assert_eq!(
store.release_pin(999, "coordinator"),
Err(StoreError::UnknownPin)
);
store.release_pin(900, "coordinator").unwrap();
assert!(store.pin_row(900).unwrap().unwrap().released);
}
#[test]
fn h035_equal_digests_different_manifest_opens_projection_completeness_incident() {
let engine = RusqliteEngine::open_in_memory().unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
ready_store(&mut store);
let first = offer();
assert!(matches!(
process_offer(
&mut store,
&first,
&expected_descriptor(),
no_committed,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
)
.unwrap(),
PublicationOutcome::Committed(_)
));
let doppel_id = object(52);
store.record_object(&doppel_id.0, 64).unwrap();
store
.add_location(&doppel_id.0, "/cas/52", Some(1), "raw", true)
.unwrap();
let doppel = OfferPreparedActionResult::build(
attempt_authority(),
manifest(),
doppel_id.clone(),
evidence(&doppel_id),
object(54),
digest("rabs.observation-stream.sha256.v1", 9),
&declared(),
Vec::new(),
)
.unwrap();
store.record_object(&object(54).0, 64).unwrap();
store
.add_location(&object(54).0, "/cas/54", Some(1), "raw", true)
.unwrap();
assert_eq!(
doppel.manifest.semantic_result_digest, first.manifest.semantic_result_digest,
"fixture precondition: equal semantic digests"
);
assert_eq!(
doppel.manifest.observable_result_digest, first.manifest.observable_result_digest,
"fixture precondition: equal observable digests"
);
assert_ne!(doppel.manifest_id, first.manifest_id);
let committed = first.manifest.clone();
let outcome = process_offer(
&mut store,
&doppel,
&expected_descriptor(),
move |_| Some(committed.clone()),
905,
6,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
)
.unwrap();
let PublicationOutcome::Quarantined(quarantine) = outcome else {
panic!("equal-digests/different-manifest must quarantine, got {outcome:?}");
};
assert_eq!(
quarantine.class,
DivergenceClass::ProjectionCompletenessIncident
);
assert!(
quarantine.escalations.is_empty(),
"consumer escalation is the semantic-divergence arm only"
);
let action = digest("rabs.action-key.sha256.v1", 7);
let action_key = digest_key(&action);
assert_eq!(
store.serving_disposition_key(&action_key).unwrap().unwrap(),
DISPOSITION_QUARANTINED
);
let incidents = store.list_divergence_incidents(&action_key).unwrap();
assert_eq!(incidents.len(), 1);
assert_eq!(incidents[0].class, "projection-completeness");
assert_eq!(
incidents[0].committed_manifest_key,
digest_key(&first.manifest_id.0)
);
assert_eq!(
incidents[0].candidate_manifest_key,
digest_key(&doppel_id.0)
);
assert_eq!(
store.published_manifest_key(&action).unwrap().unwrap(),
digest_key(&first.manifest_id.0)
);
let pin = store.pin_row(905).unwrap().unwrap();
assert_eq!(pin.class, DIVERGENCE_EVIDENCE_PIN_CLASS);
assert_eq!(pin.root_key, digest_key(&doppel_id.0));
let candidate_view = store
.list_evidence_keys_for_manifest(&digest_key(&doppel_id.0))
.unwrap();
assert!(candidate_view.contains(&digest_key(&object(54).0)));
let committed_view = store
.list_evidence_keys_for_manifest(&digest_key(&first.manifest_id.0))
.unwrap();
assert!(!committed_view.contains(&digest_key(&object(54).0)));
}
#[test]
fn h039_bundle_root_and_role_uniqueness_are_enforced_at_both_gates() {
let manifest_id = object(50);
let mut duplicated = manifest();
duplicated
.logical_outputs
.push(duplicated.logical_outputs[0].clone());
assert!(matches!(
OfferPreparedActionResult::build(
attempt_authority(),
duplicated,
manifest_id.clone(),
evidence(&manifest_id),
object(51),
digest("rabs.observation-stream.sha256.v1", 9),
&[declared()[0].clone(), declared()[0].clone()],
Vec::new(),
),
Err(OfferBuildError::ManifestInvalid(_))
));
let built = offer();
assert_eq!(
built.manifest.artifact_bundle_root,
rabs_key::logical_output_map::compute_bundle_root(&built.manifest.logical_outputs),
"worker cannot ship a contradictory root"
);
let engine = RusqliteEngine::open_in_memory().unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
ready_store(&mut store);
let mut tampered = offer();
tampered.manifest.artifact_bundle_root = Some(object(99));
tampered.manifest.semantic_result_digest = semantic_result_digest_v1(&tampered.manifest);
tampered.manifest.observable_result_digest = observable_result_digest_v1(
&tampered.manifest,
&digest("rabs.observation-stream.sha256.v1", 9),
);
assert!(matches!(
process_offer(
&mut store,
&tampered,
&expected_descriptor(),
no_committed,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
),
Err(OfferRefusal::BundleRootMismatch { .. })
));
let action = digest("rabs.action-key.sha256.v1", 7);
assert!(
!store.has_publication(&action).unwrap(),
"refusals write nothing"
);
let mut rootless = offer();
rootless.manifest.artifact_bundle_root = None;
rootless.manifest.semantic_result_digest = semantic_result_digest_v1(&rootless.manifest);
rootless.manifest.observable_result_digest = observable_result_digest_v1(
&rootless.manifest,
&digest("rabs.observation-stream.sha256.v1", 9),
);
assert!(matches!(
process_offer(
&mut store,
&rootless,
&expected_descriptor(),
no_committed,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
),
Err(OfferRefusal::BundleRootMismatch { .. })
));
assert!(matches!(
process_offer(
&mut store,
&offer(),
&expected_descriptor(),
no_committed,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
)
.unwrap(),
PublicationOutcome::Committed(_)
));
}
#[test]
fn h032_commit_ack_gates_on_the_configured_durability_profile() {
let engine = RusqliteEngine::open_in_memory().unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
ready_store(&mut store);
store
.add_location(&object(61).0, "/cas/61", Some(1), "raw", false)
.unwrap();
assert_eq!(
process_offer(
&mut store,
&offer(),
&expected_descriptor(),
no_committed,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
),
Err(OfferRefusal::ObjectNotDurable {
missing: digest_key(&object(61).0),
})
);
let action = digest("rabs.action-key.sha256.v1", 7);
assert!(!store.has_publication(&action).unwrap());
store
.set_location_quarantined(&object(61).0, "/cas/61", true)
.unwrap();
assert!(!store.has_publication(&action).unwrap());
let engine = RusqliteEngine::open_in_memory().unwrap();
let mut volatile_store = SqlMetadataStore::open(engine).unwrap();
ready_store(&mut volatile_store);
volatile_store
.add_location(&object(61).0, "/cas/61", Some(1), "raw", false)
.unwrap();
assert!(matches!(
process_offer(
&mut volatile_store,
&offer(),
&expected_descriptor(),
no_committed,
900,
1,
CommitDurabilityProfile::AcceptVolatileLocations,
|| 10,
)
.unwrap(),
PublicationOutcome::Committed(_)
));
let engine = RusqliteEngine::open_in_memory().unwrap();
let mut durable_store = SqlMetadataStore::open(engine).unwrap();
ready_store(&mut durable_store);
assert!(matches!(
process_offer(
&mut durable_store,
&offer(),
&expected_descriptor(),
no_committed,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
)
.unwrap(),
PublicationOutcome::Committed(_)
));
assert!(durable_store.has_publication(&action).unwrap());
for tag in [40, 41, 50, 51, 60, 61, 62, 63] {
assert!(
durable_store
.object_durably_located(&object(tag).0)
.unwrap(),
"committed closure object {tag} must be durable"
);
}
}
#[test]
fn h011_worker_build_validates_and_stamps_projection_digests() {
let built = offer();
assert_eq!(
built.manifest.semantic_result_digest,
semantic_result_digest_v1(&built.manifest)
);
assert_eq!(
built.manifest.observable_result_digest,
observable_result_digest_v1(&built.manifest, &built.canonical_observations)
);
let manifest_id = object(50);
assert_eq!(
OfferPreparedActionResult::build(
attempt_authority(),
manifest(),
manifest_id.clone(),
evidence(&manifest_id),
object(51),
digest("rabs.observation-stream.sha256.v1", 9),
&[],
Vec::new(),
),
Err(OfferBuildError::UndeclaredOutput {
path: "out/lib.rlib".to_owned()
})
);
assert_eq!(
OfferPreparedActionResult::build(
attempt_authority(),
manifest(),
manifest_id,
evidence(&object(99)),
object(51),
digest("rabs.observation-stream.sha256.v1", 9),
&declared(),
Vec::new(),
),
Err(OfferBuildError::EvidenceManifestMismatch)
);
let mut bad = manifest();
bad.result_kind = ResultKind::DeterministicFailure;
let manifest_id = object(50);
assert!(matches!(
OfferPreparedActionResult::build(
attempt_authority(),
bad,
manifest_id.clone(),
evidence(&manifest_id),
object(51),
digest("rabs.observation-stream.sha256.v1", 9),
&declared(),
Vec::new(),
),
Err(OfferBuildError::ManifestInvalid(_))
));
}
#[test]
fn h011_commit_writes_publication_evidence_and_pin_atomically() {
let engine = RusqliteEngine::open_in_memory().unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
ready_store(&mut store);
let outcome = process_offer(
&mut store,
&offer(),
&expected_descriptor(),
no_committed,
900,
42,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
)
.unwrap();
let PublicationOutcome::Committed(receipt) = outcome else {
panic!("expected commit, got {outcome:?}");
};
assert_eq!(receipt.committed_causal_sequence, 42);
assert_eq!(receipt.canonical_result_manifest_id, object(50));
assert_eq!(receipt.winner_evidence_bundle_id, object(51));
assert!(
store
.has_publication(&digest("rabs.action-key.sha256.v1", 7))
.unwrap()
);
let snapshot = store.differential_snapshot().unwrap();
assert!(snapshot.iter().any(|l| l.starts_with("pins|")
&& l.contains("action-publication")
&& l.contains("coordinator")));
assert!(
snapshot
.iter()
.any(|l| l.starts_with("action_evidence_index|"))
);
}
#[test]
fn h011_fence_refusals_write_nothing() {
let engine = RusqliteEngine::open_in_memory().unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
assert_eq!(
process_offer(
&mut store,
&offer(),
&expected_descriptor(),
no_committed,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
),
Err(OfferRefusal::NotActiveAuthority)
);
ready_store(&mut store);
let mut tampered = offer();
tampered
.authority
.action_generation
.created_under_authority_digest = digest(AUTHORITY_DIGEST_DOMAIN, 99);
assert_eq!(
process_offer(
&mut store,
&tampered,
&expected_descriptor(),
no_committed,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
),
Err(OfferRefusal::GenerationAuthorityMismatch)
);
let mut tampered = offer();
tampered.authority.action_generation.generation_id = ActionGenerationId(999);
assert_eq!(
process_offer(
&mut store,
&tampered,
&expected_descriptor(),
no_committed,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
),
Err(OfferRefusal::UnknownGeneration)
);
let mut tampered = offer();
tampered.authority.attempt_id = AttemptId(999);
assert_eq!(
process_offer(
&mut store,
&tampered,
&expected_descriptor(),
no_committed,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
),
Err(OfferRefusal::UnknownAttempt)
);
store.release_lease(30).unwrap();
assert_eq!(
process_offer(
&mut store,
&offer(),
&expected_descriptor(),
no_committed,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
),
Err(OfferRefusal::LeaseReleased)
);
assert!(
!store
.has_publication(&digest("rabs.action-key.sha256.v1", 7))
.unwrap()
);
}
#[test]
fn h011_expired_lease_refuses_publication_and_late_renewal_without_mutation() {
fn scenario(store: &mut dyn RabsMetadataStore) -> Vec<String> {
ready_store(store);
let prepared = offer();
let before = store.differential_snapshot().unwrap();
for now in [100, 101, u64::MAX] {
for durability in [
CommitDurabilityProfile::RequireDurableClosure,
CommitDurabilityProfile::AcceptVolatileLocations,
] {
assert_eq!(
process_offer(
store,
&prepared,
&expected_descriptor(),
no_committed,
900,
1,
durability,
|| now,
),
Err(OfferRefusal::LeaseExpired)
);
}
assert_eq!(
store.validate_attempt_lease(&prepared.authority, now),
Err(StoreError::LeaseExpired)
);
}
assert_eq!(
store.renew_attempt_lease(
&prepared.authority,
LeaseRenewal {
lease: prepared.authority.execution_lease_id,
seq: LeaseRenewalSeq(2),
},
200,
&|| 100,
),
Err(StoreError::LeaseExpired)
);
assert!(
!store
.has_publication(&prepared.authority.action_key)
.unwrap()
);
assert!(
store
.list_evidence_keys(&prepared.authority.action_key)
.unwrap()
.is_empty()
);
for tag in [40, 41, 50, 51, 60, 61, 62, 63] {
assert!(store.object_located(&object(tag).0).unwrap());
}
assert!(
store
.object_located(&prepared.manifest.artifact_bundle_root.as_ref().unwrap().0)
.unwrap()
);
let after = store.differential_snapshot().unwrap();
assert_eq!(after, before, "expiry must preserve uploaded candidates");
after
}
let mut reference =
SqlMetadataStore::open(RusqliteEngine::open(&fresh_path("expiry-ref")).unwrap())
.unwrap();
let mut candidate =
SqlMetadataStore::open(FsqliteEngine::open(&fresh_path("expiry-fsq")).unwrap())
.unwrap();
assert_eq!(scenario(&mut reference), scenario(&mut candidate));
}
#[test]
fn h011_lease_expiring_during_admission_cannot_commit() {
fn scenario(store: &mut dyn RabsMetadataStore) -> Vec<String> {
ready_store(store);
let prepared = offer();
let before = store.differential_snapshot().unwrap();
let readings = std::cell::Cell::new(0);
let clock = || {
let reading = readings.get();
readings.set(reading + 1);
if reading == 0 { 99 } else { 100 }
};
assert_eq!(
process_offer(
store,
&prepared,
&expected_descriptor(),
no_committed,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
clock,
),
Err(OfferRefusal::LeaseExpired)
);
assert!(
!store
.has_publication(&prepared.authority.action_key)
.unwrap()
);
assert!(
store
.list_evidence_keys(&prepared.authority.action_key)
.unwrap()
.is_empty()
);
assert!(store.pin_row(900).unwrap().is_none());
assert!(store.object_located(&prepared.manifest_id.0).unwrap());
assert!(store.object_located(&prepared.evidence_id.0).unwrap());
let after = store.differential_snapshot().unwrap();
assert_eq!(after, before);
after
}
let mut reference =
SqlMetadataStore::open(RusqliteEngine::open(&fresh_path("late-expiry-ref")).unwrap())
.unwrap();
let mut candidate =
SqlMetadataStore::open(FsqliteEngine::open(&fresh_path("late-expiry-fsq")).unwrap())
.unwrap();
assert_eq!(scenario(&mut reference), scenario(&mut candidate));
}
#[test]
fn h011_publication_resamples_clock_and_rolls_back_transaction_time_expiry() {
use crate::metadata_store::{SqlEngine, SqlValue};
use std::cell::Cell;
struct ExpireDuringTransaction<'a, E> {
inner: E,
clock: &'a Cell<u64>,
armed: bool,
trigger: &'static str,
}
impl<E: SqlEngine> SqlEngine for ExpireDuringTransaction<'_, E> {
fn execute(&mut self, sql: &str, params: &[SqlValue]) -> Result<usize, StoreError> {
let affected = self.inner.execute(sql, params)?;
if self.armed && sql.starts_with(self.trigger) {
self.clock.set(100);
self.armed = false;
}
Ok(affected)
}
fn query(
&mut self,
sql: &str,
params: &[SqlValue],
) -> Result<Vec<Vec<SqlValue>>, StoreError> {
self.inner.query(sql, params)
}
}
fn scenario(engine: impl SqlEngine, trigger: &'static str) -> Vec<String> {
let clock = Cell::new(99);
let engine = ExpireDuringTransaction {
inner: engine,
clock: &clock,
armed: false,
trigger,
};
let mut store = SqlMetadataStore::open(engine).unwrap();
ready_store(&mut store);
let prepared = offer();
let row = PublicationRow {
action_key: prepared.authority.action_key.clone(),
descriptor_digest: expected_descriptor(),
manifest_digest: prepared.manifest_id.0.clone(),
evidence_digest: prepared.evidence_id.0.clone(),
winner_generation: prepared.authority.action_generation.generation_id.0,
winner_attempt: prepared.authority.attempt_id.0,
result_kind: ResultKindTag::Success,
pin_id: 900,
pin_owner: "coordinator".to_owned(),
provisional_ancestors: Vec::new(),
};
let before = store.differential_snapshot().unwrap();
store
.validate_attempt_lease(&prepared.authority, clock.get())
.unwrap();
store.engine_mut().armed = true;
assert_eq!(
store.commit_publication(
&authority_digest(&prepared.authority.coordinator),
Some((&prepared.authority, &|| clock.get())),
&row,
),
Err(StoreError::LeaseExpired)
);
assert_eq!(clock.get(), 100);
assert!(!store.has_publication(&row.action_key).unwrap());
assert!(store.pin_row(row.pin_id).unwrap().is_none());
assert!(
store
.list_evidence_keys(&row.action_key)
.unwrap()
.is_empty()
);
let after = store.differential_snapshot().unwrap();
assert_eq!(after, before);
after
}
for trigger in ["BEGIN", "INSERT INTO action_publications"] {
assert_eq!(
scenario(
RusqliteEngine::open(&fresh_path("txn-expiry-ref")).unwrap(),
trigger,
),
scenario(
FsqliteEngine::open(&fresh_path("txn-expiry-fsq")).unwrap(),
trigger,
),
);
}
}
#[test]
fn h011_same_key_updates_expiring_inside_transactions_preserve_existing_state() {
use crate::metadata_store::{SqlEngine, SqlValue};
use std::cell::Cell;
#[derive(Clone, Copy)]
enum SameKeyUpdate {
AppendEvidence,
SemanticQuarantine,
ObservableAlreadyQuarantined,
}
struct ExpireDuringTransaction<'a, E> {
inner: E,
clock: &'a Cell<u64>,
armed: bool,
trigger: &'static str,
}
impl<E: SqlEngine> SqlEngine for ExpireDuringTransaction<'_, E> {
fn execute(&mut self, sql: &str, params: &[SqlValue]) -> Result<usize, StoreError> {
let affected = self.inner.execute(sql, params)?;
if self.armed && sql.starts_with(self.trigger) {
self.clock.set(100);
self.armed = false;
}
Ok(affected)
}
fn query(
&mut self,
sql: &str,
params: &[SqlValue],
) -> Result<Vec<Vec<SqlValue>>, StoreError> {
self.inner.query(sql, params)
}
}
fn scenario(
engine: impl SqlEngine,
update: SameKeyUpdate,
trigger: &'static str,
) -> Vec<String> {
let clock = Cell::new(99);
let engine = ExpireDuringTransaction {
inner: engine,
clock: &clock,
armed: false,
trigger,
};
let mut store = SqlMetadataStore::open(engine).unwrap();
ready_store(&mut store);
let winner = offer();
assert!(matches!(
process_offer(
&mut store,
&winner,
&expected_descriptor(),
no_committed,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| clock.get(),
)
.unwrap(),
PublicationOutcome::Committed(_)
));
let action = winner.authority.action_key.clone();
let action_key = digest_key(&action);
if matches!(update, SameKeyUpdate::ObservableAlreadyQuarantined) {
store
.set_serving_disposition_key(&action_key, DISPOSITION_QUARANTINED)
.unwrap();
}
let candidate = match update {
SameKeyUpdate::AppendEvidence => divergent_offer(
&mut store,
None,
50,
55,
digest("rabs.observation-stream.sha256.v1", 9),
),
SameKeyUpdate::SemanticQuarantine => divergent_offer(
&mut store,
Some(42),
52,
55,
digest("rabs.observation-stream.sha256.v1", 9),
),
SameKeyUpdate::ObservableAlreadyQuarantined => divergent_offer(
&mut store,
None,
52,
55,
digest("rabs.observation-stream.sha256.v1", 10),
),
};
let evidence_before = store.list_evidence_keys(&action).unwrap();
assert!(!evidence_before.contains(&digest_key(&candidate.evidence_id.0)));
let serving_before = store.serving_disposition_key(&action_key).unwrap();
let expected_disposition =
if matches!(update, SameKeyUpdate::ObservableAlreadyQuarantined) {
DISPOSITION_QUARANTINED
} else {
"servable"
};
assert_eq!(serving_before.as_deref(), Some(expected_disposition));
let before = store.differential_snapshot().unwrap();
store.engine_mut().armed = true;
let committed = winner.manifest.clone();
assert_eq!(
process_offer(
&mut store,
&candidate,
&expected_descriptor(),
move |_| Some(committed.clone()),
901,
2,
CommitDurabilityProfile::RequireDurableClosure,
|| clock.get(),
),
Err(OfferRefusal::LeaseExpired)
);
assert_eq!(clock.get(), 100, "expiry boundary must be exercised");
assert_eq!(
store.published_manifest_key(&action).unwrap(),
Some(digest_key(&winner.manifest_id.0))
);
assert_eq!(
store.serving_disposition_key(&action_key).unwrap(),
serving_before
);
assert_eq!(store.list_evidence_keys(&action).unwrap(), evidence_before);
assert!(store.pin_row(901).unwrap().is_none());
assert!(
store
.list_divergence_incidents(&action_key)
.unwrap()
.is_empty()
);
assert!(store.object_located(&candidate.manifest_id.0).unwrap());
assert!(store.object_located(&candidate.evidence_id.0).unwrap());
for output in &candidate.manifest.logical_outputs {
assert!(store.object_located(&output.object.0).unwrap());
}
assert!(
store
.object_located(&candidate.manifest.artifact_bundle_root.as_ref().unwrap().0)
.unwrap()
);
let after = store.differential_snapshot().unwrap();
assert_eq!(
after, before,
"expired same-key offers must leave no writes"
);
after
}
for (update, trigger) in [
(SameKeyUpdate::AppendEvidence, "BEGIN"),
(
SameKeyUpdate::AppendEvidence,
"INSERT OR IGNORE INTO action_evidence_index",
),
(SameKeyUpdate::SemanticQuarantine, "BEGIN"),
(
SameKeyUpdate::SemanticQuarantine,
"UPDATE action_serving_states",
),
(SameKeyUpdate::ObservableAlreadyQuarantined, "BEGIN"),
(
SameKeyUpdate::ObservableAlreadyQuarantined,
"UPDATE action_serving_states",
),
] {
assert_eq!(
scenario(
RusqliteEngine::open(&fresh_path("same-key-expiry-ref")).unwrap(),
update,
trigger,
),
scenario(
FsqliteEngine::open(&fresh_path("same-key-expiry-fsq")).unwrap(),
update,
trigger,
),
);
}
}
#[test]
fn h011_live_lease_commits_just_before_its_monotonic_deadline() {
fn scenario(store: &mut dyn RabsMetadataStore) -> Vec<String> {
ready_store(store);
let prepared = offer();
let outcome = process_offer(
store,
&prepared,
&expected_descriptor(),
no_committed,
900,
10_000,
CommitDurabilityProfile::RequireDurableClosure,
|| 99,
)
.unwrap();
let PublicationOutcome::Committed(receipt) = outcome else {
panic!("expected commit before the lease deadline, got {outcome:?}");
};
assert_eq!(receipt.committed_causal_sequence, 10_000);
assert!(
store
.has_publication(&prepared.authority.action_key)
.unwrap()
);
store.differential_snapshot().unwrap()
}
let mut reference =
SqlMetadataStore::open(RusqliteEngine::open(&fresh_path("live-lease-ref")).unwrap())
.unwrap();
let mut candidate =
SqlMetadataStore::open(FsqliteEngine::open(&fresh_path("live-lease-fsq")).unwrap())
.unwrap();
assert_eq!(scenario(&mut reference), scenario(&mut candidate));
}
#[test]
fn h011_renewed_lease_survives_reopen_and_allows_publication_until_new_deadline() {
fn prepare(store: &mut dyn RabsMetadataStore) -> Vec<String> {
ready_store(store);
let authority = attempt_authority();
store
.renew_attempt_lease(
&authority,
LeaseRenewal {
lease: authority.execution_lease_id,
seq: LeaseRenewalSeq(2),
},
200,
&|| 99,
)
.unwrap();
store.differential_snapshot().unwrap()
}
fn publish(store: &mut dyn RabsMetadataStore) -> Vec<String> {
let mut renewed = offer();
renewed.authority.lease_renewal_seq = LeaseRenewalSeq(2);
let state = store
.validate_attempt_lease(&renewed.authority, 150)
.unwrap();
assert_eq!(state.renewal_seq, 2);
assert_eq!(state.expires_at_own_monotonic_ms, 200);
assert!(matches!(
process_offer(
store,
&renewed,
&expected_descriptor(),
no_committed,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 150,
)
.unwrap(),
PublicationOutcome::Committed(_)
));
assert!(
store
.has_publication(&renewed.authority.action_key)
.unwrap()
);
let committed = store.differential_snapshot().unwrap();
assert_eq!(
process_offer(
store,
&renewed,
&expected_descriptor(),
no_committed,
901,
2,
CommitDurabilityProfile::RequireDurableClosure,
|| 200,
),
Err(OfferRefusal::LeaseExpired)
);
assert_eq!(store.differential_snapshot().unwrap(), committed);
committed
}
let reference_path = fresh_path("renewed-lease-ref");
let candidate_path = fresh_path("renewed-lease-fsq");
let prepared = {
let mut reference =
SqlMetadataStore::open(RusqliteEngine::open(&reference_path).unwrap()).unwrap();
let mut candidate =
SqlMetadataStore::open(FsqliteEngine::open(&candidate_path).unwrap()).unwrap();
let prepared = prepare(&mut reference);
assert_eq!(prepared, prepare(&mut candidate));
prepared
};
let mut reference =
SqlMetadataStore::open(RusqliteEngine::open(&reference_path).unwrap()).unwrap();
let mut candidate =
SqlMetadataStore::open(FsqliteEngine::open(&candidate_path).unwrap()).unwrap();
assert_eq!(reference.differential_snapshot().unwrap(), prepared);
assert_eq!(candidate.differential_snapshot().unwrap(), prepared);
assert_eq!(publish(&mut reference), publish(&mut candidate));
}
#[test]
fn h011_tombstoned_generation_refused() {
let engine = RusqliteEngine::open_in_memory().unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
ready_store(&mut store);
store.tombstone_generation(11).unwrap();
assert_eq!(
process_offer(
&mut store,
&offer(),
&expected_descriptor(),
no_committed,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
),
Err(OfferRefusal::GenerationTombstoned)
);
}
#[test]
fn h011_content_validation_refusals() {
let engine = RusqliteEngine::open_in_memory().unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
ready_store(&mut store);
assert_eq!(
process_offer(
&mut store,
&offer(),
&digest("rabs.descriptor.sha256.v1", 99),
no_committed,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
),
Err(OfferRefusal::DescriptorMismatch)
);
let mut tampered = offer();
tampered.manifest.key_epoch = 2;
tampered.manifest.semantic_result_digest = semantic_result_digest_v1(&tampered.manifest);
tampered.manifest.observable_result_digest =
observable_result_digest_v1(&tampered.manifest, &tampered.canonical_observations);
assert_eq!(
process_offer(
&mut store,
&tampered,
&expected_descriptor(),
no_committed,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
),
Err(OfferRefusal::EpochMismatch)
);
let mut tampered = offer();
tampered.manifest.semantic_result_digest = digest(SEMANTIC_PROJECTION_DOMAIN, 99);
assert_eq!(
process_offer(
&mut store,
&tampered,
&expected_descriptor(),
no_committed,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
),
Err(OfferRefusal::SemanticDigestMismatch)
);
let mut tampered = offer();
tampered.manifest.observable_result_digest = digest(OBSERVABLE_PROJECTION_DOMAIN, 99);
assert_eq!(
process_offer(
&mut store,
&tampered,
&expected_descriptor(),
no_committed,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
),
Err(OfferRefusal::ObservableDigestMismatch)
);
}
#[test]
fn h011_incomplete_object_closure_refused() {
let engine = RusqliteEngine::open_in_memory().unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
ready_store(&mut store);
let manifest_id = object(50);
let mut tampered_evidence = evidence(&manifest_id);
tampered_evidence.provenance_receipt = object(200);
let tampered = OfferPreparedActionResult::build(
attempt_authority(),
manifest(),
manifest_id,
tampered_evidence,
object(51),
digest("rabs.observation-stream.sha256.v1", 9),
&declared(),
Vec::new(),
)
.unwrap();
assert_eq!(
process_offer(
&mut store,
&tampered,
&expected_descriptor(),
no_committed,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
),
Err(OfferRefusal::IncompleteObjectClosure {
missing: digest_key(&object(200).0)
})
);
}
#[test]
fn h011_repeat_offer_is_idempotent_and_divergence_quarantines() {
let engine = RusqliteEngine::open_in_memory().unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
ready_store(&mut store);
let first = offer();
assert!(matches!(
process_offer(
&mut store,
&first,
&expected_descriptor(),
no_committed,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
)
.unwrap(),
PublicationOutcome::Committed(_)
));
assert_eq!(
process_offer(
&mut store,
&first,
&expected_descriptor(),
no_committed,
901,
2,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
)
.unwrap(),
PublicationOutcome::IdempotentEvidenceAppended
);
let mut divergent_manifest = manifest();
divergent_manifest.logical_outputs[0].object = object(42);
let divergent_id = object(52);
store.record_object(&object(42).0, 64).unwrap();
store
.add_location(&object(42).0, "/cas/42", Some(1), "raw", true)
.unwrap();
store.record_object(&divergent_id.0, 64).unwrap();
store
.add_location(&divergent_id.0, "/cas/52", Some(1), "raw", true)
.unwrap();
let divergent = OfferPreparedActionResult::build(
attempt_authority(),
divergent_manifest,
divergent_id.clone(),
evidence(&divergent_id),
object(51),
digest("rabs.observation-stream.sha256.v1", 9),
&declared(),
Vec::new(),
)
.unwrap();
locate_bundle_root(&mut store, &divergent);
let committed = first.manifest.clone();
let outcome = process_offer(
&mut store,
&divergent,
&expected_descriptor(),
move |_| Some(committed.clone()),
902,
3,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
)
.unwrap();
let PublicationOutcome::Quarantined(quarantine) = outcome else {
panic!("expected quarantine, got {outcome:?}");
};
assert_eq!(quarantine.class, DivergenceClass::SemanticDivergence);
assert_eq!(quarantine.incident_seq, 3);
assert_eq!(quarantine.candidate_pin_id, 902);
assert_eq!(
store
.published_manifest_key(&digest("rabs.action-key.sha256.v1", 7))
.unwrap()
.unwrap(),
digest_key(&object(50).0)
);
}
fn divergent_offer(
store: &mut dyn RabsMetadataStore,
output_tag: Option<u8>,
manifest_tag: u8,
evidence_tag: u8,
observations: TypedDigest,
) -> OfferPreparedActionResult {
let mut m = manifest();
if let Some(tag) = output_tag {
m.logical_outputs[0].object = object(tag);
}
let id = object(manifest_tag);
for tag in output_tag
.iter()
.copied()
.chain([manifest_tag, evidence_tag])
{
store.record_object(&object(tag).0, 64).unwrap();
store
.add_location(&object(tag).0, &format!("/cas/{tag}"), Some(1), "raw", true)
.unwrap();
}
let built = OfferPreparedActionResult::build(
attempt_authority(),
m,
id.clone(),
evidence(&id),
object(evidence_tag),
observations,
&declared(),
Vec::new(),
)
.unwrap();
locate_bundle_root(store, &built);
built
}
#[test]
fn h026_semantic_divergence_quarantines_preserves_and_escalates() {
let mut store = SqlMetadataStore::open(RusqliteEngine::open_in_memory().unwrap()).unwrap();
ready_store(&mut store);
let action = digest("rabs.action-key.sha256.v1", 7);
let action_key = digest_key(&action);
let auth = authority_digest(&coordinator_authority());
let first = offer();
assert!(matches!(
process_offer(
&mut store,
&first,
&expected_descriptor(),
no_committed,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
)
.unwrap(),
PublicationOutcome::Committed(_)
));
store
.record_served_consumer(&action_key, "consumer-b")
.unwrap();
store
.record_served_consumer(&action_key, "consumer-a")
.unwrap();
store
.append_trust_evaluation(
&auth,
&action,
&crate::metadata_store::TrustEvaluationRow {
version: 1,
state: "project-release-eligible".to_owned(),
reason: "release gate".to_owned(),
evaluated_seq: 2,
},
)
.unwrap();
let divergent = divergent_offer(
&mut store,
Some(42),
52,
55,
digest("rabs.observation-stream.sha256.v1", 9),
);
let committed = first.manifest.clone();
let outcome = process_offer(
&mut store,
&divergent,
&expected_descriptor(),
move |_| Some(committed.clone()),
902,
3,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
)
.unwrap();
let PublicationOutcome::Quarantined(quarantine) = outcome else {
panic!("expected quarantine, got {outcome:?}");
};
assert_eq!(quarantine.class, DivergenceClass::SemanticDivergence);
assert_eq!(
store.serving_disposition_key(&action_key).unwrap().unwrap(),
DISPOSITION_QUARANTINED
);
assert_eq!(
store.published_manifest_key(&action).unwrap().unwrap(),
digest_key(&object(50).0)
);
let pin = store.pin_row(902).unwrap().unwrap();
assert_eq!(pin.class, DIVERGENCE_EVIDENCE_PIN_CLASS);
assert_eq!(pin.root_key, digest_key(&object(52).0));
assert!(!pin.released);
let snapshot = store.differential_snapshot().unwrap();
assert!(
snapshot
.iter()
.any(|l| l.starts_with("object_edges|") && l.contains("divergence-candidate"))
);
assert!(
store
.list_evidence_keys(&action)
.unwrap()
.contains(&digest_key(&object(55).0))
);
let candidate_view = store
.list_evidence_keys_for_manifest(&digest_key(&object(52).0))
.unwrap();
assert!(candidate_view.contains(&digest_key(&object(55).0)));
let committed_view = store
.list_evidence_keys_for_manifest(&digest_key(&object(50).0))
.unwrap();
assert!(!committed_view.contains(&digest_key(&object(55).0)));
let incidents = store.list_divergence_incidents(&action_key).unwrap();
assert_eq!(incidents.len(), 1);
assert_eq!(incidents[0].seq, 3);
assert_eq!(incidents[0].class, "semantic");
assert_eq!(
incidents[0].committed_manifest_key,
digest_key(&object(50).0)
);
assert_eq!(
incidents[0].candidate_manifest_key,
digest_key(&object(52).0)
);
assert_eq!(incidents[0].candidate_pin_hex, format!("{:032x}", 902));
assert_eq!(
quarantine.escalations,
vec![
ConsumerEscalation {
consumer: "consumer-a".to_owned(),
trust_state: "project-release-eligible".to_owned(),
decision: "recall-and-reverify".to_owned(),
},
ConsumerEscalation {
consumer: "consumer-b".to_owned(),
trust_state: "project-release-eligible".to_owned(),
decision: "recall-and-reverify".to_owned(),
},
]
);
assert!(snapshot.iter().any(|l| {
l.starts_with(
"decision_receipts|divergence-escalation|consumer-a|3|recall-and-reverify",
)
}));
}
#[test]
fn t025_observable_only_divergence_has_presentation_not_action_quarantine() {
let mut store = SqlMetadataStore::open(RusqliteEngine::open_in_memory().unwrap()).unwrap();
ready_store(&mut store);
let action = digest("rabs.action-key.sha256.v1", 7);
let action_key = digest_key(&action);
let first = offer();
assert!(matches!(
process_offer(
&mut store,
&first,
&expected_descriptor(),
no_committed,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
)
.unwrap(),
PublicationOutcome::Committed(_)
));
store
.record_served_consumer(&action_key, "consumer-a")
.unwrap();
let divergent = divergent_offer(
&mut store,
None,
53,
56,
digest("rabs.observation-stream.sha256.v1", 10),
);
assert_eq!(
divergent.manifest.semantic_result_digest,
first.manifest.semantic_result_digest
);
let committed = first.manifest.clone();
let outcome = process_offer(
&mut store,
&divergent,
&expected_descriptor(),
move |_| Some(committed.clone()),
903,
4,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
)
.unwrap();
let PublicationOutcome::Quarantined(quarantine) = outcome else {
panic!("expected quarantine, got {outcome:?}");
};
assert_eq!(quarantine.class, DivergenceClass::ObservableOnlyDivergence);
assert_eq!(
store.serving_disposition_key(&action_key).unwrap().unwrap(),
DISPOSITION_PRESENTATION_QUARANTINED
);
assert!(store.pin_row(903).unwrap().is_some());
let incidents = store.list_divergence_incidents(&action_key).unwrap();
assert_eq!(incidents.len(), 1);
assert_eq!(incidents[0].class, "observable-only");
assert!(quarantine.escalations.is_empty());
let snapshot = store.differential_snapshot().unwrap();
assert!(
!snapshot
.iter()
.any(|line| { line.starts_with("decision_receipts|divergence-escalation|") })
);
assert!(
!snapshot.iter().any(|line| {
line.starts_with("quarantines|action-entry|") && line.contains(&action_key)
}),
"observable-only divergence must not quarantine the canonical action entry"
);
store
.set_serving_disposition_key(&action_key, DISPOSITION_QUARANTINED)
.unwrap();
store
.add_quarantine(
QuarantineScope::ActionEntry,
&action_key,
"preexisting-full-quarantine",
)
.unwrap();
let later_observable = divergent_offer(
&mut store,
None,
57,
58,
digest("rabs.observation-stream.sha256.v1", 11),
);
let committed = first.manifest.clone();
let later_outcome = process_offer(
&mut store,
&later_observable,
&expected_descriptor(),
move |_| Some(committed.clone()),
905,
6,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
)
.unwrap();
assert!(matches!(
later_outcome,
PublicationOutcome::Quarantined(DivergenceQuarantine {
class: DivergenceClass::ObservableOnlyDivergence,
..
})
));
assert_eq!(
store.serving_disposition_key(&action_key).unwrap().unwrap(),
DISPOSITION_QUARANTINED
);
assert!(store.differential_snapshot().unwrap().iter().any(|line| {
line.starts_with("quarantines|action-entry|")
&& line.contains("preexisting-full-quarantine")
}));
}
#[test]
fn t025_distinct_worker_attempts_with_different_evidence_share_canonical_result() {
let mut store = SqlMetadataStore::open(RusqliteEngine::open_in_memory().unwrap()).unwrap();
ready_store(&mut store);
let action = digest("rabs.action-key.sha256.v1", 7);
let action_key = digest_key(&action);
let first_authority = attempt_authority();
let first_timing_payload = b"worker=worker-a;started=1;elapsed_us=3100;cpu_us=2200";
let first_resource_payload = b"worker=worker-a;cores=2;peak_rss_bytes=1048576";
let first_timing = content_object(first_timing_payload);
let first_resources = content_object(first_resource_payload);
let first_placeholder_id = object(50);
let mut first_placeholder_evidence = evidence(&first_placeholder_id);
first_placeholder_evidence.execution_snapshot_root = first_resources.clone();
first_placeholder_evidence.raw_process_and_event_evidence = first_timing.clone();
let first_stamped = OfferPreparedActionResult::build(
first_authority.clone(),
manifest(),
first_placeholder_id,
first_placeholder_evidence,
object(51),
digest("rabs.observation-stream.sha256.v1", 9),
&declared(),
Vec::new(),
)
.unwrap();
let first_manifest_bytes =
crate::manifest_codec::encode_manifest_v1(&first_stamped.manifest);
let first_manifest_id = content_object(&first_manifest_bytes);
let mut first_evidence = evidence(&first_manifest_id);
first_evidence.execution_snapshot_root = first_resources.clone();
first_evidence.raw_process_and_event_evidence = first_timing.clone();
let first = OfferPreparedActionResult::build(
first_authority,
manifest(),
first_manifest_id.clone(),
first_evidence,
object(51),
digest("rabs.observation-stream.sha256.v1", 9),
&declared(),
Vec::new(),
)
.unwrap();
assert_eq!(
crate::manifest_codec::encode_manifest_v1(&first.manifest),
first_manifest_bytes
);
locate_test_object(
&mut store,
&first_manifest_id,
"/cas/manifest-worker-a",
first_manifest_bytes.len(),
);
locate_test_object(
&mut store,
&first_timing,
"/cas/timing-worker-a",
first_timing_payload.len(),
);
locate_test_object(
&mut store,
&first_resources,
"/cas/resources-worker-a",
first_resource_payload.len(),
);
let first_outcome = process_offer(
&mut store,
&first,
&expected_descriptor(),
no_committed,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
)
.unwrap();
let PublicationOutcome::Committed(first_record) = first_outcome else {
panic!("expected first commit, got {first_outcome:?}");
};
assert_eq!(first_record.canonical_result_manifest_id, first_manifest_id);
let second_authority = distinct_attempt_authority();
let auth = authority_digest(&second_authority.coordinator);
store
.admit_worker_session(
&auth,
&WorkerSessionOffer {
worker_peer_id: second_authority.worker_peer_id.clone(),
boot_generation: second_authority.worker_boot_generation,
incarnation: second_authority.worker_incarnation_id,
reenrollment_proof: None,
},
10,
)
.unwrap();
store
.admit_attempt_lease(&second_authority, 17, 200)
.unwrap();
let second_timing_payload = b"worker=worker-b;started=10;elapsed_us=8700;cpu_us=6900";
let second_resource_payload = b"worker=worker-b;cores=6;peak_rss_bytes=3145728";
let second_timing = content_object(second_timing_payload);
let second_resources = content_object(second_resource_payload);
let second_placeholder_id = object(52);
let mut second_placeholder_evidence = evidence(&second_placeholder_id);
second_placeholder_evidence.execution_snapshot_root = second_resources.clone();
second_placeholder_evidence.raw_process_and_event_evidence = second_timing.clone();
let second_stamped = OfferPreparedActionResult::build(
second_authority.clone(),
manifest(),
second_placeholder_id,
second_placeholder_evidence,
object(54),
digest("rabs.observation-stream.sha256.v1", 9),
&declared(),
Vec::new(),
)
.unwrap();
let second_manifest_bytes =
crate::manifest_codec::encode_manifest_v1(&second_stamped.manifest);
let second_manifest_id = content_object(&second_manifest_bytes);
let mut second_evidence = evidence(&second_manifest_id);
second_evidence.execution_snapshot_root = second_resources.clone();
second_evidence.raw_process_and_event_evidence = second_timing.clone();
let reoffer = OfferPreparedActionResult::build(
second_authority.clone(),
manifest(),
second_manifest_id.clone(),
second_evidence,
object(54),
digest("rabs.observation-stream.sha256.v1", 9),
&declared(),
Vec::new(),
)
.unwrap();
assert_ne!(
first.authority.worker_peer_id,
reoffer.authority.worker_peer_id
);
assert_ne!(first.authority.attempt_id, reoffer.authority.attempt_id);
assert_ne!(
first.authority.execution_lease_id,
reoffer.authority.execution_lease_id
);
assert_ne!(
first.evidence.execution_snapshot_root,
reoffer.evidence.execution_snapshot_root
);
assert_ne!(
first.evidence.raw_process_and_event_evidence,
reoffer.evidence.raw_process_and_event_evidence
);
assert_eq!(first_manifest_bytes, second_manifest_bytes);
assert_eq!(first_manifest_id, second_manifest_id);
assert_eq!(first.manifest_id, reoffer.manifest_id);
assert_eq!(first.manifest, reoffer.manifest);
locate_test_object(
&mut store,
&second_manifest_id,
"/cas/manifest-worker-b",
second_manifest_bytes.len(),
);
locate_test_object(
&mut store,
&second_timing,
"/cas/timing-worker-b",
second_timing_payload.len(),
);
locate_test_object(
&mut store,
&second_resources,
"/cas/resources-worker-b",
second_resource_payload.len(),
);
locate_test_object(&mut store, &object(54), "/cas/54", 64);
assert_eq!(
process_offer(
&mut store,
&reoffer,
&expected_descriptor(),
no_committed,
904,
18,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
)
.unwrap(),
PublicationOutcome::IdempotentEvidenceAppended
);
let evidence_keys = store.list_evidence_keys(&action).unwrap();
assert!(evidence_keys.contains(&digest_key(&object(51).0)));
assert!(evidence_keys.contains(&digest_key(&object(54).0)));
let manifest_view = store
.list_evidence_keys_for_manifest(&digest_key(&first_manifest_id.0))
.unwrap();
assert!(manifest_view.contains(&digest_key(&object(51).0)));
assert!(manifest_view.contains(&digest_key(&object(54).0)));
let snapshot = store.differential_snapshot().unwrap();
let second_attempt_hex = format!("{:032x}", second_authority.attempt_id.0);
let second_attempt_line = snapshot
.iter()
.find(|line| line.starts_with(&format!("action_attempts|{second_attempt_hex}|")))
.expect("second attempt row");
let attempt_fields: Vec<_> = second_attempt_line.split('|').collect();
assert_eq!(attempt_fields[4], "worker-b");
assert_eq!(attempt_fields[5], "17");
assert_eq!(
attempt_fields[8],
format!("{:032x}", second_authority.execution_lease_id.0)
);
let second_evidence_line = snapshot
.iter()
.find(|line| {
let fields: Vec<_> = line.split('|').collect();
fields.len() == 8
&& fields[0] == "action_evidence_index"
&& fields[6] == second_attempt_hex
})
.expect("second attempt evidence row");
let evidence_fields: Vec<_> = second_evidence_line.split('|').collect();
assert_eq!(evidence_fields[7], digest_key(&first_manifest_id.0));
assert_eq!(
store.serving_disposition_key(&action_key).unwrap().unwrap(),
"servable"
);
assert!(
store
.list_divergence_incidents(&action_key)
.unwrap()
.is_empty()
);
assert!(!snapshot.iter().any(|line| line.starts_with("quarantines|")));
}
#[test]
fn h026_incident_rows_are_append_only_and_authority_gated() {
let mut store = SqlMetadataStore::open(RusqliteEngine::open_in_memory().unwrap()).unwrap();
let auth = authority_digest(&coordinator_authority());
let row = DivergenceIncidentRow {
action_key: "k".to_owned(),
seq: 1,
class: "semantic".to_owned(),
committed_manifest_key: "m1".to_owned(),
candidate_manifest_key: "m2".to_owned(),
candidate_evidence_key: "e2".to_owned(),
candidate_pin_hex: format!("{:032x}", 7),
generation_hex: format!("{:032x}", 11),
attempt_hex: format!("{:032x}", 20),
detail: "detail".to_owned(),
};
assert_eq!(
store.record_divergence_incident(&auth, &row),
Err(StoreError::NotActiveAuthority)
);
ready_store(&mut store);
store.record_divergence_incident(&auth, &row).unwrap();
store.record_divergence_incident(&auth, &row).unwrap();
let mut tampered = row.clone();
tampered.detail = "rewritten".to_owned();
assert_eq!(
store.record_divergence_incident(&auth, &tampered),
Err(StoreError::AppendConflict("divergence_incidents".into()))
);
let incidents = store.list_divergence_incidents("k").unwrap();
assert_eq!(incidents.len(), 1);
assert_eq!(incidents[0], row);
}
#[test]
fn h026_differential_reference_vs_frankensqlite() {
fn scenario(store: &mut dyn RabsMetadataStore) -> Vec<String> {
ready_store(store);
let action = digest("rabs.action-key.sha256.v1", 7);
let action_key = digest_key(&action);
let first = offer();
assert!(matches!(
process_offer(
store,
&first,
&expected_descriptor(),
no_committed,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
)
.unwrap(),
PublicationOutcome::Committed(_)
));
store
.record_served_consumer(&action_key, "consumer-a")
.unwrap();
let divergent = divergent_offer(
store,
Some(42),
52,
55,
digest("rabs.observation-stream.sha256.v1", 9),
);
let committed = first.manifest.clone();
let outcome = process_offer(
store,
&divergent,
&expected_descriptor(),
move |_| Some(committed.clone()),
902,
3,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
)
.unwrap();
assert!(matches!(outcome, PublicationOutcome::Quarantined(_)));
store.differential_snapshot().unwrap()
}
let mut reference =
SqlMetadataStore::open(RusqliteEngine::open(&fresh_path("ref26")).unwrap()).unwrap();
let mut candidate =
SqlMetadataStore::open(FsqliteEngine::open(&fresh_path("fsq26")).unwrap()).unwrap();
assert_eq!(scenario(&mut reference), scenario(&mut candidate));
}
#[test]
fn h011_differential_full_pipeline_reference_vs_frankensqlite() {
fn scenario(store: &mut dyn RabsMetadataStore) -> Vec<String> {
ready_store(store);
let first = offer();
assert!(matches!(
process_offer(
store,
&first,
&expected_descriptor(),
no_committed,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
)
.unwrap(),
PublicationOutcome::Committed(_)
));
assert_eq!(
process_offer(
store,
&first,
&expected_descriptor(),
no_committed,
901,
2,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
)
.unwrap(),
PublicationOutcome::IdempotentEvidenceAppended
);
assert_eq!(
process_offer(
store,
&first,
&digest("rabs.descriptor.sha256.v1", 99),
no_committed,
902,
3,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
),
Err(OfferRefusal::DescriptorMismatch)
);
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));
}
#[test]
fn h034_evict_then_recompute_divergence_is_detected() {
let mut store = SqlMetadataStore::open(RusqliteEngine::open_in_memory().unwrap()).unwrap();
ready_store(&mut store);
let first = offer();
assert!(matches!(
process_offer(
&mut store,
&first,
&expected_descriptor(),
no_committed,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
)
.unwrap(),
PublicationOutcome::Committed(_)
));
retain_eviction_tombstone(&mut store, &first.manifest, 50).unwrap();
for tag in [40u8, 41, 50] {
store
.remove_location_by_key(&digest_key(&object(tag).0), &format!("/cas/{tag}"))
.unwrap();
}
let action = digest("rabs.action-key.sha256.v1", 7);
let action_key = digest_key(&action);
assert_eq!(
check_recomputation_against_tombstone(
&mut store,
&action,
&first.manifest.semantic_result_digest,
&digest(OBSERVABLE_PROJECTION_DOMAIN, 99),
)
.unwrap(),
RecomputationCheck::Divergence(DivergenceClass::ObservableOnlyDivergence)
);
assert_eq!(
store.serving_disposition_key(&action_key).unwrap().unwrap(),
DISPOSITION_PRESENTATION_QUARANTINED
);
assert!(!store.differential_snapshot().unwrap().iter().any(|line| {
line.starts_with("quarantines|action-entry|") && line.contains(&action_key)
}));
let outcome = check_recomputation_against_tombstone(
&mut store,
&action,
&digest(SEMANTIC_PROJECTION_DOMAIN, 99),
&first.manifest.observable_result_digest,
)
.unwrap();
assert_eq!(
outcome,
RecomputationCheck::Divergence(DivergenceClass::SemanticDivergence)
);
assert_eq!(
store.serving_disposition_key(&action_key).unwrap().unwrap(),
DISPOSITION_QUARANTINED
);
assert!(store.differential_snapshot().unwrap().iter().any(|line| {
line.starts_with("quarantines|action-entry|")
&& line.contains("post-eviction recomputation divergence")
}));
assert!(
store.eviction_tombstone(&action).unwrap().is_some(),
"divergence keeps the tombstone as incident evidence"
);
assert_eq!(
check_recomputation_against_tombstone(
&mut store,
&action,
&first.manifest.semantic_result_digest,
&digest(OBSERVABLE_PROJECTION_DOMAIN, 98),
)
.unwrap(),
RecomputationCheck::Divergence(DivergenceClass::ObservableOnlyDivergence)
);
assert_eq!(
store.serving_disposition_key(&action_key).unwrap().unwrap(),
DISPOSITION_QUARANTINED
);
assert_eq!(
check_recomputation_against_tombstone(
&mut store,
&action,
&first.manifest.semantic_result_digest,
&first.manifest.observable_result_digest,
)
.unwrap(),
RecomputationCheck::Reproduced
);
assert!(store.eviction_tombstone(&action).unwrap().is_none());
assert_eq!(
check_recomputation_against_tombstone(
&mut store,
&action,
&first.manifest.semantic_result_digest,
&first.manifest.observable_result_digest,
)
.unwrap(),
RecomputationCheck::NoTombstone
);
}
#[test]
fn h034_differential_reference_vs_frankensqlite() {
fn scenario(store: &mut dyn RabsMetadataStore) -> Vec<String> {
ready_store(store);
let first = offer();
assert!(matches!(
process_offer(
store,
&first,
&expected_descriptor(),
no_committed,
900,
1,
CommitDurabilityProfile::RequireDurableClosure,
|| 10,
)
.unwrap(),
PublicationOutcome::Committed(_)
));
retain_eviction_tombstone(store, &first.manifest, 50).unwrap();
let action = digest("rabs.action-key.sha256.v1", 7);
assert_eq!(
check_recomputation_against_tombstone(
store,
&action,
&digest(SEMANTIC_PROJECTION_DOMAIN, 99),
&first.manifest.observable_result_digest,
)
.unwrap(),
RecomputationCheck::Divergence(DivergenceClass::SemanticDivergence)
);
store.differential_snapshot().unwrap()
}
let mut reference =
SqlMetadataStore::open(RusqliteEngine::open(&fresh_path("ref34")).unwrap()).unwrap();
let mut candidate =
SqlMetadataStore::open(FsqliteEngine::open(&fresh_path("fsq34")).unwrap()).unwrap();
assert_eq!(scenario(&mut reference), scenario(&mut candidate));
}
}