use crate::metadata_store::{RabsMetadataStore, StoreError};
use crate::provisional_pins::ProvisionalPinError;
use rabs_protocol::generation::AttemptId;
use std::sync::Arc;
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum TerminalDelivery {
Ready,
Withheld {
pending_pin_keys: Vec<String>,
},
Refused {
refused_pin_keys: Vec<String>,
},
}
pub fn lineage_gated_terminal_delivery(
store: &mut dyn RabsMetadataStore,
consumer_attempt: AttemptId,
) -> Result<TerminalDelivery, ProvisionalPinError> {
let rows = store
.list_provisional_obligations_by_attempt_all(&format!("{:032x}", consumer_attempt.0))?;
let mut pending = Vec::new();
let mut refused = Vec::new();
for row in rows {
match row.status.as_str() {
"open" => pending.push(row.pin_key),
"cancelled" => refused.push(row.pin_key),
"resolved" => {
if row.resolution_object_key.as_deref() != Some(row.object_key.as_str()) {
refused.push(row.pin_key);
}
}
other => {
return Err(ProvisionalPinError::Store(StoreError::Corruption(format!(
"obligation status {other:?}"
))));
}
}
}
Ok(if !refused.is_empty() {
TerminalDelivery::Refused {
refused_pin_keys: refused,
}
} else if !pending.is_empty() {
TerminalDelivery::Withheld {
pending_pin_keys: pending,
}
} else {
TerminalDelivery::Ready
})
}
pub fn lineage_wait_depth(
store: &mut dyn RabsMetadataStore,
consumer_attempt: AttemptId,
) -> Result<u64, ProvisionalPinError> {
let mut depth = 0u64;
for obligation in store
.list_open_provisional_obligations_by_attempt(&format!("{:032x}", consumer_attempt.0))?
{
depth = depth.max(
store
.provisional_pin_closure_depth(&obligation.pin_key)?
.saturating_add(1),
);
}
Ok(depth)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct WaiterBounds {
pub max_concurrent: usize,
pub reserved_progress_slots: usize,
pub max_lineage_depth: u64,
}
impl Default for WaiterBounds {
fn default() -> Self {
Self {
max_concurrent: 8,
reserved_progress_slots: 1,
max_lineage_depth: 16,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum WaiterRefusal {
Saturated {
root: String,
active: usize,
capacity: usize,
},
DepthExceeded {
root: String,
depth: u64,
max: u64,
},
AlreadyAdmitted {
root: String,
},
InvalidBounds,
}
impl std::fmt::Display for WaiterRefusal {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Saturated {
root,
active,
capacity,
} => write!(
f,
"provisional-lineage waiters saturated on root {root}: {active}/{capacity}"
),
Self::DepthExceeded { root, depth, max } => {
write!(
f,
"lineage depth {depth} exceeds bound {max} on root {root}"
)
}
Self::AlreadyAdmitted { root } => {
write!(f, "attempt already admitted as waiter on root {root}")
}
Self::InvalidBounds => write!(
f,
"producer reserve must be nonzero and no greater than root capacity"
),
}
}
}
impl std::error::Error for WaiterRefusal {}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum WaiterReleaseRefusal {
ForeignRegistry,
UnknownAdmission,
}
impl std::fmt::Display for WaiterReleaseRefusal {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::ForeignRegistry => write!(f, "waiter permit belongs to another registry"),
Self::UnknownAdmission => write!(f, "waiter permit has no matching active admission"),
}
}
}
impl std::error::Error for WaiterReleaseRefusal {}
#[derive(Debug)]
pub struct WaiterPermit {
owner: Arc<()>,
root: String,
attempt: u128,
released: bool,
}
impl WaiterPermit {
#[must_use]
pub fn root(&self) -> &str {
&self.root
}
#[must_use]
pub fn is_released(&self) -> bool {
self.released
}
}
#[derive(Debug, Default)]
struct RootState {
waiters: std::collections::HashSet<u128>,
}
#[derive(Debug)]
pub struct WaiterRegistry {
owner: Arc<()>,
roots: std::collections::HashMap<String, RootState>,
bounds: WaiterBounds,
}
impl Default for WaiterRegistry {
fn default() -> Self {
Self::with_defaults()
}
}
impl WaiterRegistry {
#[must_use]
pub fn new(bounds: WaiterBounds) -> Self {
Self {
owner: Arc::new(()),
roots: std::collections::HashMap::new(),
bounds,
}
}
#[must_use]
pub fn with_defaults() -> Self {
Self::new(WaiterBounds::default())
}
fn effective_capacity(&self) -> Result<usize, WaiterRefusal> {
if self.bounds.reserved_progress_slots == 0 {
return Err(WaiterRefusal::InvalidBounds);
}
self.bounds
.max_concurrent
.checked_sub(self.bounds.reserved_progress_slots)
.ok_or(WaiterRefusal::InvalidBounds)
}
pub fn admit(
&mut self,
root: &str,
attempt: AttemptId,
lineage_depth: u64,
) -> Result<WaiterPermit, WaiterRefusal> {
let capacity = self.effective_capacity()?;
if self
.roots
.get(root)
.is_some_and(|state| state.waiters.contains(&attempt.0))
{
return Err(WaiterRefusal::AlreadyAdmitted {
root: root.to_owned(),
});
}
if lineage_depth > self.bounds.max_lineage_depth {
return Err(WaiterRefusal::DepthExceeded {
root: root.to_owned(),
depth: lineage_depth,
max: self.bounds.max_lineage_depth,
});
}
let active = self.active_waiters(root);
if active >= capacity {
return Err(WaiterRefusal::Saturated {
root: root.to_owned(),
active,
capacity,
});
}
self.roots
.entry(root.to_owned())
.or_default()
.waiters
.insert(attempt.0);
Ok(WaiterPermit {
owner: Arc::clone(&self.owner),
root: root.to_owned(),
attempt: attempt.0,
released: false,
})
}
pub fn release(&mut self, permit: &mut WaiterPermit) -> Result<(), WaiterReleaseRefusal> {
if !Arc::ptr_eq(&self.owner, &permit.owner) {
return Err(WaiterReleaseRefusal::ForeignRegistry);
}
if permit.released {
return Ok(());
}
let state = self
.roots
.get_mut(&permit.root)
.ok_or(WaiterReleaseRefusal::UnknownAdmission)?;
if !state.waiters.remove(&permit.attempt) {
return Err(WaiterReleaseRefusal::UnknownAdmission);
}
let root_is_empty = state.waiters.is_empty();
permit.released = true;
if root_is_empty {
self.roots.remove(&permit.root);
}
Ok(())
}
#[must_use]
pub fn active_waiters(&self, root: &str) -> usize {
self.roots.get(root).map_or(0, |state| state.waiters.len())
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::metadata_store::{AuthorityRow, RusqliteEngine, SqlMetadataStore, digest_key};
use crate::provisional_pins::{
AdoptionOutcome, ProducerContracts, ProvisionalIdentity, ProvisionalReader,
WinningAttemptContext, adopt_from_winning_attempt, authorize_reader, open_provisional_pin,
record_adoption, resolve_consumers_on_commit, resolve_for_reader,
};
use crate::publication::authority_digest;
use rabs_protocol::authority::{ClusterId, CoordinatorAuthority};
use rabs_protocol::generation::{ActionGenerationId, ExecutionLeaseId};
use rabs_protocol::raw_bytes::RawBytes;
use rabs_protocol::result_identity::{DigestAlgorithm, ObjectId, OutputRole, TypedDigest};
const ACTION_DOMAIN: &str = "rabs.action-key.sha256.v1";
fn tagged(tag: u8, domain: &'static str) -> TypedDigest {
let mut bytes = [0u8; 32];
bytes[0] = tag;
bytes[31] = tag;
TypedDigest {
algorithm: DigestAlgorithm::Sha256V1,
domain,
bytes,
}
}
fn action(tag: u8) -> TypedDigest {
tagged(tag, ACTION_DOMAIN)
}
fn obj(tag: u8) -> ObjectId {
ObjectId(tagged(tag, "rabs.object.sha256.v1"))
}
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(attempt_tag: u128) -> ProvisionalIdentity {
ProvisionalIdentity {
authority: authority(1),
action_key: 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 contracts() -> ProducerContracts {
ProducerContracts {
toolchain: action(200),
events: action(201),
}
}
fn dependent(worker: &str, attempt: u128) -> ProvisionalReader {
ProvisionalReader::DependentAttempt {
worker: worker.to_owned(),
attempt: AttemptId(attempt),
}
}
fn chain() -> (
SqlMetadataStore<RusqliteEngine>,
ProvisionalIdentity,
ProvisionalIdentity,
ProvisionalIdentity,
) {
let engine = RusqliteEngine::open_in_memory().unwrap();
let mut store = SqlMetadataStore::open(engine).unwrap();
let a = identity(30);
let b = identity(31);
let c = identity(32);
open_provisional_pin(&mut store, &a, &obj(141), &contracts()).unwrap();
authorize_reader(&mut store, &a, &dependent("worker-b", 31)).unwrap();
resolve_for_reader(&mut store, &a, &dependent("worker-b", 31)).unwrap();
open_provisional_pin(&mut store, &b, &obj(142), &contracts()).unwrap();
authorize_reader(&mut store, &b, &dependent("worker-c", 32)).unwrap();
resolve_for_reader(&mut store, &b, &dependent("worker-c", 32)).unwrap();
open_provisional_pin(&mut store, &c, &obj(143), &contracts()).unwrap();
(store, a, b, c)
}
fn acquire_active(store: &mut SqlMetadataStore<RusqliteEngine>, tag: u64) {
store
.acquire_authority(&AuthorityRow {
digest: authority_digest(&authority(tag)),
cluster_id: format!("cluster-{tag}"),
incarnation: 0xBB + u128::from(tag),
term: 100 + tag,
acquired_seq: tag,
})
.unwrap();
}
#[test]
fn m020_terminal_success_withheld_until_closure_completes() {
let (mut store, a, b, _c) = chain();
assert_eq!(
lineage_gated_terminal_delivery(&mut store, AttemptId(32)).unwrap(),
TerminalDelivery::Withheld {
pending_pin_keys: vec![b.pin_key()]
}
);
record_adoption(&mut store, &a, &obj(141)).unwrap();
assert_eq!(
lineage_gated_terminal_delivery(&mut store, AttemptId(32)).unwrap(),
TerminalDelivery::Withheld {
pending_pin_keys: vec![b.pin_key()]
}
);
resolve_consumers_on_commit(&mut store, &b, &obj(142)).unwrap();
assert_eq!(
lineage_gated_terminal_delivery(&mut store, AttemptId(32)).unwrap(),
TerminalDelivery::Ready
);
assert_eq!(
lineage_gated_terminal_delivery(&mut store, AttemptId(31)).unwrap(),
TerminalDelivery::Ready
);
assert_eq!(
lineage_gated_terminal_delivery(&mut store, AttemptId(30)).unwrap(),
TerminalDelivery::Ready
);
}
#[test]
fn m020_refused_dominates_after_divergent_winner_cascade() {
let (mut store, a, b, _c) = chain();
acquire_active(&mut store, 5);
let winner = WinningAttemptContext {
authority: authority(5),
action_key: a.action_key.clone(),
generation: ActionGenerationId(0x51),
attempt: AttemptId(39),
contracts: contracts(),
};
assert_eq!(
adopt_from_winning_attempt(&mut store, &a, &winner, &obj(199)).unwrap(),
AdoptionOutcome::DivergenceCancelled {
pins_invalidated: 3,
obligations_cancelled: 2,
}
);
assert_eq!(
lineage_gated_terminal_delivery(&mut store, AttemptId(31)).unwrap(),
TerminalDelivery::Refused {
refused_pin_keys: vec![a.pin_key()]
}
);
assert_eq!(
lineage_gated_terminal_delivery(&mut store, AttemptId(32)).unwrap(),
TerminalDelivery::Refused {
refused_pin_keys: vec![b.pin_key()]
}
);
}
#[test]
fn m020_resolved_to_foreign_bytes_is_divergence_not_satisfaction() {
let (mut store, _a, b, _c) = chain();
store
.resolve_provisional_obligations(&b.pin_key(), &digest_key(&obj(250).0))
.unwrap();
assert_eq!(
lineage_gated_terminal_delivery(&mut store, AttemptId(32)).unwrap(),
TerminalDelivery::Refused {
refused_pin_keys: vec![b.pin_key()]
}
);
}
#[test]
fn m020_early_metadata_keeps_flowing_while_terminal_withheld() {
let (mut store, a, b, _c) = chain();
assert!(matches!(
lineage_gated_terminal_delivery(&mut store, AttemptId(32)).unwrap(),
TerminalDelivery::Withheld { .. }
));
authorize_reader(&mut store, &b, &dependent("worker-d", 44)).unwrap();
assert_eq!(
resolve_for_reader(&mut store, &b, &dependent("worker-d", 44)).unwrap(),
obj(142)
);
authorize_reader(&mut store, &a, &dependent("worker-e", 45)).unwrap();
assert_eq!(
resolve_for_reader(&mut store, &a, &dependent("worker-e", 45)).unwrap(),
obj(141)
);
}
#[test]
fn m020_waiter_count_bounded_with_producer_reserve() {
let mut registry = WaiterRegistry::new(WaiterBounds {
max_concurrent: 3,
reserved_progress_slots: 1,
max_lineage_depth: 16,
});
let mut p1 = registry.admit("root-1", AttemptId(1), 1).unwrap();
registry.admit("root-1", AttemptId(2), 1).unwrap();
assert_eq!(registry.active_waiters("root-1"), 2);
assert_eq!(
registry.admit("root-1", AttemptId(3), 1).unwrap_err(),
WaiterRefusal::Saturated {
root: "root-1".to_owned(),
active: 2,
capacity: 2,
}
);
assert_eq!(
registry.admit("root-1", AttemptId(1), 1).unwrap_err(),
WaiterRefusal::AlreadyAdmitted {
root: "root-1".to_owned()
}
);
registry.admit("root-2", AttemptId(3), 1).unwrap();
registry.release(&mut p1).unwrap();
assert!(p1.is_released());
registry.release(&mut p1).unwrap();
registry.admit("root-1", AttemptId(3), 1).unwrap();
assert_eq!(registry.active_waiters("root-1"), 2);
}
#[test]
fn m020_waiter_depth_bounded_regardless_of_capacity() {
let mut registry = WaiterRegistry::new(WaiterBounds {
max_concurrent: 8,
reserved_progress_slots: 1,
max_lineage_depth: 4,
});
assert_eq!(
registry.admit("r", AttemptId(7), 5).unwrap_err(),
WaiterRefusal::DepthExceeded {
root: "r".to_owned(),
depth: 5,
max: 4,
}
);
registry.admit("r", AttemptId(7), 4).unwrap();
}
#[test]
fn m020_invalid_bounds_refuse_every_admission() {
let mut registry = WaiterRegistry::new(WaiterBounds {
max_concurrent: 1,
reserved_progress_slots: 2,
max_lineage_depth: 4,
});
assert_eq!(
registry.admit("r", AttemptId(1), 1).unwrap_err(),
WaiterRefusal::InvalidBounds
);
}
#[test]
fn waiter_permits_cannot_release_another_registrys_matching_slot() {
let bounds = WaiterBounds {
max_concurrent: 2,
reserved_progress_slots: 1,
max_lineage_depth: 4,
};
let mut a = WaiterRegistry::new(bounds);
let mut b = WaiterRegistry::new(bounds);
let mut pa = a.admit("root", AttemptId(1), 1).unwrap();
let mut pb = b.admit("root", AttemptId(1), 1).unwrap();
assert_eq!(
a.release(&mut pb),
Err(WaiterReleaseRefusal::ForeignRegistry)
);
assert!(
!pb.is_released(),
"the true owner must still be able to release"
);
assert_eq!(pb.root(), "root");
assert_eq!(a.active_waiters("root"), 1);
assert_eq!(b.active_waiters("root"), 1);
assert_eq!(
a.admit("root", AttemptId(2), 1).unwrap_err(),
WaiterRefusal::Saturated {
root: "root".to_owned(),
active: 1,
capacity: 1,
}
);
let mut empty = WaiterRegistry::default();
assert_eq!(
empty.release(&mut pa),
Err(WaiterReleaseRefusal::ForeignRegistry)
);
assert!(!pa.is_released());
assert!(empty.roots.is_empty());
b.release(&mut pb).unwrap();
assert_eq!(a.active_waiters("root"), 1);
a.release(&mut pa).unwrap();
assert!(a.roots.is_empty());
assert!(b.roots.is_empty());
}
#[test]
fn stale_waiter_permit_cannot_release_a_replacement_registry() {
let mut stale = {
let mut old = WaiterRegistry::default();
old.admit("root", AttemptId(1), 1).unwrap()
};
let mut replacement = WaiterRegistry::with_defaults();
let mut current = replacement.admit("root", AttemptId(1), 1).unwrap();
assert_eq!(
replacement.release(&mut stale),
Err(WaiterReleaseRefusal::ForeignRegistry)
);
assert!(!stale.is_released());
assert_eq!(replacement.active_waiters("root"), 1);
replacement.release(&mut current).unwrap();
assert!(replacement.roots.is_empty());
}
#[test]
fn waiter_release_survives_moves_and_is_idempotent_after_readmission() {
let mut registry = WaiterRegistry::default();
let mut old = registry.admit("root", AttemptId(1), 1).unwrap();
let mut moved = registry;
moved.release(&mut old).unwrap();
assert!(old.is_released());
assert_eq!(old.root(), "root", "release retains diagnostic identity");
assert!(moved.roots.is_empty());
let mut current = moved.admit("root", AttemptId(1), 1).unwrap();
moved.release(&mut old).unwrap();
assert_eq!(moved.active_waiters("root"), 1);
moved.release(&mut current).unwrap();
moved.release(&mut current).unwrap();
assert!(moved.roots.is_empty());
}
#[test]
fn missing_admission_does_not_consume_a_local_waiter_permit() {
let mut registry = WaiterRegistry::default();
let mut permit = registry.admit("root", AttemptId(1), 1).unwrap();
registry.roots.get_mut("root").unwrap().waiters.remove(&1);
assert_eq!(
registry.release(&mut permit),
Err(WaiterReleaseRefusal::UnknownAdmission)
);
assert!(!permit.is_released());
assert_eq!(permit.root(), "root");
registry.roots.get_mut("root").unwrap().waiters.insert(1);
registry.release(&mut permit).unwrap();
assert!(registry.roots.is_empty());
}
#[test]
fn refused_waiters_do_not_accumulate_empty_root_entries() {
let mut registry = WaiterRegistry::default();
let mut serial = WaiterRegistry::new(WaiterBounds {
max_concurrent: 1,
reserved_progress_slots: 1,
max_lineage_depth: 16,
});
for n in 0..64 {
let root = format!("root-{n}");
assert_eq!(
registry.admit(&root, AttemptId(n), 17).unwrap_err(),
WaiterRefusal::DepthExceeded {
root: root.clone(),
depth: 17,
max: 16,
}
);
assert_eq!(
serial.admit(&root, AttemptId(n), 1).unwrap_err(),
WaiterRefusal::Saturated {
root,
active: 0,
capacity: 0,
}
);
assert!(registry.roots.is_empty());
assert!(serial.roots.is_empty());
}
}
#[test]
fn finished_roots_are_retired_without_affecting_live_roots() {
let mut registry = WaiterRegistry::default();
let mut live = registry.admit("live", AttemptId(1), 1).unwrap();
for n in 0..64 {
let root = format!("finished-{n}");
let mut permit = registry.admit(&root, AttemptId(n), 1).unwrap();
assert_eq!(registry.roots.len(), 2);
registry.release(&mut permit).unwrap();
assert_eq!(registry.roots.len(), 1);
assert_eq!(registry.active_waiters("live"), 1);
assert_eq!(registry.active_waiters(&root), 0);
}
registry.release(&mut live).unwrap();
assert!(registry.roots.is_empty());
}
#[test]
fn waiter_bounds_cannot_disable_the_producer_reserve() {
for (total, reserve) in [(0, 0), (0, 1), (3, 0)] {
let mut registry = WaiterRegistry::new(WaiterBounds {
max_concurrent: total,
reserved_progress_slots: reserve,
max_lineage_depth: 4,
});
assert_eq!(
registry.admit("root", AttemptId(1), 1).unwrap_err(),
WaiterRefusal::InvalidBounds
);
assert!(registry.roots.is_empty());
}
}
#[test]
fn m020_lineage_wait_depth_tracks_transitive_chain() {
let (mut store, _a, _b, c) = chain();
assert_eq!(lineage_wait_depth(&mut store, AttemptId(32)).unwrap(), 2);
assert_eq!(lineage_wait_depth(&mut store, AttemptId(31)).unwrap(), 1);
assert_eq!(lineage_wait_depth(&mut store, AttemptId(30)).unwrap(), 0);
assert_eq!(
store.provisional_pin_closure_depth(&c.pin_key()).unwrap(),
2
);
assert_eq!(
store
.provisional_pin_closure_depth(&identity(99).pin_key())
.unwrap(),
0
);
}
}