use std::time::Duration;
use serde::{Deserialize, Serialize};
use super::state::{FoldEntry, FoldState, MergeAction, NoIndex, NodeId};
use super::{FoldKind, SignedAnnouncement};
pub type IslandId = u64;
pub type UnitId = u32;
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct UnitSet(Vec<UnitId>);
impl UnitSet {
pub fn new(mut units: Vec<UnitId>) -> Self {
units.sort_unstable();
units.dedup();
Self(units)
}
pub fn len(&self) -> usize {
self.0.len()
}
pub fn is_empty(&self) -> bool {
self.0.is_empty()
}
pub fn units(&self) -> &[UnitId] {
&self.0
}
pub fn intersects(&self, other: &UnitSet) -> bool {
let (mut i, mut j) = (0, 0);
while i < self.0.len() && j < other.0.len() {
match self.0[i].cmp(&other.0[j]) {
std::cmp::Ordering::Less => i += 1,
std::cmp::Ordering::Greater => j += 1,
std::cmp::Ordering::Equal => return true,
}
}
false
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct IslandRecord {
pub id: IslandId,
pub units: UnitSet,
pub host: NodeId,
pub capabilities: Vec<String>,
pub load: f32,
pub p50_latency_us: u32,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum IslandQuery {
Get(IslandId),
All,
HostedBy(NodeId),
HostedByAny(std::collections::HashSet<NodeId>),
}
pub type IslandRow = (IslandId, IslandRecord);
#[derive(Debug)]
pub struct IslandTopologyFold;
impl FoldKind for IslandTopologyFold {
const KIND_ID: u16 = 4;
const CHANNEL_PREFIX: &'static str = "fold:island:";
const DEFAULT_TTL: Duration = Duration::from_secs(30);
type Key = IslandId;
type Payload = IslandRecord;
type Query = IslandQuery;
type Result = Vec<IslandRow>;
type Index = NoIndex;
fn key_for(_publisher: NodeId, payload: &Self::Payload) -> IslandId {
payload.id
}
fn build_index() -> NoIndex {
NoIndex
}
fn merge(
existing: Option<&FoldEntry<Self>>,
incoming: &SignedAnnouncement<Self::Payload>,
) -> MergeAction {
if incoming.payload.host != incoming.node_id {
return MergeAction::Reject;
}
if !incoming.payload.load.is_finite() {
return MergeAction::Reject;
}
match existing {
None => MergeAction::Insert,
Some(entry) => {
if entry.node_id != incoming.node_id {
return MergeAction::Reject;
}
if incoming.generation > entry.generation {
MergeAction::Replace
} else {
MergeAction::Reject
}
}
}
}
fn query(state: &FoldState<Self>, _index: &NoIndex, query: IslandQuery) -> Vec<IslandRow> {
match query {
IslandQuery::Get(id) => state
.entries
.get(&id)
.map(|e| vec![(id, e.payload.clone())])
.unwrap_or_default(),
IslandQuery::All => state
.entries
.iter()
.map(|(k, e)| (*k, e.payload.clone()))
.collect(),
IslandQuery::HostedBy(host) => state
.entries
.iter()
.filter(|(_, e)| e.payload.host == host)
.map(|(k, e)| (*k, e.payload.clone()))
.collect(),
IslandQuery::HostedByAny(hosts) => state
.entries
.iter()
.filter(|(_, e)| hosts.contains(&e.payload.host))
.map(|(k, e)| (*k, e.payload.clone()))
.collect(),
}
}
}
#[cfg(test)]
mod tests {
use std::sync::Arc;
use std::time::Duration;
use super::*;
use crate::adapter::net::behavior::fold::{
ApplyOutcome, EnvelopeMeta, Fold, FoldRegistry, SignedAnnouncement,
};
use crate::adapter::net::identity::EntityKeypair;
fn sign_island(
keypair: &EntityKeypair,
node_id: NodeId,
generation: u64,
record: IslandRecord,
) -> SignedAnnouncement<IslandRecord> {
SignedAnnouncement::sign(
keypair,
IslandTopologyFold::KIND_ID,
0, node_id,
generation,
EnvelopeMeta::default(),
record,
)
.expect("sign succeeds")
}
fn record(id: IslandId, host: NodeId, load: f32) -> IslandRecord {
IslandRecord {
id,
units: UnitSet::new(vec![0, 1, 2, 3]),
host,
capabilities: vec!["model:a1".into()],
load,
p50_latency_us: 1_500,
}
}
fn new_fold() -> Fold<IslandTopologyFold> {
Fold::with_sweep_interval(Duration::ZERO)
}
#[test]
fn first_announcement_installs_the_island() {
let fold = new_fold();
let kp = EntityKeypair::generate();
let outcome = fold
.apply(sign_island(&kp, 0xAA, 1, record(0x10, 0xAA, 0.25)))
.expect("apply");
assert_eq!(outcome, ApplyOutcome::Inserted);
let q = fold.query(IslandQuery::Get(0x10));
assert_eq!(q.len(), 1);
assert_eq!(q[0].1.host, 0xAA);
assert_eq!(q[0].1.load, 0.25);
}
#[test]
fn host_re_announce_replaces_live_axes() {
let fold = new_fold();
let kp = EntityKeypair::generate();
fold.apply(sign_island(&kp, 0xAA, 1, record(0x10, 0xAA, 0.10)))
.expect("first");
let outcome = fold
.apply(sign_island(&kp, 0xAA, 2, record(0x10, 0xAA, 0.90)))
.expect("heartbeat");
assert_eq!(outcome, ApplyOutcome::Replaced);
assert_eq!(fold.query(IslandQuery::Get(0x10))[0].1.load, 0.90);
}
#[test]
fn stale_generation_from_host_is_rejected() {
let fold = new_fold();
let kp = EntityKeypair::generate();
fold.apply(sign_island(&kp, 0xAA, 5, record(0x10, 0xAA, 0.5)))
.expect("gen=5");
assert_eq!(
fold.apply(sign_island(&kp, 0xAA, 5, record(0x10, 0xAA, 0.1)))
.unwrap(),
ApplyOutcome::Rejected,
);
assert_eq!(
fold.apply(sign_island(&kp, 0xAA, 4, record(0x10, 0xAA, 0.1)))
.unwrap(),
ApplyOutcome::Rejected,
);
assert_eq!(fold.query(IslandQuery::Get(0x10))[0].1.load, 0.5);
}
#[test]
fn announcement_for_a_non_self_host_is_rejected() {
let fold = new_fold();
let kp = EntityKeypair::generate();
let outcome = fold
.apply(sign_island(&kp, 0xAA, 1, record(0x10, 0xBB, 0.0)))
.expect("apply");
assert_eq!(outcome, ApplyOutcome::Rejected);
assert!(fold.query(IslandQuery::Get(0x10)).is_empty());
}
#[test]
fn foreign_publisher_cannot_take_over_an_island_key() {
let fold = new_fold();
let kp_a = EntityKeypair::generate();
let kp_b = EntityKeypair::generate();
fold.apply(sign_island(&kp_a, 0xAA, 1, record(0x10, 0xAA, 0.2)))
.expect("A installs");
let outcome = fold
.apply(sign_island(&kp_b, 0xBB, 99, record(0x10, 0xBB, 0.0)))
.expect("B attempts takeover");
assert_eq!(outcome, ApplyOutcome::Rejected);
assert_eq!(fold.query(IslandQuery::Get(0x10))[0].1.host, 0xAA);
}
#[test]
fn query_all_and_hosted_by() {
let fold = new_fold();
let kp_a = EntityKeypair::generate();
let kp_b = EntityKeypair::generate();
fold.apply(sign_island(&kp_a, 0xAA, 1, record(0x10, 0xAA, 0.2)))
.unwrap();
fold.apply(sign_island(&kp_a, 0xAA, 1, record(0x11, 0xAA, 0.3)))
.unwrap();
fold.apply(sign_island(&kp_b, 0xBB, 1, record(0x20, 0xBB, 0.4)))
.unwrap();
let mut all: Vec<IslandId> = fold
.query(IslandQuery::All)
.into_iter()
.map(|(id, _)| id)
.collect();
all.sort();
assert_eq!(all, vec![0x10, 0x11, 0x20]);
let mut by_a: Vec<IslandId> = fold
.query(IslandQuery::HostedBy(0xAA))
.into_iter()
.map(|(id, _)| id)
.collect();
by_a.sort();
assert_eq!(by_a, vec![0x10, 0x11]);
let hosts: std::collections::HashSet<NodeId> = [0xAA, 0xBB].into_iter().collect();
let mut any: Vec<IslandId> = fold
.query(IslandQuery::HostedByAny(hosts))
.into_iter()
.map(|(id, _)| id)
.collect();
any.sort();
assert_eq!(any, vec![0x10, 0x11, 0x20]);
let only_b: std::collections::HashSet<NodeId> = [0xBB].into_iter().collect();
assert_eq!(
fold.query(IslandQuery::HostedByAny(only_b))
.into_iter()
.map(|(id, _)| id)
.collect::<Vec<_>>(),
vec![0x20],
);
}
#[test]
fn non_finite_load_is_rejected() {
let fold = new_fold();
let kp = EntityKeypair::generate();
for bad in [f32::NAN, f32::INFINITY, f32::NEG_INFINITY] {
let outcome = fold
.apply(sign_island(&kp, 0xAA, 1, record(0x10, 0xAA, bad)))
.expect("apply");
assert_eq!(outcome, ApplyOutcome::Rejected, "load {bad} rejected");
}
assert!(fold.query(IslandQuery::Get(0x10)).is_empty());
}
#[test]
fn runtime_ttl_sweeps_stale_islands() {
let fold = new_fold();
let kp = EntityKeypair::generate();
let ann = SignedAnnouncement::sign(
&kp,
IslandTopologyFold::KIND_ID,
0,
0xAA,
1,
EnvelopeMeta {
ttl_secs: Some(0),
..Default::default()
},
record(0x10, 0xAA, 0.5),
)
.unwrap();
fold.apply(ann).unwrap();
assert_eq!(fold.metrics().entries(), 1);
std::thread::sleep(Duration::from_millis(10));
let n = fold.sweep_expired_now();
assert_eq!(n, 1);
assert!(fold.query(IslandQuery::Get(0x10)).is_empty());
}
#[test]
fn unit_set_normalizes_and_intersects() {
let a = UnitSet::new(vec![3, 1, 1, 2]);
assert_eq!(a.units(), &[1, 2, 3]);
assert_eq!(a.len(), 3);
assert!(a.intersects(&UnitSet::new(vec![5, 3])));
assert!(!a.intersects(&UnitSet::new(vec![4, 5, 6])));
assert!(!a.intersects(&UnitSet::default()));
}
#[test]
fn island_fold_plugs_into_registry_and_dispatches_signed_envelopes() {
let registry = FoldRegistry::new();
let fold: Arc<Fold<IslandTopologyFold>> = Arc::new(new_fold());
registry.register(fold.clone());
let kp = EntityKeypair::generate();
let nid = kp.entity_id().node_id();
let ann = sign_island(&kp, nid, 1, record(0x10, nid, 0.5));
let bytes = ann.encode().expect("encode");
let outcome = registry.dispatch(&bytes, kp.entity_id()).expect("dispatch");
assert_eq!(outcome, ApplyOutcome::Inserted);
}
}