use std::collections::HashMap;
use crate::adapter::net::behavior::fold::{ApplyOutcome, IslandId, JobId, NodeId};
use super::claim::{activate_announcement, ClaimError, Claimant};
use super::quorum::{Epoch, FenceLedger, QuorumWitness, ReplicaSet};
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ActiveCommitOutcome {
Committed,
NoQuorum {
acks: usize,
needed: usize,
},
LostReservation,
}
#[derive(Debug, Clone)]
pub struct ReplicaCohort {
fences: HashMap<NodeId, FenceLedger>,
}
impl ReplicaCohort {
pub fn new(members: &[NodeId]) -> Self {
Self {
fences: members.iter().map(|&m| (m, FenceLedger::new())).collect(),
}
}
fn vote(&mut self, replica: NodeId, island: IslandId, epoch: Epoch) -> bool {
match self.fences.get_mut(&replica) {
Some(fence) => fence.accept_active(island, epoch),
None => false,
}
}
pub fn highest_witnessed(&self, replica: NodeId, island: IslandId) -> Epoch {
self.fences
.get(&replica)
.map(|f| f.highest_witnessed(island))
.unwrap_or(0)
}
}
pub fn commit_active(
claimant: &Claimant,
cohort: &mut ReplicaCohort,
set: &ReplicaSet,
reachable: &[NodeId],
island: IslandId,
job_id: JobId,
epoch: Epoch,
) -> Result<ActiveCommitOutcome, ClaimError> {
let mut witness = QuorumWitness::new(set);
for &replica in reachable {
if cohort.vote(replica, island, epoch) {
witness.record_ack(replica);
}
}
if !witness.has_quorum() {
return Ok(ActiveCommitOutcome::NoQuorum {
acks: witness.ack_count(),
needed: set.quorum_threshold(),
});
}
let ann = activate_announcement(claimant.keypair, claimant.node_id, epoch, island, job_id)?;
match claimant.reservations.apply(ann)? {
ApplyOutcome::Inserted | ApplyOutcome::Replaced => Ok(ActiveCommitOutcome::Committed),
ApplyOutcome::Rejected => Ok(ActiveCommitOutcome::LostReservation),
}
}
#[cfg(test)]
mod tests {
use std::time::Duration;
use super::*;
use crate::adapter::net::behavior::fold::{
Fold, ReservationFold, ReservationQuery, ReservationState,
};
use crate::adapter::net::behavior::gang::single_island_claim;
use crate::adapter::net::current_timestamp_micros;
use crate::adapter::net::identity::EntityKeypair;
fn new_reservations() -> Fold<ReservationFold> {
Fold::with_sweep_interval(Duration::ZERO)
}
fn fresh() -> u64 {
current_timestamp_micros() + 60_000_000
}
fn state_of(fold: &Fold<ReservationFold>, island: IslandId) -> ReservationState {
fold.query(ReservationQuery::State(island))[0].1.clone()
}
#[test]
fn quorum_side_commits_and_applies_active() {
let fold = new_reservations();
let set = ReplicaSet::new([1, 2, 3, 4, 5]);
let mut cohort = ReplicaCohort::new(set.members());
let leader = EntityKeypair::generate();
let ln = leader.entity_id().node_id();
single_island_claim(&fold, &leader, ln, 1, 0xA0, fresh()).unwrap();
let claimant = Claimant::new(&fold, &leader, ln);
let out = commit_active(
&claimant,
&mut cohort,
&set,
&[1, 2, 3], 0xA0,
7, 2, )
.unwrap();
assert_eq!(out, ActiveCommitOutcome::Committed);
assert!(matches!(
state_of(&fold, 0xA0),
ReservationState::Active { job_id: 7, holder } if holder == ln
));
}
#[test]
fn partition_split_lets_at_most_one_side_commit_active() {
let fold = new_reservations();
let set = ReplicaSet::new([1, 2, 3, 4, 5]);
let mut cohort = ReplicaCohort::new(set.members());
let a = EntityKeypair::generate();
let an = a.entity_id().node_id();
single_island_claim(&fold, &a, an, 1, 0xA0, fresh()).unwrap();
let claimant_a = Claimant::new(&fold, &a, an);
let majority =
commit_active(&claimant_a, &mut cohort, &set, &[1, 2, 3], 0xA0, 7, 2).unwrap();
assert_eq!(majority, ActiveCommitOutcome::Committed);
let b = EntityKeypair::generate();
let bn = b.entity_id().node_id();
let claimant_b = Claimant::new(&fold, &b, bn);
let minority = commit_active(&claimant_b, &mut cohort, &set, &[4, 5], 0xA0, 9, 3).unwrap();
assert_eq!(
minority,
ActiveCommitOutcome::NoQuorum { acks: 2, needed: 3 },
"minority side must never reach Active",
);
assert!(matches!(
state_of(&fold, 0xA0),
ReservationState::Active { holder, .. } if holder == an
));
}
#[test]
fn stale_ex_leader_active_is_fenced_even_with_a_majority_reachable() {
let fold = new_reservations();
let set = ReplicaSet::new([1, 2, 3]);
let mut cohort = ReplicaCohort::new(set.members());
let n = EntityKeypair::generate();
let nn = n.entity_id().node_id();
single_island_claim(&fold, &n, nn, 1, 0xA0, fresh()).unwrap();
let claimant_n = Claimant::new(&fold, &n, nn);
assert_eq!(
commit_active(&claimant_n, &mut cohort, &set, &[1, 2, 3], 0xA0, 1, 5).unwrap(),
ActiveCommitOutcome::Committed,
);
for r in [1, 2, 3] {
assert_eq!(cohort.highest_witnessed(r, 0xA0), 5);
}
let o = EntityKeypair::generate();
let on = o.entity_id().node_id();
let claimant_o = Claimant::new(&fold, &o, on);
let stale = commit_active(&claimant_o, &mut cohort, &set, &[1, 2, 3], 0xA0, 2, 4).unwrap();
assert_eq!(
stale,
ActiveCommitOutcome::NoQuorum { acks: 0, needed: 2 },
"stale ex-leader's epoch-4 Active fenced at every replica",
);
assert!(matches!(
state_of(&fold, 0xA0),
ReservationState::Active { holder, .. } if holder == nn
));
}
#[test]
fn quorum_but_lost_reservation_does_not_commit() {
let fold = new_reservations();
let set = ReplicaSet::new([1, 2, 3]);
let mut cohort = ReplicaCohort::new(set.members());
let a = EntityKeypair::generate();
let an = a.entity_id().node_id();
let b = EntityKeypair::generate();
let bn = b.entity_id().node_id();
single_island_claim(&fold, &b, bn, 1, 0xA0, fresh()).unwrap();
let claimant_a = Claimant::new(&fold, &a, an);
let out = commit_active(&claimant_a, &mut cohort, &set, &[1, 2, 3], 0xA0, 7, 2).unwrap();
assert_eq!(out, ActiveCommitOutcome::LostReservation);
assert!(matches!(
state_of(&fold, 0xA0),
ReservationState::Reserved { holder, .. } if holder == bn
));
}
}