use std::collections::BTreeMap;
use parking_lot::Mutex;
use crate::adapter::net::behavior::capability::MAX_CAPABILITY_HOPS;
use crate::adapter::net::behavior::org_scoped_ann::ScopedCapabilityAnnouncement;
use crate::adapter::net::identity::EntityId;
const RELAY_EXPIRY_SKEW_SECS: u64 = 300;
pub struct ScopedCapabilityRelayFrame<'a> {
pub hop_count: u8,
pub envelope: &'a [u8],
}
impl<'a> ScopedCapabilityRelayFrame<'a> {
pub const VERSION: u8 = 1;
const PREFIX_LEN: usize = 2;
pub fn encode(hop_count: u8, envelope: &[u8]) -> Vec<u8> {
let mut out = Vec::with_capacity(Self::PREFIX_LEN + envelope.len());
out.push(Self::VERSION);
out.push(hop_count);
out.extend_from_slice(envelope);
out
}
pub fn decode(bytes: &'a [u8]) -> Option<Self> {
if bytes.len() < Self::PREFIX_LEN + 1 || bytes[0] != Self::VERSION {
return None;
}
Some(Self {
hop_count: bytes[1],
envelope: &bytes[Self::PREFIX_LEN..],
})
}
}
#[derive(Clone, PartialEq, Eq, PartialOrd, Ord, Debug)]
pub struct RelayDedupKey {
pub provider: EntityId,
pub grant_id: [u8; 32],
pub audience_handle: [u8; 32],
pub generation: u64,
}
#[derive(Default)]
pub struct ScopedAnnRelayGate {
inner: Mutex<RelayGateInner>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RelayAdmission {
Fresh,
ShorterPath,
Drop,
}
impl RelayAdmission {
pub fn forwards(self) -> bool {
matches!(self, Self::Fresh | Self::ShorterPath)
}
pub fn ingests_locally(self) -> bool {
matches!(self, Self::Fresh)
}
}
#[derive(Debug)]
struct SeenEntry {
deadline: u64,
admitted_by: u64,
min_hop: u8,
}
#[derive(Default)]
struct RelayGateInner {
seen: BTreeMap<RelayDedupKey, SeenEntry>,
per_peer: BTreeMap<u64, usize>,
}
impl RelayGateInner {
fn sweep(&mut self, now_secs: u64) {
let per_peer = &mut self.per_peer;
self.seen.retain(|_, entry| {
if now_secs < entry.deadline {
return true;
}
if let Some(count) = per_peer.get_mut(&entry.admitted_by) {
*count = count.saturating_sub(1);
if *count == 0 {
per_peer.remove(&entry.admitted_by);
}
}
false
});
}
}
impl ScopedAnnRelayGate {
const MAX_ENTRIES: usize = 8192;
const MAX_ENTRIES_PER_PEER: usize = Self::MAX_ENTRIES / 8;
const RETENTION_SECS: u64 = 600;
pub fn new() -> Self {
Self::default()
}
pub fn admit(
&self,
from_node: u64,
key: RelayDedupKey,
hop_count: u8,
now_secs: u64,
) -> RelayAdmission {
let mut inner = self.inner.lock();
inner.sweep(now_secs);
if let Some(existing) = inner.seen.get_mut(&key) {
if hop_count < existing.min_hop {
existing.min_hop = hop_count;
return RelayAdmission::ShorterPath;
}
return RelayAdmission::Drop; }
if inner.per_peer.get(&from_node).copied().unwrap_or(0) >= Self::MAX_ENTRIES_PER_PEER {
return RelayAdmission::Drop;
}
if inner.seen.len() >= Self::MAX_ENTRIES {
return RelayAdmission::Drop;
}
inner.seen.insert(
key,
SeenEntry {
deadline: now_secs.saturating_add(Self::RETENTION_SECS),
admitted_by: from_node,
min_hop: hop_count,
},
);
*inner.per_peer.entry(from_node).or_insert(0) += 1;
RelayAdmission::Fresh
}
pub fn release(&self, key: &RelayDedupKey) {
let mut inner = self.inner.lock();
if let Some(entry) = inner.seen.remove(key) {
if let Some(count) = inner.per_peer.get_mut(&entry.admitted_by) {
*count = count.saturating_sub(1);
if *count == 0 {
inner.per_peer.remove(&entry.admitted_by);
}
}
}
}
pub(crate) fn len(&self) -> usize {
self.inner.lock().seen.len()
}
#[cfg(test)]
fn admit_fresh(&self, from_node: u64, key: RelayDedupKey, now_secs: u64) -> bool {
self.admit(from_node, key, 0, now_secs) == RelayAdmission::Fresh
}
}
pub struct ScopedRelayAdmit {
pub envelope: ScopedCapabilityAnnouncement,
pub forward: Option<Vec<u8>>,
pub dedup_key: RelayDedupKey,
pub ingest_locally: bool,
}
pub fn decide_scoped_relay(
frame_bytes: &[u8],
from_node: u64,
gate: &ScopedAnnRelayGate,
now_secs: u64,
) -> Option<ScopedRelayAdmit> {
if from_node == 0 {
return None;
}
let frame = ScopedCapabilityRelayFrame::decode(frame_bytes)?;
let envelope = ScopedCapabilityAnnouncement::from_bytes(frame.envelope).ok()?;
if frame.hop_count == 0 && envelope.provider().node_id() != from_node {
return None;
}
if envelope.expires_at().saturating_add(RELAY_EXPIRY_SKEW_SECS) <= now_secs {
return None;
}
let key = RelayDedupKey {
provider: envelope.provider().clone(),
grant_id: *envelope.grant_id(),
audience_handle: *envelope.audience_handle(),
generation: envelope.generation(),
};
let key_for_result = key.clone();
let admission = gate.admit(from_node, key, frame.hop_count, now_secs);
if !admission.forwards() {
return None;
}
let forward = if frame.hop_count < MAX_CAPABILITY_HOPS - 1 {
Some(ScopedCapabilityRelayFrame::encode(
frame.hop_count.saturating_add(1),
frame.envelope,
))
} else {
None
};
Some(ScopedRelayAdmit {
envelope,
forward,
dedup_key: key_for_result,
ingest_locally: admission.ingests_locally(),
})
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_released_identity_is_admissible_again_and_a_retained_one_is_not() {
let gate = ScopedAnnRelayGate::new();
const PEER: u64 = 9;
let key = key_n(1, 7);
assert_eq!(
gate.admit(PEER, key.clone(), 0, 1_000),
RelayAdmission::Fresh,
);
assert_eq!(
gate.admit(PEER, key.clone(), 0, 1_000),
RelayAdmission::Drop,
"without a release, a re-delivery is still a duplicate — this is \
the loop suppression the gate exists for",
);
gate.release(&key);
assert_eq!(gate.len(), 0, "the identity must leave the gate");
assert_eq!(
gate.admit(PEER, key.clone(), 0, 1_000),
RelayAdmission::Fresh,
"the same generation must be reconsidered after a retryable \
refusal, not swallowed for the retention horizon",
);
gate.release(&key);
for index in 0..ScopedAnnRelayGate::MAX_ENTRIES_PER_PEER as u64 {
assert_eq!(
gate.admit(PEER, key_n(index + 100, 1), 0, 1_000),
RelayAdmission::Fresh,
"peer budget leaked across release at index {index}",
);
}
}
#[test]
fn releasing_an_unknown_identity_is_harmless() {
let gate = ScopedAnnRelayGate::new();
gate.release(&key_n(42, 1));
assert_eq!(gate.len(), 0);
assert_eq!(gate.admit(7, key_n(42, 1), 0, 1_000), RelayAdmission::Fresh,);
}
#[test]
fn a_shorter_path_duplicate_is_re_forwarded() {
let gate = ScopedAnnRelayGate::new();
const MALICIOUS: u64 = 0xBAD;
const HONEST: u64 = 0x600D;
let key = key_n(1, 7);
assert_eq!(
gate.admit(MALICIOUS, key.clone(), MAX_CAPABILITY_HOPS - 1, 1_000),
RelayAdmission::Fresh,
);
let honest = gate.admit(HONEST, key.clone(), 0, 1_000);
assert_eq!(honest, RelayAdmission::ShorterPath);
assert!(
honest.forwards(),
"the shorter-path copy must be forwarded, or the subtree behind \
this node stays truncated for the whole generation",
);
assert!(
!honest.ingests_locally(),
"it must NOT be re-ingested — the store already holds this \
identity, and re-opening the AEAD would let a peer solicit \
repeated crypto work by replaying at ever-lower hops",
);
assert_eq!(
gate.admit(HONEST, key.clone(), 0, 1_000),
RelayAdmission::Drop,
);
assert_eq!(gate.admit(MALICIOUS, key, 5, 1_000), RelayAdmission::Drop);
}
#[test]
fn an_equal_or_worse_hop_duplicate_is_still_dropped() {
let gate = ScopedAnnRelayGate::new();
const PEER: u64 = 7;
let key = key_n(2, 3);
assert_eq!(
gate.admit(PEER, key.clone(), 4, 1_000),
RelayAdmission::Fresh
);
assert_eq!(
gate.admit(PEER, key.clone(), 4, 1_000),
RelayAdmission::Drop,
"equal hop is not an improvement",
);
assert_eq!(
gate.admit(PEER, key.clone(), 9, 1_000),
RelayAdmission::Drop,
"a worse hop is not an improvement",
);
assert_eq!(
gate.admit(PEER, key, 3, 1_000),
RelayAdmission::ShorterPath,
"a strictly better hop is",
);
}
fn provider_n(index: u64) -> EntityId {
let mut bytes = [0u8; 32];
bytes[..8].copy_from_slice(&index.to_le_bytes());
EntityId::from_bytes(bytes)
}
fn key_n(index: u64, generation: u64) -> RelayDedupKey {
RelayDedupKey {
provider: provider_n(index),
grant_id: [0u8; 32],
audience_handle: [0x11; 32],
generation,
}
}
#[test]
fn frame_round_trips_and_preserves_hop() {
let env = b"opaque-envelope-bytes";
let framed = ScopedCapabilityRelayFrame::encode(3, env);
let decoded = ScopedCapabilityRelayFrame::decode(&framed).expect("decode");
assert_eq!(decoded.hop_count, 3);
assert_eq!(decoded.envelope, env);
}
#[test]
fn frame_decode_is_strict() {
assert!(
ScopedCapabilityRelayFrame::decode(&[ScopedCapabilityRelayFrame::VERSION, 0]).is_none()
);
assert!(ScopedCapabilityRelayFrame::decode(&[]).is_none());
let mut framed = ScopedCapabilityRelayFrame::encode(0, b"x");
framed[0] = ScopedCapabilityRelayFrame::VERSION.wrapping_add(1);
assert!(ScopedCapabilityRelayFrame::decode(&framed).is_none());
}
const PEER: u64 = 0xA1;
#[test]
fn gate_admits_a_fresh_identity_once() {
let gate = ScopedAnnRelayGate::new();
assert!(
gate.admit_fresh(PEER, key_n(1, 7), 1_000),
"first sighting admits"
);
assert!(
!gate.admit_fresh(PEER, key_n(1, 7), 1_000),
"the identical identity is a duplicate"
);
assert!(
gate.admit_fresh(PEER, key_n(1, 8), 1_000),
"newer generation is fresh"
);
assert!(
gate.admit_fresh(PEER, key_n(2, 7), 1_000),
"different provider is fresh"
);
}
#[test]
fn dedup_is_global_across_ingress_peers() {
let gate = ScopedAnnRelayGate::new();
assert!(gate.admit_fresh(PEER, key_n(1, 7), 1_000));
assert!(
!gate.admit_fresh(PEER + 1, key_n(1, 7), 1_000),
"a second peer delivering the SAME identity is still a duplicate"
);
assert_eq!(gate.len(), 1);
}
#[test]
fn gate_expires_on_the_local_horizon() {
let gate = ScopedAnnRelayGate::new();
assert!(gate.admit_fresh(PEER, key_n(1, 7), 1_000));
assert!(!gate.admit_fresh(
PEER,
key_n(1, 7),
1_000 + ScopedAnnRelayGate::RETENTION_SECS - 1
));
assert!(gate.admit_fresh(
PEER,
key_n(1, 7),
1_000 + ScopedAnnRelayGate::RETENTION_SECS
));
}
#[test]
fn gate_is_bounded_fail_closed_and_never_evicts_active() {
let gate = ScopedAnnRelayGate::new();
let per_peer = ScopedAnnRelayGate::MAX_ENTRIES_PER_PEER;
for index in 0..ScopedAnnRelayGate::MAX_ENTRIES {
let peer = (index / per_peer) as u64;
assert!(gate.admit_fresh(peer, key_n(index as u64, 1), 1_000));
}
assert_eq!(gate.len(), ScopedAnnRelayGate::MAX_ENTRIES);
assert!(
!gate.admit_fresh(u64::MAX, key_n(u64::MAX, 1), 1_000),
"an unseen identity is refused when full"
);
assert_eq!(gate.len(), ScopedAnnRelayGate::MAX_ENTRIES);
assert!(!gate.admit_fresh(0, key_n(0, 1), 1_000));
assert!(gate.admit_fresh(
u64::MAX,
key_n(u64::MAX, 1),
1_000 + ScopedAnnRelayGate::RETENTION_SECS
));
}
#[test]
fn one_peer_cannot_starve_another_out_of_the_gate() {
let gate = ScopedAnnRelayGate::new();
const FLOODER: u64 = 0xBAD;
const HONEST: u64 = 0x600D;
let mut admitted = 0usize;
for index in 0..ScopedAnnRelayGate::MAX_ENTRIES as u64 {
if gate.admit_fresh(FLOODER, key_n(index, 1), 1_000) {
admitted += 1;
}
}
assert_eq!(
admitted,
ScopedAnnRelayGate::MAX_ENTRIES_PER_PEER,
"a single peer is capped at its own budget",
);
assert!(
gate.len() < ScopedAnnRelayGate::MAX_ENTRIES,
"the flooder must not have consumed the whole gate",
);
assert!(
gate.admit_fresh(HONEST, key_n(u64::MAX, 1), 1_000),
"an honest peer must still be admitted after a flood",
);
}
#[test]
fn per_peer_budget_is_returned_when_entries_expire() {
let gate = ScopedAnnRelayGate::new();
const PEER_A: u64 = 7;
for index in 0..ScopedAnnRelayGate::MAX_ENTRIES_PER_PEER as u64 {
assert!(gate.admit_fresh(PEER_A, key_n(index, 1), 1_000));
}
assert!(!gate.admit_fresh(PEER_A, key_n(u64::MAX, 1), 1_000));
let later = 1_000 + ScopedAnnRelayGate::RETENTION_SECS;
assert!(gate.admit_fresh(PEER_A, key_n(u64::MAX, 1), later));
assert_eq!(gate.len(), 1, "the expired window was fully reclaimed");
}
fn build_owner_frame(
hop: u8,
provider_seed: u8,
generation: u64,
expires_at: u64,
) -> (Vec<u8>, u64) {
use crate::adapter::net::behavior::org::{OrgKeypair, OrgMembershipCert};
use crate::adapter::net::behavior::org_authority::OwnerAudienceCredential;
use crate::adapter::net::identity::EntityKeypair;
let provider = EntityKeypair::from_bytes([provider_seed; 32]);
let org = OrgKeypair::from_bytes([1u8; 32]);
let credential = OwnerAudienceCredential::generate(org.org_id());
let cert = OrgMembershipCert::issue_at(
&org,
provider.entity_id().clone(),
5,
0,
1_000_000,
0x1234,
);
let envelope = ScopedCapabilityAnnouncement::build_owner(
&provider,
org.org_id(),
cert,
credential.audience_handle,
credential.discovery_key(),
generation,
expires_at,
b"owner-descriptor",
)
.expect("build owner envelope");
let provider_node = provider.entity_id().node_id();
(
ScopedCapabilityRelayFrame::encode(hop, &envelope.to_bytes()),
provider_node,
)
}
#[test]
fn decide_relay_rejects_the_unresolved_from_node() {
let gate = ScopedAnnRelayGate::new();
let (frame, _) = build_owner_frame(1, 0x20, 1, 1_000_000);
assert!(
decide_scoped_relay(&frame, 0, &gate, 1_000).is_none(),
"from_node == 0 (unresolved) is refused"
);
assert_eq!(gate.len(), 0, "a rejected frame never primes the gate");
assert!(decide_scoped_relay(&frame, 0xABCD, &gate, 1_000).is_some());
}
#[test]
fn decide_relay_drops_a_malformed_frame_without_priming_the_gate() {
let gate = ScopedAnnRelayGate::new();
let bad = ScopedCapabilityRelayFrame::encode(1, b"not-a-valid-envelope");
assert!(decide_scoped_relay(&bad, 0xABCD, &gate, 1_000).is_none());
assert!(decide_scoped_relay(&[0xFF, 0x00, 0x01], 0xABCD, &gate, 1_000).is_none());
assert_eq!(
gate.len(),
0,
"malformed / unsigned frames never prime the dedup gate"
);
}
#[test]
fn decide_relay_binds_the_direct_origin_at_hop_zero() {
let gate = ScopedAnnRelayGate::new();
let (frame, provider_node) = build_owner_frame(0, 0x21, 1, 1_000_000);
assert!(
decide_scoped_relay(&frame, provider_node ^ 1, &gate, 1_000).is_none(),
"a hop-0 frame from the wrong session peer is refused"
);
assert_eq!(gate.len(), 0);
assert!(decide_scoped_relay(&frame, provider_node, &gate, 1_000).is_some());
}
#[test]
fn decide_relay_rejects_an_expired_envelope() {
let gate = ScopedAnnRelayGate::new();
let (frame, _) = build_owner_frame(1, 0x22, 1, 1_000);
assert!(
decide_scoped_relay(&frame, 0xABCD, &gate, 1_000 + RELAY_EXPIRY_SKEW_SECS + 1)
.is_none(),
"an envelope dead past the relay skew is dropped"
);
assert_eq!(gate.len(), 0);
}
#[test]
fn decide_relay_admits_once_and_forwards_below_the_cap() {
let gate = ScopedAnnRelayGate::new();
let (frame, _) = build_owner_frame(1, 0x23, 1, 1_000_000);
let admit = decide_scoped_relay(&frame, 0xABCD, &gate, 1_000).expect("admitted");
let fwd = admit.forward.expect("forwarded below the cap");
let decoded = ScopedCapabilityRelayFrame::decode(&fwd).expect("decode forward");
assert_eq!(decoded.hop_count, 2, "hop incremented on forward");
let orig = ScopedCapabilityRelayFrame::decode(&frame).unwrap();
assert_eq!(decoded.envelope, orig.envelope);
assert!(
decide_scoped_relay(&frame, 0xABCD, &gate, 1_000).is_none(),
"a duplicate frame does not re-forward"
);
}
#[test]
fn decide_relay_admits_but_does_not_forward_at_the_hop_boundary() {
let gate = ScopedAnnRelayGate::new();
let (frame, _) = build_owner_frame(MAX_CAPABILITY_HOPS - 1, 0x24, 1, 1_000_000);
let admit =
decide_scoped_relay(&frame, 0xABCD, &gate, 1_000).expect("admitted at boundary");
assert!(
admit.forward.is_none(),
"a frame at the hop boundary is not forwarded"
);
}
}