pub mod active;
pub mod claim;
pub mod contention;
pub mod filter;
pub mod multi;
#[cfg(feature = "redex")]
pub mod placement;
pub mod quorum;
pub mod schedule;
#[cfg(test)]
mod proptest;
pub use active::{commit_active, ActiveCommitOutcome, ReplicaCohort};
pub use claim::{
activate_announcement, activate_island, release_announcement, release_island,
reserve_announcement, single_island_claim, ClaimError, ClaimOutcome, Claimant,
};
pub use contention::claim_first_available;
pub use filter::{
candidate_hosts, candidate_hosts_for, numeric_filter, select_islands, select_with_affinity,
NumericFilter, SelectionPolicy,
};
pub use multi::{acquire_gang, try_acquire_gang, AcquireAttempt, GangClaim, GangOutcome};
#[cfg(feature = "redex")]
pub use placement::{colocated_island_config, pinned_island_replicas, COLOCATE_WITH_STRICT_KEY};
pub use quorum::{Epoch, FenceLedger, QuorumWitness, ReplicaSet};
pub use schedule::{
schedule_gang, schedule_single, GangRequest, GangScheduler, ScheduleError, Scheduled,
};
use std::collections::HashSet;
use crate::adapter::net::behavior::fold::{
CapabilityFold, CapabilityQuery, Fold, IslandId, IslandQuery, IslandRecord, IslandTopologyFold,
NodeId,
};
#[derive(Debug, Clone)]
pub struct MatchCriteria {
pub capability: CapabilityQuery,
pub numeric: NumericFilter,
pub selection: SelectionPolicy,
pub prefer_capability: Option<String>,
}
pub fn match_islands(
capability_fold: &Fold<CapabilityFold>,
topology_fold: &Fold<IslandTopologyFold>,
criteria: &MatchCriteria,
down_nodes: &HashSet<NodeId>,
) -> Vec<IslandId> {
let candidates = match_island_records(capability_fold, topology_fold, criteria, down_nodes);
select_with_affinity(
candidates,
criteria.selection,
criteria.prefer_capability.clone(),
)
}
fn match_island_records(
capability_fold: &Fold<CapabilityFold>,
topology_fold: &Fold<IslandTopologyFold>,
criteria: &MatchCriteria,
down_nodes: &HashSet<NodeId>,
) -> Vec<IslandRecord> {
let mut hosts = candidate_hosts_for(capability_fold, &criteria.capability);
if !down_nodes.is_empty() {
hosts.retain(|host| !down_nodes.contains(host));
}
if hosts.is_empty() {
return Vec::new();
}
topology_fold
.query(IslandQuery::HostedByAny(hosts))
.into_iter()
.map(|(_, record)| record)
.filter(|record| criteria.numeric.accepts(record))
.collect()
}
pub fn match_islands_sensed(
capability_fold: &Fold<CapabilityFold>,
topology_fold: &Fold<IslandTopologyFold>,
criteria: &MatchCriteria,
down_nodes: &HashSet<NodeId>,
sensed_non_viable: &HashSet<NodeId>,
sensed_viable_order: &[NodeId],
) -> Vec<IslandId> {
let pruned: std::borrow::Cow<'_, HashSet<NodeId>> = if sensed_non_viable.is_empty() {
std::borrow::Cow::Borrowed(down_nodes)
} else {
std::borrow::Cow::Owned(down_nodes.union(sensed_non_viable).copied().collect())
};
let candidates = match_island_records(capability_fold, topology_fold, criteria, &pruned);
if sensed_viable_order.is_empty() || candidates.len() < 2 {
return select_with_affinity(
candidates,
criteria.selection,
criteria.prefer_capability.clone(),
);
}
let hosts: std::collections::HashMap<IslandId, NodeId> =
candidates.iter().map(|r| (r.id, r.host)).collect();
let mut ordered = select_with_affinity(
candidates,
criteria.selection,
criteria.prefer_capability.clone(),
);
let mut provider_rank: std::collections::HashMap<NodeId, usize> =
std::collections::HashMap::with_capacity(sensed_viable_order.len());
for (rank, provider) in sensed_viable_order.iter().enumerate() {
provider_rank.entry(*provider).or_insert(rank);
}
let bands: std::collections::HashMap<IslandId, usize> = ordered
.iter()
.map(|island| {
let band = hosts
.get(island)
.and_then(|host| provider_rank.get(host).copied())
.unwrap_or(usize::MAX);
(*island, band)
})
.collect();
ordered.sort_by_key(|island| bands.get(island).copied().unwrap_or(usize::MAX));
ordered
}
#[cfg(test)]
mod tests {
use std::collections::BTreeMap;
use std::time::Duration;
use super::*;
use crate::adapter::net::behavior::fold::{
CapabilityFilter, CapabilityMembership, EnvelopeMeta, Fold, FoldKind, IslandRecord,
IslandTopologyFold, NodeState, ReservationFold, ReservationQuery, ReservationState,
SignedAnnouncement, UnitSet,
};
use crate::adapter::net::current_timestamp_micros;
use crate::adapter::net::identity::EntityKeypair;
fn announce_capability(
fold: &Fold<CapabilityFold>,
kp: &EntityKeypair,
node: u64,
tags: Vec<String>,
) {
announce_capability_in(fold, kp, node, tags, None);
}
fn announce_capability_in(
fold: &Fold<CapabilityFold>,
kp: &EntityKeypair,
node: u64,
tags: Vec<String>,
region: Option<String>,
) {
let membership = CapabilityMembership {
class_hash: 0x67_70_75, tags,
hardware: None,
state: NodeState::Idle,
region,
price_quote: None,
reflex_addr: None,
allowed_nodes: Vec::new(),
allowed_subnets: Vec::new(),
allowed_groups: Vec::new(),
metadata: BTreeMap::new(),
owner: None,
};
let ann = SignedAnnouncement::sign(
kp,
CapabilityFold::KIND_ID,
membership.class_hash,
node,
1,
EnvelopeMeta::default(),
membership,
)
.expect("sign cap");
fold.apply(ann).expect("apply cap");
}
fn announce_island_gen(
fold: &Fold<IslandTopologyFold>,
kp: &EntityKeypair,
node: u64,
id: IslandId,
units: usize,
load: f32,
generation: u64,
) {
let record = IslandRecord {
id,
units: UnitSet::new((0..units as u32).collect()),
host: node,
capabilities: vec!["model:a1".into()],
load,
p50_latency_us: 1_500,
};
let ann = SignedAnnouncement::sign(
kp,
IslandTopologyFold::KIND_ID,
0,
node,
generation,
EnvelopeMeta::default(),
record,
)
.expect("sign island");
fold.apply(ann).expect("apply island");
}
fn announce_island(
fold: &Fold<IslandTopologyFold>,
kp: &EntityKeypair,
node: u64,
id: IslandId,
units: usize,
load: f32,
) {
let record = IslandRecord {
id,
units: UnitSet::new((0..units as u32).collect()),
host: node,
capabilities: vec!["model:a1".into()],
load,
p50_latency_us: 1_500,
};
let ann = SignedAnnouncement::sign(
kp,
IslandTopologyFold::KIND_ID,
0,
node,
1,
EnvelopeMeta::default(),
record,
)
.expect("sign island");
fold.apply(ann).expect("apply island");
}
fn new_fold<K: crate::adapter::net::behavior::fold::FoldKind>() -> Fold<K> {
Fold::with_sweep_interval(Duration::ZERO)
}
#[test]
fn match_islands_narrows_by_capability_then_numeric_then_orders() {
let caps: Fold<CapabilityFold> = new_fold();
let topo: Fold<IslandTopologyFold> = new_fold();
let kp_a = EntityKeypair::generate();
let kp_b = EntityKeypair::generate();
let kp_c = EntityKeypair::generate();
let (na, nb, nc) = (
kp_a.entity_id().node_id(),
kp_b.entity_id().node_id(),
kp_c.entity_id().node_id(),
);
announce_capability(&caps, &kp_a, na, vec!["gpu:h100".into()]);
announce_capability(&caps, &kp_b, nb, vec!["gpu:h100".into()]);
announce_capability(&caps, &kp_c, nc, vec!["gpu:a10".into()]);
announce_island(&topo, &kp_a, na, 0xA0, 8, 0.6);
announce_island(&topo, &kp_a, na, 0xA5, 8, 0.2);
announce_island(&topo, &kp_b, nb, 0xB0, 8, 0.4);
announce_island(&topo, &kp_c, nc, 0xC0, 8, 0.0);
let criteria = MatchCriteria {
capability: CapabilityQuery::Composite(CapabilityFilter {
tags_all: vec!["gpu:h100".into()],
..Default::default()
}),
numeric: NumericFilter {
min_units: 8,
max_load: Some(0.5),
..Default::default()
},
selection: SelectionPolicy::LeastLoaded,
prefer_capability: None,
};
let order = match_islands(&caps, &topo, &criteria, &HashSet::new());
assert_eq!(order, vec![0xA5, 0xB0]);
}
#[test]
fn sensed_match_prunes_non_viable_and_ranks_viable_first() {
let caps: Fold<CapabilityFold> = new_fold();
let topo: Fold<IslandTopologyFold> = new_fold();
let kp_a = EntityKeypair::generate();
let kp_b = EntityKeypair::generate();
let kp_c = EntityKeypair::generate();
let (na, nb, nc) = (
kp_a.entity_id().node_id(),
kp_b.entity_id().node_id(),
kp_c.entity_id().node_id(),
);
for (kp, node) in [(&kp_a, na), (&kp_b, nb), (&kp_c, nc)] {
announce_capability(&caps, kp, node, vec!["gpu:h100".into()]);
}
announce_island(&topo, &kp_a, na, 0xA0, 8, 0.1);
announce_island(&topo, &kp_b, nb, 0xB0, 8, 0.2);
announce_island(&topo, &kp_c, nc, 0xC0, 8, 0.3);
let criteria = MatchCriteria {
capability: CapabilityQuery::Composite(CapabilityFilter {
tags_all: vec!["gpu:h100".into()],
..Default::default()
}),
numeric: NumericFilter {
min_units: 8,
..Default::default()
},
selection: SelectionPolicy::LeastLoaded,
prefer_capability: None,
};
assert_eq!(
match_islands(&caps, &topo, &criteria, &HashSet::new()),
vec![0xA0, 0xB0, 0xC0],
);
let non_viable: HashSet<NodeId> = [nb].into_iter().collect();
let order = match_islands_sensed(
&caps,
&topo,
&criteria,
&HashSet::new(),
&non_viable,
&[nc, na],
);
assert_eq!(
order,
vec![0xC0, 0xA0],
"NotReady host pruned; sensed rank leads the claim order",
);
assert_eq!(
match_islands(&caps, &topo, &criteria, &HashSet::new()),
vec![0xA0, 0xB0, 0xC0],
"one interest's NotReady never suspends the entry",
);
}
fn sensed_fixture() -> (
Fold<CapabilityFold>,
Fold<IslandTopologyFold>,
Vec<EntityKeypair>,
Vec<NodeId>,
MatchCriteria,
) {
let caps: Fold<CapabilityFold> = new_fold();
let topo: Fold<IslandTopologyFold> = new_fold();
let kps: Vec<EntityKeypair> = (0..3).map(|_| EntityKeypair::generate()).collect();
let nodes: Vec<NodeId> = kps.iter().map(|k| k.entity_id().node_id()).collect();
for (i, (kp, node)) in kps.iter().zip(nodes.iter()).enumerate() {
announce_capability(&caps, kp, *node, vec!["gpu:h100".into()]);
announce_island(&topo, kp, *node, 0xA0 + i as u64, 8, 0.1 + 0.1 * i as f32);
}
let criteria = MatchCriteria {
capability: CapabilityQuery::Composite(CapabilityFilter {
tags_all: vec!["gpu:h100".into()],
..Default::default()
}),
numeric: NumericFilter {
min_units: 8,
..Default::default()
},
selection: SelectionPolicy::LeastLoaded,
prefer_capability: None,
};
(caps, topo, kps, nodes, criteria)
}
#[test]
fn sensed_match_takes_exactly_one_topology_snapshot() {
let (caps, topo, _kps, nodes, criteria) = sensed_fixture();
let before = topo.metrics().queries();
let order = match_islands_sensed(
&caps,
&topo,
&criteria,
&HashSet::new(),
&HashSet::new(),
&[nodes[2], nodes[0]],
);
let delta = topo.metrics().queries() - before;
assert_eq!(order.len(), 3, "the re-ranking path must have run");
assert_eq!(
delta, 1,
"sensed banding must reuse the selection snapshot, not take a second one",
);
}
#[test]
fn same_host_heartbeat_does_not_change_band_assignment() {
let (caps, topo, kps, nodes, criteria) = sensed_fixture();
let viable = [nodes[2], nodes[0]];
let before = match_islands_sensed(
&caps,
&topo,
&criteria,
&HashSet::new(),
&HashSet::new(),
&viable,
);
announce_island_gen(&topo, &kps[2], nodes[2], 0xA2, 8, 0.99, 2);
let after = match_islands_sensed(
&caps,
&topo,
&criteria,
&HashSet::new(),
&HashSet::new(),
&viable,
);
assert_eq!(
before, after,
"a same-host heartbeat must not change band assignment",
);
}
#[test]
fn eviction_lands_on_the_next_invocation_not_the_completed_one() {
let (caps, topo, _kps, nodes, criteria) = sensed_fixture();
let viable = [nodes[2], nodes[0]];
let first = match_islands_sensed(
&caps,
&topo,
&criteria,
&HashSet::new(),
&HashSet::new(),
&viable,
);
assert_eq!(first, vec![0xA2, 0xA0, 0xA1]);
topo.evict_node(nodes[2], "test");
let second = match_islands_sensed(
&caps,
&topo,
&criteria,
&HashSet::new(),
&HashSet::new(),
&viable,
);
assert_eq!(
first,
vec![0xA2, 0xA0, 0xA1],
"the completed result must be unchanged by a later eviction",
);
assert_eq!(
second,
vec![0xA0, 0xA1],
"the next invocation must reflect the eviction",
);
}
#[test]
fn island_reinserted_under_a_new_host_bands_by_that_host() {
let (caps, topo, kps, nodes, criteria) = sensed_fixture();
let viable = [nodes[1], nodes[0]];
assert_eq!(
match_islands_sensed(
&caps,
&topo,
&criteria,
&HashSet::new(),
&HashSet::new(),
&viable,
),
vec![0xA1, 0xA0, 0xA2],
"B's island leads, then A's, then unsensed C's",
);
topo.evict_node(nodes[2], "test");
announce_island_gen(&topo, &kps[1], nodes[1], 0xA2, 8, 0.35, 1);
let after = match_islands_sensed(
&caps,
&topo,
&criteria,
&HashSet::new(),
&HashSet::new(),
&viable,
);
assert_eq!(
after,
vec![0xA1, 0xA2, 0xA0],
"0xA2 must band by its NEW host B (leading band), not its old host C",
);
}
#[test]
fn sensed_match_with_empty_delta_is_identical_and_potential_is_never_pruned() {
let caps: Fold<CapabilityFold> = new_fold();
let topo: Fold<IslandTopologyFold> = new_fold();
let kp_a = EntityKeypair::generate();
let kp_b = EntityKeypair::generate();
let (na, nb) = (kp_a.entity_id().node_id(), kp_b.entity_id().node_id());
announce_capability(&caps, &kp_a, na, vec!["gpu:h100".into()]);
announce_capability(&caps, &kp_b, nb, vec!["gpu:h100".into()]);
announce_island(&topo, &kp_a, na, 0xA0, 8, 0.1);
announce_island(&topo, &kp_b, nb, 0xB0, 8, 0.2);
let criteria = MatchCriteria {
capability: CapabilityQuery::Composite(CapabilityFilter {
tags_all: vec!["gpu:h100".into()],
..Default::default()
}),
numeric: NumericFilter {
min_units: 8,
..Default::default()
},
selection: SelectionPolicy::LeastLoaded,
prefer_capability: None,
};
let plain = match_islands(&caps, &topo, &criteria, &HashSet::new());
assert_eq!(
match_islands_sensed(
&caps,
&topo,
&criteria,
&HashSet::new(),
&HashSet::new(),
&[],
),
plain,
"empty sensed delta ⇒ byte-identical to match_islands",
);
assert_eq!(
match_islands_sensed(
&caps,
&topo,
&criteria,
&HashSet::new(),
&HashSet::new(),
&[nb],
),
vec![0xB0, 0xA0],
"potential hosts trail the viable band but are never pruned",
);
}
#[test]
fn match_islands_empty_when_no_capability_match() {
let caps: Fold<CapabilityFold> = new_fold();
let topo: Fold<IslandTopologyFold> = new_fold();
let kp = EntityKeypair::generate();
let n = kp.entity_id().node_id();
announce_capability(&caps, &kp, n, vec!["gpu:a10".into()]);
announce_island(&topo, &kp, n, 0xA0, 8, 0.1);
let criteria = MatchCriteria {
capability: CapabilityQuery::Composite(CapabilityFilter {
tags_all: vec!["gpu:h100".into()],
..Default::default()
}),
numeric: NumericFilter::default(),
selection: SelectionPolicy::LeastLoaded,
prefer_capability: None,
};
assert!(match_islands(&caps, &topo, &criteria, &HashSet::new()).is_empty());
}
#[test]
fn dead_host_islands_are_pruned_from_matching() {
let caps: Fold<CapabilityFold> = new_fold();
let topo: Fold<IslandTopologyFold> = new_fold();
let kp_a = EntityKeypair::generate();
let kp_b = EntityKeypair::generate();
let na = kp_a.entity_id().node_id();
let nb = kp_b.entity_id().node_id();
announce_capability(&caps, &kp_a, na, vec!["gpu:h100".into()]);
announce_capability(&caps, &kp_b, nb, vec!["gpu:h100".into()]);
announce_island(&topo, &kp_a, na, 0xA0, 8, 0.1);
announce_island(&topo, &kp_b, nb, 0xB0, 8, 0.2);
let criteria = MatchCriteria {
capability: CapabilityQuery::Composite(CapabilityFilter {
tags_all: vec!["gpu:h100".into()],
..Default::default()
}),
numeric: NumericFilter::default(),
selection: SelectionPolicy::LeastLoaded,
prefer_capability: None,
};
assert_eq!(
match_islands(&caps, &topo, &criteria, &HashSet::new()),
vec![0xA0, 0xB0],
);
let a_down: HashSet<NodeId> = [na].into_iter().collect();
assert_eq!(match_islands(&caps, &topo, &criteria, &a_down), vec![0xB0]);
let both_down: HashSet<NodeId> = [na, nb].into_iter().collect();
assert!(match_islands(&caps, &topo, &criteria, &both_down).is_empty());
}
#[test]
fn region_filters_at_the_host_stage_not_the_island() {
let caps: Fold<CapabilityFold> = new_fold();
let topo: Fold<IslandTopologyFold> = new_fold();
let kp_east = EntityKeypair::generate();
let kp_west = EntityKeypair::generate();
let ne = kp_east.entity_id().node_id();
let nw = kp_west.entity_id().node_id();
announce_capability_in(
&caps,
&kp_east,
ne,
vec!["gpu:h100".into()],
Some("us-east".into()),
);
announce_capability_in(
&caps,
&kp_west,
nw,
vec!["gpu:h100".into()],
Some("us-west".into()),
);
announce_island(&topo, &kp_east, ne, 0xE0, 8, 0.1);
announce_island(&topo, &kp_west, nw, 0xF0, 8, 0.1);
let east_only = MatchCriteria {
capability: CapabilityQuery::Composite(CapabilityFilter {
tags_all: vec!["gpu:h100".into()],
region: Some("us-east".into()),
..Default::default()
}),
numeric: NumericFilter::default(),
selection: SelectionPolicy::LeastLoaded,
prefer_capability: None,
};
assert_eq!(
match_islands(&caps, &topo, &east_only, &HashSet::new()),
vec![0xE0]
);
let any_region = MatchCriteria {
capability: CapabilityQuery::Composite(CapabilityFilter {
tags_all: vec!["gpu:h100".into()],
..Default::default()
}),
..east_only.clone()
};
let mut both = match_islands(&caps, &topo, &any_region, &HashSet::new());
both.sort_unstable();
assert_eq!(both, vec![0xE0, 0xF0]);
let nowhere = MatchCriteria {
capability: CapabilityQuery::Composite(CapabilityFilter {
tags_all: vec!["gpu:h100".into()],
region: Some("ap-south".into()),
..Default::default()
}),
..east_only.clone()
};
assert!(match_islands(&caps, &topo, &nowhere, &HashSet::new()).is_empty());
}
#[test]
fn pipeline_then_claim_run_release() {
let caps: Fold<CapabilityFold> = new_fold();
let topo: Fold<IslandTopologyFold> = new_fold();
let reservations: Fold<ReservationFold> = new_fold();
let kp = EntityKeypair::generate();
let node = kp.entity_id().node_id();
announce_capability(&caps, &kp, node, vec!["gpu:h100".into()]);
announce_island(&topo, &kp, node, 0xA0, 8, 0.3);
let criteria = MatchCriteria {
capability: CapabilityQuery::Composite(CapabilityFilter {
tags_all: vec!["gpu:h100".into()],
..Default::default()
}),
numeric: NumericFilter {
min_units: 8,
..Default::default()
},
selection: SelectionPolicy::LeastLoaded,
prefer_capability: None,
};
let order = match_islands(&caps, &topo, &criteria, &HashSet::new());
let island = *order.first().expect("a candidate island");
assert_eq!(island, 0xA0);
let deadline = current_timestamp_micros() + 60_000_000;
assert_eq!(
single_island_claim(&reservations, &kp, node, 1, island, deadline).unwrap(),
ClaimOutcome::Won,
);
assert_eq!(
activate_island(&reservations, &kp, node, 2, island, 0x42).unwrap(),
ClaimOutcome::Won,
);
assert!(matches!(
reservations.query(ReservationQuery::State(island))[0].1,
ReservationState::Active { job_id: 0x42, .. }
));
assert_eq!(
release_island(&reservations, &kp, node, 3, island).unwrap(),
ClaimOutcome::Won,
);
assert_eq!(
reservations.query(ReservationQuery::State(island))[0].1,
ReservationState::Free,
);
}
}