use crate::adapter::net::behavior::fold::{Fold, IslandId, JobId, NodeId, ReservationFold};
use crate::adapter::net::identity::EntityKeypair;
use super::claim::{release_island, single_island_claim, ClaimError, ClaimOutcome, Claimant};
#[derive(Debug, Clone)]
pub struct GangClaim {
pub job: JobId,
pub islands: Vec<IslandId>,
pub deadline_us: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum AcquireAttempt {
Held,
Blocked(IslandId),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum GangOutcome {
Held(Vec<IslandId>),
DeadlineExceeded,
}
pub fn try_acquire_gang(
reservations: &Fold<ReservationFold>,
keypair: &EntityKeypair,
node_id: NodeId,
generation: &mut u64,
islands_sorted: &[IslandId],
until_unix_us: u64,
) -> Result<AcquireAttempt, ClaimError> {
let mut held: Vec<IslandId> = Vec::with_capacity(islands_sorted.len());
for &island in islands_sorted {
let gen = *generation;
*generation += 1;
match single_island_claim(reservations, keypair, node_id, gen, island, until_unix_us)? {
ClaimOutcome::Won => held.push(island),
ClaimOutcome::Lost => {
for &grabbed in held.iter().rev() {
let gen = *generation;
*generation += 1;
let _ = release_island(reservations, keypair, node_id, gen, grabbed);
}
return Ok(AcquireAttempt::Blocked(island));
}
}
}
Ok(AcquireAttempt::Held)
}
pub fn acquire_gang(
claimant: &mut Claimant,
claim: &GangClaim,
reserve_ttl_us: u64,
now_us: impl Fn() -> u64,
mut backoff: impl FnMut(u32),
) -> Result<GangOutcome, ClaimError> {
let mut islands = claim.islands.clone();
islands.sort_unstable();
islands.dedup();
let mut attempt = 0u32;
loop {
let until = now_us().saturating_add(reserve_ttl_us);
match try_acquire_gang(
claimant.reservations,
claimant.keypair,
claimant.node_id,
&mut claimant.generation,
&islands,
until,
)? {
AcquireAttempt::Held => return Ok(GangOutcome::Held(islands)),
AcquireAttempt::Blocked(_) => {
if now_us() >= claim.deadline_us {
return Ok(GangOutcome::DeadlineExceeded);
}
backoff(attempt);
attempt = attempt.saturating_add(1);
}
}
}
}
#[cfg(test)]
mod tests {
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, Barrier};
use std::time::Duration;
use super::*;
use crate::adapter::net::behavior::fold::ReservationQuery;
use crate::adapter::net::current_timestamp_micros;
fn new_reservations() -> Fold<ReservationFold> {
Fold::with_sweep_interval(Duration::ZERO)
}
fn fresh() -> u64 {
current_timestamp_micros() + 60_000_000
}
fn holder_of(fold: &Fold<ReservationFold>, island: IslandId) -> Option<NodeId> {
fold.query(ReservationQuery::State(island))
.first()
.and_then(|(_, s)| s.holder())
}
#[test]
fn try_acquire_holds_full_set_when_free() {
let fold = new_reservations();
let kp = EntityKeypair::generate();
let n = kp.entity_id().node_id();
let mut g = 1;
let r = try_acquire_gang(&fold, &kp, n, &mut g, &[1, 2, 3], fresh()).unwrap();
assert_eq!(r, AcquireAttempt::Held);
for island in [1, 2, 3] {
assert_eq!(holder_of(&fold, island), Some(n));
}
}
#[test]
fn try_acquire_releases_all_on_block_and_reports_blocker() {
let fold = new_reservations();
let a = EntityKeypair::generate();
let b = EntityKeypair::generate();
let (na, nb) = (a.entity_id().node_id(), b.entity_id().node_id());
single_island_claim(&fold, &b, nb, 1, 2, fresh()).unwrap();
let mut g = 1;
let r = try_acquire_gang(&fold, &a, na, &mut g, &[1, 2, 3], fresh()).unwrap();
assert_eq!(r, AcquireAttempt::Blocked(2));
assert_eq!(holder_of(&fold, 1), None, "island 1 must be released");
assert_eq!(holder_of(&fold, 2), Some(nb), "B still holds 2");
assert_eq!(holder_of(&fold, 3), None, "never reached 3");
}
#[test]
fn ascending_lock_order_prevents_the_classic_two_gang_deadlock() {
let fold = new_reservations();
let a = EntityKeypair::generate();
let b = EntityKeypair::generate();
let (na, nb) = (a.entity_id().node_id(), b.entity_id().node_id());
let mut ca = Claimant::new(&fold, &a, na);
let ra = acquire_gang(
&mut ca,
&GangClaim {
job: 1,
islands: vec![2, 1], deadline_us: fresh(),
},
60_000_000,
current_timestamp_micros,
|_| {},
)
.unwrap();
assert_eq!(
ra,
GangOutcome::Held(vec![1, 2]),
"A holds the set in id order"
);
let mut cb = Claimant::new(&fold, &b, nb);
let rb = acquire_gang(
&mut cb,
&GangClaim {
job: 2,
islands: vec![1, 2],
deadline_us: 0, },
60_000_000,
current_timestamp_micros,
|_| {},
)
.unwrap();
assert_eq!(rb, GangOutcome::DeadlineExceeded);
assert_eq!(holder_of(&fold, 1), Some(na));
assert_eq!(holder_of(&fold, 2), Some(na));
}
#[test]
fn retry_succeeds_after_the_blocker_releases() {
let fold = new_reservations();
let a = EntityKeypair::generate();
let b = EntityKeypair::generate();
let (na, nb) = (a.entity_id().node_id(), b.entity_id().node_id());
single_island_claim(&fold, &a, na, 1, 2, fresh()).unwrap();
let clock = AtomicU64::new(1);
let now = || clock.fetch_add(1, Ordering::Relaxed);
let mut ga_rel = 100u64;
let backoff = |_attempt: u32| {
release_island(&fold, &a, na, ga_rel, 2).unwrap();
ga_rel += 1;
};
let mut cb = Claimant::new(&fold, &b, nb);
let rb = acquire_gang(
&mut cb,
&GangClaim {
job: 2,
islands: vec![1, 2],
deadline_us: u64::MAX,
},
60_000_000,
now,
backoff,
)
.unwrap();
assert_eq!(rb, GangOutcome::Held(vec![1, 2]));
assert_eq!(holder_of(&fold, 1), Some(nb));
assert_eq!(holder_of(&fold, 2), Some(nb));
}
#[test]
fn dead_claimants_expired_reserves_are_taken_over() {
let fold = new_reservations();
let a = EntityKeypair::generate();
let b = EntityKeypair::generate();
let (na, nb) = (a.entity_id().node_id(), b.entity_id().node_id());
let expired = current_timestamp_micros().saturating_sub(60_000_000);
let mut ga = 1;
let ra = try_acquire_gang(&fold, &a, na, &mut ga, &[2, 3], expired).unwrap();
assert_eq!(ra, AcquireAttempt::Held);
assert_eq!(holder_of(&fold, 2), Some(na));
let mut cb = Claimant::new(&fold, &b, nb);
let rb = acquire_gang(
&mut cb,
&GangClaim {
job: 9,
islands: vec![2, 3, 4],
deadline_us: fresh(),
},
60_000_000,
current_timestamp_micros,
|_| {},
)
.unwrap();
assert_eq!(rb, GangOutcome::Held(vec![2, 3, 4]));
for island in [2, 3, 4] {
assert_eq!(holder_of(&fold, island), Some(nb));
}
}
#[test]
fn two_overlapping_gangs_make_bounded_deadlock_free_progress() {
let fold = Arc::new(new_reservations());
let barrier = Arc::new(Barrier::new(2));
let gangs = [vec![1, 2, 3, 4], vec![3, 4, 5, 6]];
let handles: Vec<_> = gangs
.into_iter()
.map(|islands| {
let fold = fold.clone();
let barrier = barrier.clone();
std::thread::spawn(move || {
let kp = EntityKeypair::generate();
let node = kp.entity_id().node_id();
let mut claimant = Claimant::new(&fold, &kp, node);
let claim = GangClaim {
job: 1,
islands: islands.clone(),
deadline_us: current_timestamp_micros() + 5_000_000,
};
barrier.wait();
let outcome = acquire_gang(
&mut claimant,
&claim,
2_000_000,
current_timestamp_micros,
|_| std::thread::sleep(Duration::from_micros(200)),
)
.expect("acquire");
if let GangOutcome::Held(ref held) = outcome {
let mut sorted = islands.clone();
sorted.sort_unstable();
assert_eq!(held, &sorted, "Held set is the full gang, in order");
for &island in held {
release_island(&fold, &kp, node, claimant.next_gen(), island).unwrap();
}
}
matches!(outcome, GangOutcome::Held(_))
})
})
.collect();
let results: Vec<bool> = handles.into_iter().map(|h| h.join().unwrap()).collect();
assert!(
results.iter().all(|&held| held),
"both gangs eventually assembled their full set (bounded, deadlock-free): {results:?}",
);
for island in [1, 2, 3, 4, 5, 6] {
assert_eq!(
holder_of(&fold, island),
None,
"island {island} released at rest"
);
}
}
}