use std::collections::hash_map::Entry;
use std::collections::HashMap;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use std::time::Duration;
use parking_lot::Mutex;
use super::identity::{
strictest_sample_interval, AudienceScopeCommitment, Digest256, InterestSpec,
};
const MAX_LEASED_INTERESTS: usize = 256;
const MAX_HOLDERS_PER_INTEREST: usize = 64;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum LeaseRefused {
NodeAtCapacity,
InterestAtCapacity,
}
#[derive(Default)]
struct LeaseMetrics {
refused_node_at_capacity: AtomicU64,
refused_interest_at_capacity: AtomicU64,
reconcile_failures: AtomicU64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub struct LeaseToken(u64);
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum SensingLeaseKey {
ProviderFree {
audience: AudienceScopeCommitment,
interest_digest: Digest256,
},
ExactProvider {
audience: AudienceScopeCommitment,
interest_digest: Digest256,
provider: u64,
},
}
#[derive(Debug, Clone, PartialEq)]
pub enum LeaseAction {
Register {
spec: Arc<InterestSpec>,
interval: Duration,
},
Reregister {
spec: Arc<InterestSpec>,
interval: Duration,
},
Unchanged,
Deregister {
spec: Arc<InterestSpec>,
},
}
#[derive(Debug, Clone, Copy)]
pub struct SensingLeaseTicket {
pub(crate) key: SensingLeaseKey,
pub(crate) token: LeaseToken,
}
struct LeaseEntry {
spec: Arc<InterestSpec>,
registrations: HashMap<LeaseToken, Duration>,
installed_interval: Duration,
}
#[derive(Default)]
pub struct SensingInterestLeases {
entries: Mutex<HashMap<SensingLeaseKey, LeaseEntry>>,
next_token: AtomicU64,
metrics: LeaseMetrics,
}
impl SensingInterestLeases {
fn mint_token(&self) -> LeaseToken {
LeaseToken(self.next_token.fetch_add(1, Ordering::Relaxed))
}
pub fn acquire(
&self,
key: SensingLeaseKey,
spec: &InterestSpec,
interval: Duration,
) -> Result<(LeaseToken, LeaseAction), LeaseRefused> {
let mut entries = self.entries.lock();
if !entries.contains_key(&key) && entries.len() >= MAX_LEASED_INTERESTS {
self.metrics
.refused_node_at_capacity
.fetch_add(1, Ordering::AcqRel);
return Err(LeaseRefused::NodeAtCapacity);
}
match entries.entry(key) {
Entry::Vacant(v) => {
let token = self.mint_token();
let spec = Arc::new(spec.clone());
let mut registrations = HashMap::new();
registrations.insert(token, interval);
v.insert(LeaseEntry {
spec: Arc::clone(&spec),
registrations,
installed_interval: interval,
});
Ok((token, LeaseAction::Register { spec, interval }))
}
Entry::Occupied(mut o) => {
let entry = o.get_mut();
if entry.registrations.len() >= MAX_HOLDERS_PER_INTEREST {
self.metrics
.refused_interest_at_capacity
.fetch_add(1, Ordering::AcqRel);
return Err(LeaseRefused::InterestAtCapacity);
}
let token = self.mint_token();
entry.registrations.insert(token, interval);
let new_min = interval.min(entry.installed_interval);
if new_min < entry.installed_interval {
entry.installed_interval = new_min;
Ok((
token,
LeaseAction::Reregister {
spec: Arc::clone(&entry.spec),
interval: new_min,
},
))
} else {
Ok((token, LeaseAction::Unchanged))
}
}
}
}
pub fn note_reconcile_failure(&self) {
self.metrics
.reconcile_failures
.fetch_add(1, Ordering::AcqRel);
}
pub fn reconcile_failures(&self) -> u64 {
self.metrics.reconcile_failures.load(Ordering::Acquire)
}
pub fn refusals(&self) -> (u64, u64) {
(
self.metrics
.refused_node_at_capacity
.load(Ordering::Acquire),
self.metrics
.refused_interest_at_capacity
.load(Ordering::Acquire),
)
}
pub fn release(&self, ticket: SensingLeaseTicket) -> LeaseAction {
let mut entries = self.entries.lock();
let Entry::Occupied(mut o) = entries.entry(ticket.key) else {
return LeaseAction::Unchanged;
};
let entry = o.get_mut();
if entry.registrations.remove(&ticket.token).is_none() {
return LeaseAction::Unchanged;
}
if entry.registrations.is_empty() {
let spec = Arc::clone(&entry.spec);
o.remove();
return LeaseAction::Deregister { spec };
}
let new_min = strictest_sample_interval(entry.registrations.values().copied())
.unwrap_or(entry.installed_interval);
if new_min > entry.installed_interval {
entry.installed_interval = new_min;
LeaseAction::Reregister {
spec: Arc::clone(&entry.spec),
interval: new_min,
}
} else {
LeaseAction::Unchanged
}
}
#[doc(hidden)]
#[cfg(any(test, feature = "fixtures"))]
pub fn len(&self) -> usize {
self.entries.lock().len()
}
#[doc(hidden)]
#[cfg(any(test, feature = "fixtures"))]
pub fn is_empty(&self) -> bool {
self.len() == 0
}
#[doc(hidden)]
#[cfg(any(test, feature = "fixtures"))]
pub fn entry_for_test(&self, key: &SensingLeaseKey) -> Option<(usize, Duration)> {
self.entries
.lock()
.get(key)
.map(|e| (e.registrations.len(), e.installed_interval))
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::adapter::net::behavior::sensing::identity::{
CanonicalConstraints, CapabilityId, DisclosureClass, ProviderSelector, ResultMode,
WorkLatencyEnvelope,
};
fn audience(byte: u8) -> AudienceScopeCommitment {
AudienceScopeCommitment::from_bytes([byte; 32])
}
fn spec(cap: &str) -> InterestSpec {
InterestSpec {
capability_id: CapabilityId::new(cap),
constraints: CanonicalConstraints::from_entries([("k", "v")]).unwrap(),
work_latency: WorkLatencyEnvelope::start_within(Duration::from_secs(2)),
providers: ProviderSelector::Node(7),
result_mode: ResultMode::Any,
disclosure_class: DisclosureClass::Owner,
audience: audience(1),
}
}
fn key_for(s: &InterestSpec, provider: u64) -> SensingLeaseKey {
SensingLeaseKey::ExactProvider {
audience: s.audience,
interest_digest: s.interest_digest(),
provider,
}
}
fn ticket(key: SensingLeaseKey, token: LeaseToken) -> SensingLeaseTicket {
SensingLeaseTicket { key, token }
}
fn ms(n: u64) -> Duration {
Duration::from_millis(n)
}
#[test]
fn the_two_hundred_and_fifty_seventh_interest_is_refused_without_evicting_any() {
let leases = SensingInterestLeases::default();
let s = spec("gpu.infer");
let mut held = Vec::new();
for provider in 0..MAX_LEASED_INTERESTS as u64 {
let key = key_for(&s, provider);
let (token, action) = leases.acquire(key, &s, ms(100)).expect("within capacity");
assert!(matches!(action, LeaseAction::Register { .. }));
held.push(ticket(key, token));
}
assert_eq!(leases.len(), MAX_LEASED_INTERESTS);
let overflow = key_for(&s, MAX_LEASED_INTERESTS as u64);
assert_eq!(
leases.acquire(overflow, &s, ms(100)),
Err(LeaseRefused::NodeAtCapacity)
);
assert_eq!(
leases.entry_for_test(&overflow),
None,
"a refused acquisition records nothing"
);
assert_eq!(
leases.len(),
MAX_LEASED_INTERESTS,
"and evicts nothing — the first 256 are intact"
);
assert_eq!(leases.refusals(), (1, 0));
let existing = held[0].key;
let (_t, action) = leases
.acquire(existing, &s, ms(500))
.expect("an existing interest is not node-bounded");
assert_eq!(action, LeaseAction::Unchanged);
assert!(matches!(
leases.release(held.pop().expect("held")),
LeaseAction::Deregister { .. }
));
assert!(leases.acquire(overflow, &s, ms(100)).is_ok());
}
#[test]
fn the_sixty_fifth_holder_is_refused_and_cannot_tighten_the_cadence() {
let leases = SensingInterestLeases::default();
let s = spec("gpu.infer");
let key = key_for(&s, 7);
for _ in 0..MAX_HOLDERS_PER_INTEREST {
leases.acquire(key, &s, ms(100)).expect("within capacity");
}
assert_eq!(
leases.entry_for_test(&key),
Some((MAX_HOLDERS_PER_INTEREST, ms(100)))
);
assert_eq!(
leases.acquire(key, &s, ms(10)),
Err(LeaseRefused::InterestAtCapacity)
);
assert_eq!(
leases.entry_for_test(&key),
Some((MAX_HOLDERS_PER_INTEREST, ms(100))),
"a refused holder neither joins nor tightens the installed cadence"
);
assert_eq!(leases.refusals(), (0, 1));
assert_eq!(
leases.len(),
1,
"and the refusal creates no second interest"
);
}
#[test]
fn first_acquire_registers_the_spec_at_its_interval() {
let leases = SensingInterestLeases::default();
let s = spec("gpu.infer");
let key = key_for(&s, 7);
let (_t, action) = leases.acquire(key, &s, ms(100)).expect("within capacity");
match action {
LeaseAction::Register { spec, interval } => {
assert_eq!(*spec, s);
assert_eq!(interval, ms(100));
}
other => panic!("expected Register, got {other:?}"),
}
assert_eq!(leases.entry_for_test(&key), Some((1, ms(100))));
}
#[test]
fn looser_second_acquire_is_unchanged() {
let leases = SensingInterestLeases::default();
let s = spec("gpu.infer");
let key = key_for(&s, 7);
leases.acquire(key, &s, ms(100)).expect("within capacity");
let (_t, action) = leases.acquire(key, &s, ms(500)).expect("within capacity");
assert_eq!(action, LeaseAction::Unchanged);
assert_eq!(leases.entry_for_test(&key), Some((2, ms(100))));
}
#[test]
fn stricter_second_acquire_reregisters_tighter() {
let leases = SensingInterestLeases::default();
let s = spec("gpu.infer");
let key = key_for(&s, 7);
leases.acquire(key, &s, ms(500)).expect("within capacity");
let (_t, action) = leases.acquire(key, &s, ms(100)).expect("within capacity");
match action {
LeaseAction::Reregister { spec, interval } => {
assert_eq!(*spec, s);
assert_eq!(interval, ms(100));
}
other => panic!("expected Reregister, got {other:?}"),
}
}
#[test]
fn releasing_a_non_strictest_holder_makes_no_wire_change() {
let leases = SensingInterestLeases::default();
let s = spec("gpu.infer");
let key = key_for(&s, 7);
let (strict, _) = leases.acquire(key, &s, ms(100)).expect("within capacity");
let (loose, _) = leases.acquire(key, &s, ms(500)).expect("within capacity");
let _ = strict;
let action = leases.release(ticket(key, loose));
assert_eq!(action, LeaseAction::Unchanged);
assert_eq!(leases.entry_for_test(&key), Some((1, ms(100))));
}
#[test]
fn releasing_the_strictest_holder_relaxes_the_cadence() {
let leases = SensingInterestLeases::default();
let s = spec("gpu.infer");
let key = key_for(&s, 7);
let (strict, _) = leases.acquire(key, &s, ms(100)).expect("within capacity");
leases.acquire(key, &s, ms(500)).expect("within capacity");
match leases.release(ticket(key, strict)) {
LeaseAction::Reregister { spec, interval } => {
assert_eq!(*spec, s);
assert_eq!(interval, ms(500));
}
other => panic!("expected Reregister, got {other:?}"),
}
assert_eq!(leases.entry_for_test(&key), Some((1, ms(500))));
}
#[test]
fn last_release_deregisters_and_drops_the_entry() {
let leases = SensingInterestLeases::default();
let s = spec("gpu.infer");
let key = key_for(&s, 7);
let (only, _) = leases.acquire(key, &s, ms(100)).expect("within capacity");
match leases.release(ticket(key, only)) {
LeaseAction::Deregister { spec } => assert_eq!(*spec, s),
other => panic!("expected Deregister, got {other:?}"),
}
assert!(leases.is_empty());
}
#[test]
fn equal_interval_holders_share_one_registration() {
let leases = SensingInterestLeases::default();
let s = spec("gpu.infer");
let key = key_for(&s, 7);
let (a, first) = leases.acquire(key, &s, ms(100)).expect("within capacity");
let (b, second) = leases.acquire(key, &s, ms(100)).expect("within capacity");
assert!(matches!(first, LeaseAction::Register { .. }));
assert_eq!(second, LeaseAction::Unchanged);
assert_eq!(leases.release(ticket(key, a)), LeaseAction::Unchanged);
assert!(matches!(
leases.release(ticket(key, b)),
LeaseAction::Deregister { .. }
));
assert!(leases.is_empty());
}
#[test]
fn releasing_an_unknown_or_repeated_token_is_a_noop() {
let leases = SensingInterestLeases::default();
let s = spec("gpu.infer");
let key = key_for(&s, 7);
let k2 = SensingLeaseKey::ExactProvider {
audience: audience(1),
interest_digest: s.interest_digest(),
provider: 9,
};
let (k1_tok, _) = leases.acquire(key, &s, ms(100)).expect("within capacity");
let (k2_tok, _) = leases.acquire(k2, &s, ms(100)).expect("within capacity");
assert_eq!(leases.release(ticket(key, k2_tok)), LeaseAction::Unchanged);
assert_eq!(leases.entry_for_test(&key), Some((1, ms(100))));
assert!(matches!(
leases.release(ticket(key, k1_tok)),
LeaseAction::Deregister { .. }
));
assert_eq!(leases.release(ticket(key, k1_tok)), LeaseAction::Unchanged);
}
#[test]
fn distinct_keys_never_alias() {
let leases = SensingInterestLeases::default();
let s = spec("gpu.infer");
let k1 = key_for(&s, 7);
let k2 = SensingLeaseKey::ExactProvider {
audience: s.audience,
interest_digest: s.interest_digest(),
provider: 8,
};
leases.acquire(k1, &s, ms(100)).expect("within capacity");
leases.acquire(k2, &s, ms(100)).expect("within capacity");
assert_eq!(leases.len(), 2);
assert_eq!(leases.entry_for_test(&k1), Some((1, ms(100))));
assert_eq!(leases.entry_for_test(&k2), Some((1, ms(100))));
}
}