use crate::metadata_store::{
ProvisionalObligationInsert, ProvisionalPinInsert, RabsMetadataStore, StoreError, digest_key,
};
use crate::pin_leases::{ReleaseOutcome, Releaser, release_pin_scoped};
use crate::publication::{Framing, authority_digest, output_role_name_for_tag, output_role_tag};
use rabs_protocol::authority::CoordinatorAuthority;
use rabs_protocol::generation::{ActionGenerationId, AttemptId, ExecutionLeaseId};
use rabs_protocol::raw_bytes::RawBytes;
use rabs_protocol::reconnect::SubscriberId;
use rabs_protocol::result_identity::{ObjectId, OutputRole, TypedDigest};
pub const PROVISIONAL_PIN_DOMAIN: &str = "rabs.provisional-pin.sha256.v1";
pub const PROVISIONAL_PIN_CLASS: &str = "provisional-metadata";
const GRANT_DEPENDENT_ATTEMPT: &str = "dependent-attempt";
const GRANT_AWAITING_EDGE: &str = "awaiting-edge";
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ProvisionalIdentity {
pub authority: CoordinatorAuthority,
pub action_key: TypedDigest,
pub generation: ActionGenerationId,
pub attempt: AttemptId,
pub lease: ExecutionLeaseId,
pub role: OutputRole,
pub virtual_path: RawBytes,
}
impl ProvisionalIdentity {
#[must_use]
pub fn digest(&self) -> TypedDigest {
let mut framing = Framing::new(PROVISIONAL_PIN_DOMAIN);
let authority = authority_digest(&self.authority);
framing
.digest_field(&authority)
.digest_field(&self.action_key)
.field(&self.generation.0.to_be_bytes())
.field(&self.attempt.0.to_be_bytes())
.field(&self.lease.0.to_be_bytes())
.u64(output_role_tag(self.role))
.field(self.virtual_path.as_bytes());
framing.finish(PROVISIONAL_PIN_DOMAIN)
}
#[must_use]
pub fn pin_key(&self) -> String {
digest_key(&self.digest())
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ProvisionalReader {
DependentAttempt {
worker: String,
attempt: AttemptId,
},
AwaitingEdge {
subscriber: SubscriberId,
},
}
impl ProvisionalReader {
fn grant_kind(&self) -> &'static str {
match self {
Self::DependentAttempt { .. } => GRANT_DEPENDENT_ATTEMPT,
Self::AwaitingEdge { .. } => GRANT_AWAITING_EDGE,
}
}
fn grant_id(&self) -> String {
match self {
Self::DependentAttempt { worker, attempt } => format!("{worker}/{:032x}", attempt.0),
Self::AwaitingEdge { subscriber } => format!("edge/{:032x}", subscriber.0),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ProducerContracts {
pub toolchain: TypedDigest,
pub events: TypedDigest,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ProvisionalPinError {
Collision {
pin_key: String,
pinned: String,
offered: String,
},
UnknownPin {
pin_key: String,
},
Unauthorized {
pin_key: String,
},
Closed {
pin_key: String,
},
ProducerInvalidated {
pin_key: String,
reason: String,
},
NonMonotonicRenewal {
pin_key: String,
},
AdoptionMismatch {
pin_key: String,
pinned: String,
committed: String,
},
NotActiveAuthority,
UnresolvedConsumerDebt {
pin_key: String,
open_count: usize,
},
Store(StoreError),
SameWinningAttempt {
pin_key: String,
},
ForeignAction {
pin_key: String,
},
ContractMismatch {
pin_key: String,
},
AncestorLineageUnresolved {
pin_key: String,
ancestor_pin_key: String,
detail: String,
},
AncestorLineageDiverged {
pin_key: String,
ancestor_pin_key: String,
detail: String,
},
}
impl From<StoreError> for ProvisionalPinError {
fn from(value: StoreError) -> Self {
Self::Store(value)
}
}
impl std::fmt::Display for ProvisionalPinError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Collision {
pin_key,
pinned,
offered,
} => write!(
f,
"provisional pin collision on {pin_key}: pinned {pinned}, offered {offered}"
),
Self::UnknownPin { pin_key } => write!(f, "unknown provisional pin {pin_key}"),
Self::Unauthorized { pin_key } => {
write!(f, "reader not authorized for provisional pin {pin_key}")
}
Self::Closed { pin_key } => write!(f, "provisional pin {pin_key} is closed"),
Self::ProducerInvalidated { pin_key, reason } => {
write!(f, "producer lineage invalidated ({reason}): {pin_key}")
}
Self::NonMonotonicRenewal { pin_key } => {
write!(f, "non-monotonic renewal for provisional pin {pin_key}")
}
Self::AdoptionMismatch {
pin_key,
pinned,
committed,
} => write!(
f,
"adoption mismatch on {pin_key}: pinned {pinned}, committed {committed}"
),
Self::NotActiveAuthority => {
write!(f, "coordinator authority is not active for this release")
}
Self::UnresolvedConsumerDebt {
pin_key,
open_count,
} => write!(
f,
"provisional pin {pin_key} still has {open_count} open consumer obligation(s)"
),
Self::SameWinningAttempt { pin_key } => write!(
f,
"same producer attempt already owns commit resolution for {pin_key}"
),
Self::ForeignAction { pin_key } => write!(
f,
"winning attempt serves a different action than pinned producer of {pin_key}"
),
Self::ContractMismatch { pin_key } => write!(
f,
"winner toolchain/event contracts differ from pin binding on {pin_key}"
),
Self::AncestorLineageUnresolved {
pin_key,
ancestor_pin_key,
detail,
} => write!(
f,
"ancestor lineage unresolved for {pin_key} at {ancestor_pin_key}: {detail}"
),
Self::AncestorLineageDiverged {
pin_key,
ancestor_pin_key,
detail,
} => write!(
f,
"ancestor lineage diverged for {pin_key} at {ancestor_pin_key}: {detail}"
),
Self::Store(e) => write!(f, "store error: {e:?}"),
}
}
}
impl std::error::Error for ProvisionalPinError {}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum OpenOutcome {
Created,
AlreadyPinned,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CloseOutcome {
Released,
AlreadyReleased,
}
pub fn open_provisional_pin(
store: &mut dyn RabsMetadataStore,
identity: &ProvisionalIdentity,
object: &ObjectId,
contracts: &ProducerContracts,
) -> Result<OpenOutcome, ProvisionalPinError> {
let pin_key = identity.pin_key();
if let Some(existing) = store.provisional_pin_row(&pin_key)? {
let offered = digest_key(&object.0);
if existing.object_key == offered {
return if existing.released {
Err(ProvisionalPinError::Closed { pin_key })
} else {
Ok(OpenOutcome::AlreadyPinned)
};
}
return Err(ProvisionalPinError::Collision {
pin_key,
pinned: existing.object_key,
offered,
});
}
let mut best: std::collections::BTreeMap<String, u64> = std::collections::BTreeMap::new();
let mut frontier: Vec<(String, u64)> = store
.list_open_provisional_obligations_by_attempt(&format!("{:032x}", identity.attempt.0))?
.into_iter()
.map(|obligation| (obligation.pin_key, 1))
.collect();
while let Some((key, hops)) = frontier.pop() {
if best.get(&key).is_some_and(|&known| known <= hops) {
continue;
}
best.insert(key.clone(), hops);
for (ancestor_key, stored_hops) in store.list_provisional_pin_ancestors(&key)? {
frontier.push((ancestor_key, hops + stored_hops));
}
}
let ancestor_pin_keys: Vec<(String, u64)> = best.into_iter().collect();
let digest = identity.digest();
let mut be16 = [0u8; 16];
be16.copy_from_slice(&digest.bytes[0..16]);
let protective_pin_id = u128::from_be_bytes(be16);
store.insert_provisional_pin(&ProvisionalPinInsert {
pin_key: pin_key.clone(),
authority_key: digest_key(&authority_digest(&identity.authority)),
action_key: digest_key(&identity.action_key),
generation: identity.generation.0,
attempt: identity.attempt.0,
lease: identity.lease.0,
role_tag: i64::try_from(output_role_tag(identity.role)).expect("role tags fit i64"),
virtual_path: identity.virtual_path.as_bytes().to_vec(),
object: object.0.clone(),
protective_pin_id,
reason: format!(
"M004 candidate pin: action {}, gen {:032x}, attempt {:032x}",
digest_key(&identity.action_key),
identity.generation.0,
identity.attempt.0
),
toolchain_contract_key: digest_key(&contracts.toolchain),
event_contract_key: digest_key(&contracts.events),
ancestor_pin_keys,
})?;
Ok(OpenOutcome::Created)
}
pub fn authorize_reader(
store: &mut dyn RabsMetadataStore,
identity: &ProvisionalIdentity,
reader: &ProvisionalReader,
) -> Result<(), ProvisionalPinError> {
let pin_key = identity.pin_key();
let Some(row) = store.provisional_pin_row(&pin_key)? else {
return Err(ProvisionalPinError::UnknownPin { pin_key });
};
if row.released {
return Err(ProvisionalPinError::Closed { pin_key });
}
store.record_provisional_grant(
&pin_key,
reader.grant_kind(),
&reader.grant_id(),
row.renewal_seq,
)?;
Ok(())
}
pub fn resolve_for_reader(
store: &mut dyn RabsMetadataStore,
identity: &ProvisionalIdentity,
reader: &ProvisionalReader,
) -> Result<ObjectId, ProvisionalPinError> {
let pin_key = identity.pin_key();
let Some(row) = store.provisional_pin_row(&pin_key)? else {
return Err(ProvisionalPinError::UnknownPin { pin_key });
};
if let Some(reason) = row.invalidated_reason {
return Err(ProvisionalPinError::ProducerInvalidated { pin_key, reason });
}
if row.released {
return Err(ProvisionalPinError::Closed { pin_key });
}
let authorized = store
.list_provisional_grants(&pin_key)?
.iter()
.any(|(kind, id, _)| kind.as_str() == reader.grant_kind() && id == &reader.grant_id());
if !authorized {
return Err(ProvisionalPinError::Unauthorized { pin_key });
}
if let ProvisionalReader::DependentAttempt { worker, attempt } = reader {
store.record_provisional_consumption(&ProvisionalObligationInsert {
consumer_worker: worker.clone(),
consumer_attempt: attempt.0,
pin_key: pin_key.clone(),
producer_action_key: digest_key(&identity.action_key),
producer_generation: identity.generation.0,
producer_attempt: identity.attempt.0,
role_tag: i64::try_from(output_role_tag(identity.role)).expect("role tags fit i64"),
virtual_path: identity.virtual_path.as_bytes().to_vec(),
object_key: row.object_key.clone(),
created_seq: row.renewal_seq,
})?;
}
Ok(ObjectId(row.object))
}
pub fn record_adoption(
store: &mut dyn RabsMetadataStore,
identity: &ProvisionalIdentity,
committed_object: &ObjectId,
) -> Result<(), ProvisionalPinError> {
let pin_key = identity.pin_key();
let Some(row) = store.provisional_pin_row(&pin_key)? else {
return Err(ProvisionalPinError::UnknownPin { pin_key });
};
let committed = digest_key(&committed_object.0);
if row.object_key != committed {
return Err(ProvisionalPinError::AdoptionMismatch {
pin_key,
pinned: row.object_key,
committed,
});
}
verify_transitive_lineage(store, identity)?;
store.adopt_provisional_pin(&pin_key, &committed)?;
store.resolve_provisional_obligations(&pin_key, &committed)?;
Ok(())
}
pub fn resolve_consumers_on_commit(
store: &mut dyn RabsMetadataStore,
identity: &ProvisionalIdentity,
committed_object: &ObjectId,
) -> Result<usize, ProvisionalPinError> {
let pin_key = identity.pin_key();
let Some(row) = store.provisional_pin_row(&pin_key)? else {
return Err(ProvisionalPinError::UnknownPin { pin_key });
};
let committed = digest_key(&committed_object.0);
if row.object_key != committed {
return Err(ProvisionalPinError::AdoptionMismatch {
pin_key,
pinned: row.object_key.clone(),
committed: committed.clone(),
});
}
verify_transitive_lineage(store, identity)?;
store.adopt_provisional_pin(&pin_key, &committed)?;
Ok(store.resolve_provisional_obligations(&pin_key, &committed)?)
}
pub fn verify_transitive_lineage(
store: &mut dyn RabsMetadataStore,
identity: &ProvisionalIdentity,
) -> Result<(), ProvisionalPinError> {
let pin_key = identity.pin_key();
let Some(row) = store.provisional_pin_row(&pin_key)? else {
return Err(ProvisionalPinError::UnknownPin { pin_key });
};
if let Some(reason) = row.invalidated_reason {
return Err(ProvisionalPinError::ProducerInvalidated { pin_key, reason });
}
if row.released {
return Err(ProvisionalPinError::Closed { pin_key });
}
if let Some(obligation) = store
.list_open_provisional_obligations_by_attempt(&format!("{:032x}", identity.attempt.0))?
.into_iter()
.next()
{
if obligation.status == "cancelled" {
return Err(ProvisionalPinError::AncestorLineageDiverged {
pin_key,
ancestor_pin_key: obligation.pin_key,
detail: format!(
"producing attempt consumed cancelled object {}",
obligation.object_key
),
});
}
return Err(ProvisionalPinError::AncestorLineageUnresolved {
pin_key: pin_key.clone(),
ancestor_pin_key: obligation.pin_key.clone(),
detail: format!(
"inbound obligation status {} on consumed object {}",
obligation.status, obligation.object_key
),
});
}
for (ancestor_pin_key, _min_hops) in store.list_provisional_pin_ancestors(&pin_key)? {
let Some(ancestor) = store.provisional_pin_row(&ancestor_pin_key)? else {
return Err(ProvisionalPinError::AncestorLineageUnresolved {
pin_key: pin_key.clone(),
ancestor_pin_key,
detail: "closure member missing from registry".to_owned(),
});
};
if let Some(reason) = ancestor.invalidated_reason {
return Err(ProvisionalPinError::AncestorLineageDiverged {
pin_key: pin_key.clone(),
ancestor_pin_key,
detail: format!("invalidated: {reason}"),
});
}
match &ancestor.adopted_object_key {
Some(adopted) if *adopted == ancestor.object_key => {}
Some(adopted) => {
return Err(ProvisionalPinError::AncestorLineageDiverged {
pin_key: pin_key.clone(),
ancestor_pin_key,
detail: format!("adopted foreign bytes {adopted}"),
});
}
None if ancestor.released => {}
None => {
return Err(ProvisionalPinError::AncestorLineageUnresolved {
pin_key: pin_key.clone(),
ancestor_pin_key,
detail: "ancestor not yet adopted".to_owned(),
});
}
}
}
Ok(())
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct WinningAttemptContext {
pub authority: CoordinatorAuthority,
pub action_key: TypedDigest,
pub generation: ActionGenerationId,
pub attempt: AttemptId,
pub contracts: ProducerContracts,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum AdoptionOutcome {
Adopted {
obligations_resolved: usize,
},
DivergenceCancelled {
pins_invalidated: usize,
obligations_cancelled: usize,
},
}
pub fn adopt_from_winning_attempt(
store: &mut dyn RabsMetadataStore,
identity: &ProvisionalIdentity,
winner: &WinningAttemptContext,
committed_object: &ObjectId,
) -> Result<AdoptionOutcome, ProvisionalPinError> {
let pin_key = identity.pin_key();
let Some(row) = store.provisional_pin_row(&pin_key)? else {
return Err(ProvisionalPinError::UnknownPin { pin_key });
};
if let Some(reason) = row.invalidated_reason {
return Err(ProvisionalPinError::ProducerInvalidated { pin_key, reason });
}
if row.released {
return Err(ProvisionalPinError::Closed { pin_key });
}
if winner.action_key != identity.action_key {
return Err(ProvisionalPinError::ForeignAction { pin_key });
}
if winner.attempt == identity.attempt && winner.generation == identity.generation {
return Err(ProvisionalPinError::SameWinningAttempt { pin_key });
}
if row.toolchain_contract_key.is_empty()
|| row.event_contract_key.is_empty()
|| row.toolchain_contract_key != digest_key(&winner.contracts.toolchain)
|| row.event_contract_key != digest_key(&winner.contracts.events)
{
return Err(ProvisionalPinError::ContractMismatch { pin_key });
}
let committed = digest_key(&committed_object.0);
if committed != row.object_key {
let reason = format!(
"winning attempt {:032x}/{:032x} committed divergent object {} \
for logical output of {}",
winner.generation.0, winner.attempt.0, committed, row.object_key
);
let counts = close_and_cancel_cascading(store, &pin_key, &reason)?;
return Ok(AdoptionOutcome::DivergenceCancelled {
pins_invalidated: counts.pins,
obligations_cancelled: counts.obligations,
});
}
verify_transitive_lineage(store, identity)?;
store.record_adoption_edge(
&authority_digest(&winner.authority),
&digest_key(&identity.action_key),
output_role_name_for_tag(row.role_tag),
&row.virtual_path,
&row.object_key,
&committed,
)?;
store.adopt_provisional_pin(&pin_key, &committed)?;
let resolved = store.resolve_provisional_obligations(&pin_key, &committed)?;
Ok(AdoptionOutcome::Adopted {
obligations_resolved: resolved,
})
}
pub fn invalidate_lineage(
store: &mut dyn RabsMetadataStore,
identity: &ProvisionalIdentity,
reason: &str,
) -> Result<CloseOutcome, ProvisionalPinError> {
let pin_key = identity.pin_key();
if store.provisional_pin_row(&pin_key)?.is_none() {
return Err(ProvisionalPinError::UnknownPin { pin_key });
}
close_and_cancel_cascading(store, &pin_key, reason)?;
Ok(CloseOutcome::Released)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum TerminalGate {
Clear,
Blocked {
pending_pin_keys: Vec<String>,
},
Refused {
cancelled_pin_keys: Vec<String>,
},
}
pub fn descendant_terminal_gate(
store: &mut dyn RabsMetadataStore,
consumer_worker: &str,
consumer_attempt: AttemptId,
) -> Result<TerminalGate, ProvisionalPinError> {
let rows = store.list_open_provisional_obligations(
consumer_worker,
&format!("{:032x}", consumer_attempt.0),
)?;
let mut blocked = Vec::new();
let mut refused = Vec::new();
for row in rows {
match row.status.as_str() {
"open" => blocked.push(row.pin_key),
"cancelled" => refused.push(row.pin_key),
other => {
return Err(ProvisionalPinError::Store(StoreError::Corruption(format!(
"obligation status {other:?}"
))));
}
}
}
Ok(if !refused.is_empty() {
TerminalGate::Refused {
cancelled_pin_keys: refused,
}
} else if !blocked.is_empty() {
TerminalGate::Blocked {
pending_pin_keys: blocked,
}
} else {
TerminalGate::Clear
})
}
pub fn release_after_drain(
store: &mut dyn RabsMetadataStore,
identity: &ProvisionalIdentity,
releaser: &Releaser,
) -> Result<CloseOutcome, ProvisionalPinError> {
let pin_key = identity.pin_key();
let Some(row) = store.provisional_pin_row(&pin_key)? else {
return Err(ProvisionalPinError::UnknownPin { pin_key });
};
let presented = match releaser {
Releaser::Worker(_) => {
return Err(ProvisionalPinError::Unauthorized { pin_key });
}
Releaser::Coordinator(authority) => authority_digest(authority),
};
let active = store
.active_authority()?
.map(|authority_row| authority_row.digest);
if active.as_ref() != Some(&presented) {
return Err(ProvisionalPinError::NotActiveAuthority);
}
let open_debt = store.count_open_provisional_obligations(&pin_key)?;
if open_debt > 0 {
return Err(ProvisionalPinError::UnresolvedConsumerDebt {
pin_key,
open_count: open_debt,
});
}
let outcome = close_internal(store, identity, None)?;
let protective_id = u128::from_str_radix(&row.protective_pin_hex, 16).map_err(|_| {
ProvisionalPinError::Store(StoreError::Corruption(format!(
"protective pin hex {}",
row.protective_pin_hex
)))
})?;
match release_pin_scoped(store, protective_id, releaser)? {
ReleaseOutcome::Released | ReleaseOutcome::AlreadyReleased => {}
ReleaseOutcome::RefusedNotActiveAuthority => {
return Err(ProvisionalPinError::NotActiveAuthority);
}
other => {
return Err(ProvisionalPinError::Store(StoreError::Corruption(format!(
"protective pin release refused: {other:?}"
))));
}
}
Ok(outcome)
}
pub fn renew_provisional_pin(
store: &mut dyn RabsMetadataStore,
identity: &ProvisionalIdentity,
renewal_seq: u64,
) -> Result<(), ProvisionalPinError> {
let pin_key = identity.pin_key();
store
.renew_provisional_pin(&pin_key, renewal_seq)
.map_err(|e| match e {
StoreError::NonMonotonicPinRenewal => {
ProvisionalPinError::NonMonotonicRenewal { pin_key }
}
StoreError::PinReleased => ProvisionalPinError::Closed { pin_key },
StoreError::UnknownPin => ProvisionalPinError::UnknownPin { pin_key },
other => ProvisionalPinError::Store(other),
})
}
fn close_internal(
store: &mut dyn RabsMetadataStore,
identity: &ProvisionalIdentity,
invalidation_reason: Option<&str>,
) -> Result<CloseOutcome, ProvisionalPinError> {
let pin_key = identity.pin_key();
if store.provisional_pin_row(&pin_key)?.is_none() {
return Err(ProvisionalPinError::UnknownPin { pin_key });
}
store.close_provisional_pin(&pin_key, invalidation_reason)?;
Ok(CloseOutcome::Released)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct InvalidationSummary {
pub pins_invalidated: usize,
pub obligations_cancelled: usize,
}
struct CascadeCounts {
pins: usize,
obligations: usize,
}
fn close_and_cancel_cascading(
store: &mut dyn RabsMetadataStore,
root_pin_key: &str,
reason: &str,
) -> Result<CascadeCounts, ProvisionalPinError> {
let mut pending: Vec<String> = vec![root_pin_key.to_owned()];
pending.extend(store.list_provisional_pin_descendants(root_pin_key)?);
pending.sort();
pending.dedup();
let mut counts = CascadeCounts {
pins: 0,
obligations: 0,
};
for key in pending {
let Some(row) = store.provisional_pin_row(&key)? else {
continue;
};
if row.released {
continue;
}
store.close_provisional_pin(&key, Some(reason))?;
counts.pins += 1;
counts.obligations += store.cancel_provisional_obligations(&key)?;
}
Ok(counts)
}
pub fn invalidate_lineage_for_generation_failure(
store: &mut dyn RabsMetadataStore,
action_key: &TypedDigest,
generation: ActionGenerationId,
reason: &str,
) -> Result<InvalidationSummary, ProvisionalPinError> {
let rows = store.list_open_provisional_pins_for_action_generation(
&digest_key(action_key),
&format!("{:032x}", generation.0),
)?;
let mut summary = InvalidationSummary {
pins_invalidated: 0,
obligations_cancelled: 0,
};
for row in &rows {
let counts = close_and_cancel_cascading(store, &row.pin_key, reason)?;
summary.pins_invalidated += counts.pins;
summary.obligations_cancelled += counts.obligations;
}
Ok(summary)
}
pub fn invalidate_lineage_for_authority_loss(
store: &mut dyn RabsMetadataStore,
authority: &CoordinatorAuthority,
reason: &str,
) -> Result<InvalidationSummary, ProvisionalPinError> {
let rows = store
.list_open_provisional_pins_for_authority(&digest_key(&authority_digest(authority)))?;
let mut summary = InvalidationSummary {
pins_invalidated: 0,
obligations_cancelled: 0,
};
for row in &rows {
let counts = close_and_cancel_cascading(store, &row.pin_key, reason)?;
summary.pins_invalidated += counts.pins;
summary.obligations_cancelled += counts.obligations;
}
Ok(summary)
}
pub fn invalidate_unadopted_lineage_for_action(
store: &mut dyn RabsMetadataStore,
action_key: &TypedDigest,
committed_output_keys: &std::collections::BTreeSet<String>,
reason: &str,
) -> Result<InvalidationSummary, ProvisionalPinError> {
let rows = store.list_open_provisional_pins_for_action(&digest_key(action_key))?;
let mut summary = InvalidationSummary {
pins_invalidated: 0,
obligations_cancelled: 0,
};
for row in &rows {
if committed_output_keys.contains(&row.object_key) {
continue;
}
let counts = close_and_cancel_cascading(store, &row.pin_key, reason)?;
summary.pins_invalidated += counts.pins;
summary.obligations_cancelled += counts.obligations;
}
Ok(summary)
}
pub fn provisional_causal_trace(
store: &mut dyn RabsMetadataStore,
pin_key: &str,
) -> Result<Vec<crate::metadata_store::ProvisionalObligationRow>, ProvisionalPinError> {
Ok(store.list_provisional_obligations_for_pin(pin_key)?)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::metadata_store::{RusqliteEngine, SqlMetadataStore};
use rabs_protocol::authority::ClusterId;
use rabs_protocol::result_identity::DigestAlgorithm;
const ACTION_DOMAIN: &str = "rabs.action-key.sha256.v1";
struct Fixture {
store: SqlMetadataStore<RusqliteEngine>,
}
fn tagged_action(tag: u8) -> TypedDigest {
let mut bytes = [0u8; 32];
bytes[0] = tag;
bytes[31] = tag;
TypedDigest {
algorithm: DigestAlgorithm::Sha256V1,
domain: ACTION_DOMAIN,
bytes,
}
}
fn tagged_object(tag: u8) -> ObjectId {
let mut d = tagged_action(tag);
d.domain = "rabs.object.sha256.v1";
ObjectId(d)
}
fn authority(tag: u64) -> CoordinatorAuthority {
CoordinatorAuthority {
cluster_id: ClusterId(format!("cluster-{tag}")),
credential_generation: tag,
term: 100 + tag,
incarnation_id: rabs_protocol::authority::CoordinatorIncarnationId(
0xAA00_0000_0000_0000 + u128::from(tag),
),
}
}
fn identity(authority_tag: u64, attempt_tag: u128) -> ProvisionalIdentity {
ProvisionalIdentity {
authority: authority(authority_tag),
action_key: tagged_action(10),
generation: ActionGenerationId(0x50),
attempt: AttemptId(attempt_tag),
lease: ExecutionLeaseId(attempt_tag + 1),
role: OutputRole::ProvisionalMetadata,
virtual_path: RawBytes::new(b"target/debug/deps/libfeat.rmeta".to_vec()),
}
}
fn fixture(name: &str) -> Fixture {
let engine = RusqliteEngine::open_in_memory().unwrap();
let store = SqlMetadataStore::open(engine).unwrap();
let _ = name;
Fixture { store }
}
fn dependent(worker: &str, attempt: u128) -> ProvisionalReader {
ProvisionalReader::DependentAttempt {
worker: worker.to_owned(),
attempt: AttemptId(attempt),
}
}
fn edge(subscriber: u128) -> ProvisionalReader {
ProvisionalReader::AwaitingEdge {
subscriber: SubscriberId(subscriber),
}
}
#[test]
fn m004_open_is_idempotent_and_collision_refuses_different_object() {
let mut f = fixture("open-idem");
let id = identity(1, 20);
let obj_a = tagged_object(41);
assert_eq!(
open_provisional_pin(&mut f.store, &id, &obj_a, &m017_ctx()).unwrap(),
OpenOutcome::Created
);
assert_eq!(
open_provisional_pin(&mut f.store, &id, &obj_a, &m017_ctx()).unwrap(),
OpenOutcome::AlreadyPinned
);
let err =
open_provisional_pin(&mut f.store, &id, &tagged_object(42), &m017_ctx()).unwrap_err();
let ProvisionalPinError::Collision {
pinned, offered, ..
} = &err
else {
panic!("expected collision, got {err:?}");
};
assert_eq!(pinned, &digest_key(&obj_a.0));
assert_eq!(offered, &digest_key(&tagged_object(42).0));
authorize_reader(&mut f.store, &id, &edge(7)).unwrap();
assert_eq!(
resolve_for_reader(&mut f.store, &id, &edge(7)).unwrap(),
obj_a
);
}
#[test]
fn m004_visibility_is_closed_by_default_and_grant_scoped() {
let mut f = fixture("visibility");
let id = identity(1, 21);
let obj = tagged_object(43);
open_provisional_pin(&mut f.store, &id, &obj, &m017_ctx()).unwrap();
assert_eq!(
resolve_for_reader(&mut f.store, &id, &dependent("worker-a", 99)).unwrap_err(),
ProvisionalPinError::Unauthorized {
pin_key: id.pin_key()
}
);
authorize_reader(&mut f.store, &id, &dependent("worker-a", 99)).unwrap();
authorize_reader(&mut f.store, &id, &edge(7)).unwrap();
assert_eq!(
resolve_for_reader(&mut f.store, &id, &dependent("worker-a", 99)).unwrap(),
obj
);
assert_eq!(
resolve_for_reader(&mut f.store, &id, &edge(7)).unwrap(),
obj
);
assert_eq!(
resolve_for_reader(&mut f.store, &id, &dependent("worker-a", 100)).unwrap_err(),
ProvisionalPinError::Unauthorized {
pin_key: id.pin_key()
}
);
assert_eq!(
resolve_for_reader(&mut f.store, &id, &edge(8)).unwrap_err(),
ProvisionalPinError::Unauthorized {
pin_key: id.pin_key()
}
);
let other_generation = ProvisionalIdentity {
generation: ActionGenerationId(0x51),
..identity(1, 21)
};
assert_eq!(
resolve_for_reader(&mut f.store, &other_generation, &edge(7)).unwrap_err(),
ProvisionalPinError::UnknownPin {
pin_key: other_generation.pin_key()
}
);
let other_authority = ProvisionalIdentity {
authority: authority(2),
..identity(1, 21)
};
assert_ne!(other_authority.pin_key(), id.pin_key());
}
#[test]
fn m004_worker_cannot_release_and_active_coordinator_can() {
let mut f = fixture("release-auth");
let id = identity(1, 22);
open_provisional_pin(&mut f.store, &id, &tagged_object(44), &m017_ctx()).unwrap();
let err = release_after_drain(&mut f.store, &id, &Releaser::Worker("worker-a".to_owned()))
.unwrap_err();
assert_eq!(
err,
ProvisionalPinError::Unauthorized {
pin_key: id.pin_key()
}
);
f.store
.acquire_authority(&crate::metadata_store::AuthorityRow {
digest: crate::publication::authority_digest(&authority(9)),
cluster_id: "cluster-9".to_owned(),
incarnation: 0xEE,
term: 109,
acquired_seq: 1,
})
.unwrap();
let err = release_after_drain(&mut f.store, &id, &Releaser::Coordinator(authority(2)))
.unwrap_err();
assert_eq!(err, ProvisionalPinError::NotActiveAuthority);
authorize_reader(&mut f.store, &id, &edge(7)).unwrap();
f.store
.release_authority(&crate::publication::authority_digest(&authority(9)))
.unwrap();
f.store
.acquire_authority(&crate::metadata_store::AuthorityRow {
digest: crate::publication::authority_digest(&authority(1)),
cluster_id: "cluster-1".to_owned(),
incarnation: 0xAA,
term: 101,
acquired_seq: 2,
})
.unwrap();
assert_eq!(
release_after_drain(&mut f.store, &id, &Releaser::Coordinator(authority(1))).unwrap(),
CloseOutcome::Released
);
authorize_reader(&mut f.store, &id, &edge(7)).unwrap_err();
}
#[test]
fn m004_invalidation_refuses_readers_and_fails_toward_retention_order() {
let mut f = fixture("invalidate");
let id = identity(1, 23);
let obj = tagged_object(45);
open_provisional_pin(&mut f.store, &id, &obj, &m017_ctx()).unwrap();
authorize_reader(&mut f.store, &id, &dependent("worker-b", 5)).unwrap();
assert_eq!(
invalidate_lineage(&mut f.store, &id, "producer generation failed").unwrap(),
CloseOutcome::Released
);
assert_eq!(
resolve_for_reader(&mut f.store, &id, &dependent("worker-b", 5)).unwrap_err(),
ProvisionalPinError::ProducerInvalidated {
pin_key: id.pin_key(),
reason: "producer generation failed".to_owned()
}
);
invalidate_lineage(&mut f.store, &id, "second reason").unwrap();
let row = f.store.provisional_pin_row(&id.pin_key()).unwrap().unwrap();
assert_eq!(
row.invalidated_reason.as_deref(),
Some("producer generation failed")
);
assert!(row.released);
let protective = u128::from_str_radix(&row.protective_pin_hex, 16).unwrap();
let pin = f.store.pin_row(protective).unwrap().unwrap();
assert!(!pin.released);
assert_eq!(pin.class, PROVISIONAL_PIN_CLASS);
}
#[test]
fn m004_adoption_requires_exact_object_and_other_success_does_not_stabilize() {
let mut f = fixture("adoption");
let id = identity(1, 24);
let pinned_obj = tagged_object(46);
open_provisional_pin(&mut f.store, &id, &pinned_obj, &m017_ctx()).unwrap();
authorize_reader(&mut f.store, &id, &edge(11)).unwrap();
let before = f.store.provisional_pin_row(&id.pin_key()).unwrap().unwrap();
assert!(!before.released && before.adopted_object_key.is_none());
let err = record_adoption(&mut f.store, &id, &tagged_object(47)).unwrap_err();
assert_eq!(
err,
ProvisionalPinError::AdoptionMismatch {
pin_key: id.pin_key(),
pinned: digest_key(&pinned_obj.0),
committed: digest_key(&tagged_object(47).0),
}
);
record_adoption(&mut f.store, &id, &pinned_obj).unwrap();
let after = f.store.provisional_pin_row(&id.pin_key()).unwrap().unwrap();
assert_eq!(
after.adopted_object_key.as_deref(),
Some(digest_key(&pinned_obj.0).as_str())
);
assert!(!after.released);
assert_eq!(
resolve_for_reader(&mut f.store, &id, &edge(12)).unwrap_err(),
ProvisionalPinError::Unauthorized {
pin_key: id.pin_key()
}
);
assert_eq!(
resolve_for_reader(&mut f.store, &id, &edge(11)).unwrap(),
pinned_obj
);
}
#[test]
fn m004_renewal_is_monotonic_and_closed_pins_refuse() {
let mut f = fixture("renew");
let id = identity(1, 25);
open_provisional_pin(&mut f.store, &id, &tagged_object(48), &m017_ctx()).unwrap();
renew_provisional_pin(&mut f.store, &id, 4).unwrap();
renew_provisional_pin(&mut f.store, &id, 6).unwrap();
assert_eq!(
renew_provisional_pin(&mut f.store, &id, 6).unwrap_err(),
ProvisionalPinError::NonMonotonicRenewal {
pin_key: id.pin_key()
}
);
assert_eq!(
renew_provisional_pin(&mut f.store, &id, 3).unwrap_err(),
ProvisionalPinError::NonMonotonicRenewal {
pin_key: id.pin_key()
}
);
release_authority_fixture(&mut f.store, 1);
release_after_drain(&mut f.store, &id, &Releaser::Coordinator(authority(1))).unwrap();
assert_eq!(
renew_provisional_pin(&mut f.store, &id, 9).unwrap_err(),
ProvisionalPinError::Closed {
pin_key: id.pin_key()
}
);
assert_eq!(
open_provisional_pin(&mut f.store, &id, &tagged_object(48), &m017_ctx()).unwrap_err(),
ProvisionalPinError::Closed {
pin_key: id.pin_key()
}
);
}
#[test]
fn m004_identity_tuple_binds_every_component_into_the_pin_address() {
let base = identity(1, 30);
let variants = [
ProvisionalIdentity {
attempt: AttemptId(31),
lease: ExecutionLeaseId(32),
..base.clone()
},
ProvisionalIdentity {
generation: ActionGenerationId(0x52),
..base.clone()
},
ProvisionalIdentity {
role: OutputRole::DepInfo,
..base.clone()
},
ProvisionalIdentity {
virtual_path: RawBytes::new(b"target/debug/deps/libother.rmeta".to_vec()),
..base.clone()
},
ProvisionalIdentity {
authority: authority(3),
..base.clone()
},
];
let base_key = base.pin_key();
for v in &variants {
assert_ne!(v.pin_key(), base_key);
}
assert_eq!(base.pin_key(), base_key);
}
#[test]
fn m006_consumption_creates_exactly_one_lineage_obligation() {
let mut f = fixture("m006-create");
let producer = identity(1, 60);
let obj = tagged_object(61);
open_provisional_pin(&mut f.store, &producer, &obj, &m017_ctx()).unwrap();
authorize_reader(&mut f.store, &producer, &dependent("worker-c", 70)).unwrap();
authorize_reader(&mut f.store, &producer, &edge(80)).unwrap();
resolve_for_reader(&mut f.store, &producer, &dependent("worker-c", 70)).unwrap();
let rows = f
.store
.list_open_provisional_obligations("worker-c", &format!("{:032x}", 70))
.unwrap();
assert_eq!(rows.len(), 1);
let row = &rows[0];
assert_eq!(row.pin_key, producer.pin_key());
assert_eq!(row.producer_action_key, digest_key(&producer.action_key));
assert_eq!(
row.producer_generation_hex,
format!("{:032x}", producer.generation.0)
);
assert_eq!(
row.producer_attempt_hex,
format!("{:032x}", producer.attempt.0)
);
assert_eq!(row.object_key, digest_key(&obj.0));
assert_eq!(row.status, "open");
resolve_for_reader(&mut f.store, &producer, &dependent("worker-c", 70)).unwrap();
assert_eq!(
f.store
.list_open_provisional_obligations("worker-c", &format!("{:032x}", 70))
.unwrap()
.len(),
1
);
resolve_for_reader(&mut f.store, &producer, &edge(80)).unwrap();
let debt = f
.store
.count_open_provisional_obligations(&producer.pin_key())
.unwrap();
assert_eq!(debt, 1, "only the dependent attempt owes a commit");
}
#[test]
fn m006_terminal_gate_blocks_until_exact_object_commit_resolves() {
let mut f = fixture("m006-gate");
let producer = identity(1, 62);
let obj = tagged_object(63);
open_provisional_pin(&mut f.store, &producer, &obj, &m017_ctx()).unwrap();
authorize_reader(&mut f.store, &producer, &dependent("worker-d", 71)).unwrap();
resolve_for_reader(&mut f.store, &producer, &dependent("worker-d", 71)).unwrap();
assert_eq!(
descendant_terminal_gate(&mut f.store, "worker-d", AttemptId(71)).unwrap(),
TerminalGate::Blocked {
pending_pin_keys: vec![producer.pin_key()]
}
);
let err =
resolve_consumers_on_commit(&mut f.store, &producer, &tagged_object(64)).unwrap_err();
assert!(matches!(err, ProvisionalPinError::AdoptionMismatch { .. }));
assert!(matches!(
descendant_terminal_gate(&mut f.store, "worker-d", AttemptId(71)).unwrap(),
TerminalGate::Blocked { .. }
));
let resolved = resolve_consumers_on_commit(&mut f.store, &producer, &obj).unwrap();
assert_eq!(resolved, 1);
assert_eq!(
descendant_terminal_gate(&mut f.store, "worker-d", AttemptId(71)).unwrap(),
TerminalGate::Clear
);
let rows = f
.store
.list_open_provisional_obligations("worker-d", &format!("{:032x}", 71))
.unwrap();
assert!(
rows.is_empty(),
"resolved obligations leave the non-resolved set"
);
}
#[test]
fn m006_invalidation_permanently_refuses_the_descendant() {
let mut f = fixture("m006-refuse");
let producer = identity(1, 65);
open_provisional_pin(&mut f.store, &producer, &tagged_object(66), &m017_ctx()).unwrap();
authorize_reader(&mut f.store, &producer, &dependent("worker-e", 72)).unwrap();
resolve_for_reader(&mut f.store, &producer, &dependent("worker-e", 72)).unwrap();
invalidate_lineage(&mut f.store, &producer, "producer generation failed").unwrap();
assert_eq!(
descendant_terminal_gate(&mut f.store, "worker-e", AttemptId(72)).unwrap(),
TerminalGate::Refused {
cancelled_pin_keys: vec![producer.pin_key()]
}
);
assert!(matches!(
resolve_for_reader(&mut f.store, &producer, &dependent("worker-e", 72)),
Err(ProvisionalPinError::ProducerInvalidated { .. })
));
}
#[test]
fn m006_drain_refuses_while_consumer_debt_is_open() {
let mut f = fixture("m006-debt");
let producer = identity(1, 67);
open_provisional_pin(&mut f.store, &producer, &tagged_object(68), &m017_ctx()).unwrap();
authorize_reader(&mut f.store, &producer, &dependent("worker-g", 73)).unwrap();
resolve_for_reader(&mut f.store, &producer, &dependent("worker-g", 73)).unwrap();
release_authority_fixture(&mut f.store, 1);
let err = release_after_drain(
&mut f.store,
&producer,
&Releaser::Coordinator(authority(1)),
)
.unwrap_err();
assert_eq!(
err,
ProvisionalPinError::UnresolvedConsumerDebt {
pin_key: producer.pin_key(),
open_count: 1
}
);
assert!(
!f.store
.provisional_pin_row(&producer.pin_key())
.unwrap()
.unwrap()
.released
);
resolve_consumers_on_commit(&mut f.store, &producer, &tagged_object(68)).unwrap();
assert_eq!(
release_after_drain(
&mut f.store,
&producer,
&Releaser::Coordinator(authority(1))
)
.unwrap(),
CloseOutcome::Released
);
}
#[test]
fn m007_generation_failure_invalidates_pins_and_refuses_dependents() {
let mut f = fixture("m007-genfail");
let producer_a = identity(1, 80);
let producer_b = ProvisionalIdentity {
virtual_path: RawBytes::new(b"target/debug/deps/libother.rmeta".to_vec()),
..identity(1, 80)
};
let bystander = ProvisionalIdentity {
action_key: tagged_action(11),
..identity(1, 81)
};
open_provisional_pin(&mut f.store, &producer_a, &tagged_object(90), &m017_ctx()).unwrap();
open_provisional_pin(&mut f.store, &producer_b, &tagged_object(91), &m017_ctx()).unwrap();
open_provisional_pin(&mut f.store, &bystander, &tagged_object(92), &m017_ctx()).unwrap();
for (producer, consumer_attempt) in [
(&producer_a, 91_u128),
(&producer_b, 91_u128),
(&bystander, 92_u128),
] {
authorize_reader(
&mut f.store,
producer,
&dependent("worker-h", consumer_attempt),
)
.unwrap();
resolve_for_reader(
&mut f.store,
producer,
&dependent("worker-h", consumer_attempt),
)
.unwrap();
}
let summary = invalidate_lineage_for_generation_failure(
&mut f.store,
&tagged_action(10),
ActionGenerationId(0x50),
"generation tombstoned after worker loss",
)
.unwrap();
assert_eq!(summary.pins_invalidated, 2);
assert_eq!(summary.obligations_cancelled, 2);
let mut expected = vec![producer_a.pin_key(), producer_b.pin_key()];
expected.sort();
assert_eq!(
descendant_terminal_gate(&mut f.store, "worker-h", AttemptId(91)).unwrap(),
TerminalGate::Refused {
cancelled_pin_keys: expected
}
);
assert!(matches!(
descendant_terminal_gate(&mut f.store, "worker-h", AttemptId(92)).unwrap(),
TerminalGate::Blocked { .. }
));
}
#[test]
fn m007_authority_loss_invalidates_only_that_authoritys_pins() {
let mut f = fixture("m007-authloss");
let dead_authority_pin = identity(1, 82);
let live_authority_pin = identity(2, 83);
open_provisional_pin(
&mut f.store,
&dead_authority_pin,
&tagged_object(93),
&m017_ctx(),
)
.unwrap();
open_provisional_pin(
&mut f.store,
&live_authority_pin,
&tagged_object(94),
&m017_ctx(),
)
.unwrap();
let summary = invalidate_lineage_for_authority_loss(
&mut f.store,
&authority(1),
"operator reset superseded term",
)
.unwrap();
assert_eq!(summary.pins_invalidated, 1);
authorize_reader(&mut f.store, &dead_authority_pin, &edge(21)).unwrap_err();
assert!(matches!(
resolve_for_reader(&mut f.store, &dead_authority_pin, &edge(22)),
Err(ProvisionalPinError::ProducerInvalidated { .. })
));
authorize_reader(&mut f.store, &live_authority_pin, &edge(23)).unwrap();
assert_eq!(
resolve_for_reader(&mut f.store, &live_authority_pin, &edge(23)).unwrap(),
tagged_object(94)
);
}
#[test]
fn m007_supersession_keeps_exact_objects_and_invalidates_divergent() {
let mut f = fixture("m007-supersede");
let kept = identity(1, 84);
let divergent = ProvisionalIdentity {
virtual_path: RawBytes::new(b"target/debug/deps/libother.rmeta".to_vec()),
..identity(1, 84)
};
let winner_object = tagged_object(95);
open_provisional_pin(&mut f.store, &kept, &winner_object, &m017_ctx()).unwrap();
open_provisional_pin(&mut f.store, &divergent, &tagged_object(96), &m017_ctx()).unwrap();
authorize_reader(&mut f.store, &divergent, &dependent("worker-i", 95)).unwrap();
resolve_for_reader(&mut f.store, &divergent, &dependent("worker-i", 95)).unwrap();
let mut committed = std::collections::BTreeSet::new();
committed.insert(digest_key(&winner_object.0));
let summary = invalidate_unadopted_lineage_for_action(
&mut f.store,
&tagged_action(10),
&committed,
"superseded without compatible adoption",
)
.unwrap();
assert_eq!(summary.pins_invalidated, 1);
assert_eq!(summary.obligations_cancelled, 1);
authorize_reader(&mut f.store, &kept, &edge(31)).unwrap();
assert_eq!(
resolve_for_reader(&mut f.store, &kept, &edge(31)).unwrap(),
winner_object
);
assert_eq!(
descendant_terminal_gate(&mut f.store, "worker-i", AttemptId(95)).unwrap(),
TerminalGate::Refused {
cancelled_pin_keys: vec![divergent.pin_key()]
}
);
}
#[test]
fn m007_causal_trace_records_dependents_across_statuses() {
let mut f = fixture("m007-trace");
let producer = identity(1, 85);
let obj = tagged_object(97);
open_provisional_pin(&mut f.store, &producer, &obj, &m017_ctx()).unwrap();
authorize_reader(&mut f.store, &producer, &dependent("worker-j", 96)).unwrap();
resolve_for_reader(&mut f.store, &producer, &dependent("worker-j", 96)).unwrap();
resolve_consumers_on_commit(&mut f.store, &producer, &obj).unwrap();
let doomed = identity(1, 85);
let _ = doomed;
authorize_reader(&mut f.store, &producer, &dependent("worker-k", 97)).unwrap();
resolve_for_reader(&mut f.store, &producer, &dependent("worker-k", 97)).unwrap();
invalidate_lineage(&mut f.store, &producer, "superseded").unwrap();
let trace = provisional_causal_trace(&mut f.store, &producer.pin_key()).unwrap();
assert_eq!(trace.len(), 2);
assert_eq!(trace[0].consumer_worker, "worker-j");
assert_eq!(trace[0].status, "resolved");
assert_eq!(
trace[0].resolution_object_key.as_deref(),
Some(digest_key(&obj.0).as_str())
);
assert_eq!(trace[1].consumer_worker, "worker-k");
assert_eq!(trace[1].status, "cancelled");
assert_eq!(
trace[0].producer_action_key,
digest_key(&producer.action_key)
);
}
fn release_authority_fixture(store: &mut SqlMetadataStore<RusqliteEngine>, tag: u64) {
store
.acquire_authority(&crate::metadata_store::AuthorityRow {
digest: crate::publication::authority_digest(&authority(tag)),
cluster_id: format!("cluster-{tag}"),
incarnation: 0xAA + u128::from(tag),
term: 100 + tag,
acquired_seq: tag,
})
.unwrap();
}
fn m017_ctx() -> ProducerContracts {
ProducerContracts {
toolchain: tagged_action(200),
events: tagged_action(201),
}
}
fn m017_other_ctx() -> ProducerContracts {
ProducerContracts {
toolchain: tagged_action(202),
events: tagged_action(203),
}
}
fn m017_chain(
name: &str,
) -> (
Fixture,
ProvisionalIdentity,
ProvisionalIdentity,
ProvisionalIdentity,
) {
let mut f = fixture(name);
let a = identity(1, 30);
let b = identity(1, 31);
let c = identity(1, 32);
open_provisional_pin(&mut f.store, &a, &tagged_object(141), &m017_ctx()).unwrap();
authorize_reader(&mut f.store, &a, &dependent("worker-b", 31)).unwrap();
resolve_for_reader(&mut f.store, &a, &dependent("worker-b", 31)).unwrap();
open_provisional_pin(&mut f.store, &b, &tagged_object(142), &m017_ctx()).unwrap();
authorize_reader(&mut f.store, &b, &dependent("worker-c", 32)).unwrap();
resolve_for_reader(&mut f.store, &b, &dependent("worker-c", 32)).unwrap();
open_provisional_pin(&mut f.store, &c, &tagged_object(143), &m017_ctx()).unwrap();
(f, a, b, c)
}
#[test]
fn m017_transitive_closure_materialized_at_open() {
let (mut f, a, b, c) = m017_chain("m017-closure");
assert_eq!(
f.store
.list_provisional_pin_ancestors(&b.pin_key())
.unwrap(),
vec![(a.pin_key(), 1)]
);
assert_eq!(
f.store
.list_provisional_pin_ancestors(&c.pin_key())
.unwrap(),
vec![(a.pin_key(), 2), (b.pin_key(), 1)]
);
assert_eq!(
f.store
.list_provisional_pin_descendants(&a.pin_key())
.unwrap(),
{
let mut v = vec![b.pin_key(), c.pin_key()];
v.sort();
v
}
);
}
#[test]
fn m017_commit_resolution_gates_on_whole_closure() {
let (mut f, a, b, c) = m017_chain("m017-gate");
assert_eq!(
resolve_consumers_on_commit(&mut f.store, &c, &tagged_object(143)).unwrap_err(),
ProvisionalPinError::AncestorLineageUnresolved {
pin_key: c.pin_key(),
ancestor_pin_key: b.pin_key(),
detail: format!(
"inbound obligation status open on consumed object {}",
digest_key(&tagged_object(142).0)
),
}
);
record_adoption(&mut f.store, &a, &tagged_object(141)).unwrap();
assert_eq!(
resolve_consumers_on_commit(&mut f.store, &b, &tagged_object(142)).unwrap(),
1
);
assert_eq!(
resolve_consumers_on_commit(&mut f.store, &c, &tagged_object(143)).unwrap(),
0
);
assert_eq!(
descendant_terminal_gate(&mut f.store, "worker-b", AttemptId(31)).unwrap(),
TerminalGate::Clear
);
assert_eq!(
descendant_terminal_gate(&mut f.store, "worker-c", AttemptId(32)).unwrap(),
TerminalGate::Clear
);
}
#[test]
fn m017_different_winner_same_object_adopts_with_explicit_edge() {
let (mut f, a, _b, _c) = m017_chain("m017-adopt");
release_authority_fixture(&mut f.store, 5);
let winner = WinningAttemptContext {
authority: authority(5),
action_key: a.action_key.clone(),
generation: ActionGenerationId(0x51),
attempt: AttemptId(39),
contracts: m017_ctx(),
};
assert_eq!(
adopt_from_winning_attempt(&mut f.store, &a, &winner, &tagged_object(141)).unwrap(),
AdoptionOutcome::Adopted {
obligations_resolved: 1
}
);
assert!(
f.store
.has_adoption_edge(
&digest_key(&a.action_key),
output_role_name_for_tag(i64::try_from(output_role_tag(a.role)).unwrap()),
a.virtual_path.as_bytes(),
&digest_key(&tagged_object(141).0),
&digest_key(&tagged_object(141).0)
)
.unwrap()
);
assert_eq!(
descendant_terminal_gate(&mut f.store, "worker-b", AttemptId(31)).unwrap(),
TerminalGate::Clear
);
}
#[test]
fn m017_divergent_winner_cascades_refusal_to_descendants() {
let (mut f, a, b, c) = m017_chain("m017-diverge");
release_authority_fixture(&mut f.store, 5);
let winner = WinningAttemptContext {
authority: authority(5),
action_key: a.action_key.clone(),
generation: ActionGenerationId(0x51),
attempt: AttemptId(39),
contracts: m017_ctx(),
};
assert_eq!(
adopt_from_winning_attempt(&mut f.store, &a, &winner, &tagged_object(199)).unwrap(),
AdoptionOutcome::DivergenceCancelled {
pins_invalidated: 3,
obligations_cancelled: 2,
}
);
for pin in [&a, &b, &c] {
assert!(matches!(
resolve_for_reader(&mut f.store, pin, &edge(77)),
Err(ProvisionalPinError::ProducerInvalidated { .. })
));
}
for (worker, attempt) in [("worker-b", 31_u128), ("worker-c", 32)] {
assert!(matches!(
descendant_terminal_gate(&mut f.store, worker, AttemptId(attempt)).unwrap(),
TerminalGate::Refused { .. }
));
}
}
#[test]
fn m017_adoption_refuses_foreign_actions_same_attempt_and_contracts() {
let (mut f, a, _b, _c) = m017_chain("m017-refuse");
release_authority_fixture(&mut f.store, 5);
let base = WinningAttemptContext {
authority: authority(5),
action_key: a.action_key.clone(),
generation: ActionGenerationId(0x51),
attempt: AttemptId(39),
contracts: m017_ctx(),
};
assert_eq!(
adopt_from_winning_attempt(
&mut f.store,
&a,
&WinningAttemptContext {
action_key: tagged_action(11),
..base.clone()
},
&tagged_object(141)
)
.unwrap_err(),
ProvisionalPinError::ForeignAction {
pin_key: a.pin_key()
}
);
assert_eq!(
adopt_from_winning_attempt(
&mut f.store,
&a,
&WinningAttemptContext {
generation: a.generation,
attempt: a.attempt,
..base.clone()
},
&tagged_object(141)
)
.unwrap_err(),
ProvisionalPinError::SameWinningAttempt {
pin_key: a.pin_key()
}
);
assert_eq!(
adopt_from_winning_attempt(
&mut f.store,
&a,
&WinningAttemptContext {
contracts: m017_other_ctx(),
..base
},
&tagged_object(141)
)
.unwrap_err(),
ProvisionalPinError::ContractMismatch {
pin_key: a.pin_key()
}
);
}
#[test]
fn m017_release_after_drain_does_not_cascade() {
let mut f = fixture("m017-drain");
let a = identity(1, 40);
let b = identity(1, 41);
open_provisional_pin(&mut f.store, &a, &tagged_object(144), &m017_ctx()).unwrap();
authorize_reader(&mut f.store, &a, &dependent("worker-x", 41)).unwrap();
resolve_for_reader(&mut f.store, &a, &dependent("worker-x", 41)).unwrap();
open_provisional_pin(&mut f.store, &b, &tagged_object(145), &m017_ctx()).unwrap();
record_adoption(&mut f.store, &a, &tagged_object(144)).unwrap();
release_authority_fixture(&mut f.store, 1);
release_after_drain(&mut f.store, &a, &Releaser::Coordinator(authority(1))).unwrap();
authorize_reader(&mut f.store, &b, &edge(78)).unwrap();
assert_eq!(
resolve_for_reader(&mut f.store, &b, &edge(78)).unwrap(),
tagged_object(145)
);
}
#[test]
fn m017_batch_supersession_trigger_cascades_through_pins() {
let (mut f, _a, _b, _c) = m017_chain("m017-batch");
let committed = std::collections::BTreeSet::from([digest_key(&tagged_object(198).0)]);
let summary = invalidate_unadopted_lineage_for_action(
&mut f.store,
&tagged_action(10),
&committed,
"superseded without compatible adoption",
)
.unwrap();
assert_eq!(summary.pins_invalidated, 3);
assert_eq!(summary.obligations_cancelled, 2);
assert!(matches!(
descendant_terminal_gate(&mut f.store, "worker-c", AttemptId(32)).unwrap(),
TerminalGate::Refused { .. }
));
}
}