use super::super::org::{OrgError, OrgId, OrgMembershipCert};
use super::super::org_authority::{NodeAuthority, OrgAuthorityError};
use super::super::org_revocation::{BarrieredGeneration, OrgRevocationState, OrgRevocationStore};
use super::frames::{FrameSpecError, SensingInterestFrame};
use super::identity::{AudienceScopeCommitment, InterestSpec};
use super::SensingCounters;
use crate::adapter::net::identity::EntityId;
use arc_swap::ArcSwapOption;
use parking_lot::Mutex;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use std::time::Duration;
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub struct OrgAuthorityView {
pub owner_org: OrgId,
pub verification_skew_secs: u64,
}
const ORG_SENSING_AUDIENCE_DOMAIN: &str = "net.sensing.org-audience.v1";
pub fn canonical_org_sensing_commitment(org_id: &OrgId) -> AudienceScopeCommitment {
let mut hasher = blake3::Hasher::new_derive_key(ORG_SENSING_AUDIENCE_DOMAIN);
hasher.update(org_id.as_bytes());
AudienceScopeCommitment::from_bytes(*hasher.finalize().as_bytes())
}
#[derive(Debug, PartialEq, Eq)]
pub struct GateProof(());
#[derive(Debug, PartialEq, Eq)]
enum ValidatedInner {
Capability {
spec: InterestSpec,
consumer: u64,
requested_sample_interval: Duration,
soft_state_ttl: Duration,
subscriber: EntityId,
org_id: OrgId,
gate_proof: GateProof,
},
Provider {
spec: InterestSpec,
target: u64,
requested_sample_interval: Duration,
soft_state_ttl: Duration,
subscriber: EntityId,
org_id: OrgId,
gate_proof: GateProof,
},
}
#[derive(Debug, PartialEq, Eq)]
pub struct ValidatedOrgSensingRegistration(ValidatedInner);
impl Clone for ValidatedOrgSensingRegistration {
fn clone(&self) -> Self {
Self(match &self.0 {
ValidatedInner::Capability {
spec,
consumer,
requested_sample_interval,
soft_state_ttl,
subscriber,
org_id,
..
} => ValidatedInner::Capability {
spec: spec.clone(),
consumer: *consumer,
requested_sample_interval: *requested_sample_interval,
soft_state_ttl: *soft_state_ttl,
subscriber: subscriber.clone(),
org_id: *org_id,
gate_proof: GateProof(()),
},
ValidatedInner::Provider {
spec,
target,
requested_sample_interval,
soft_state_ttl,
subscriber,
org_id,
..
} => ValidatedInner::Provider {
spec: spec.clone(),
target: *target,
requested_sample_interval: *requested_sample_interval,
soft_state_ttl: *soft_state_ttl,
subscriber: subscriber.clone(),
org_id: *org_id,
gate_proof: GateProof(()),
},
})
}
}
#[cfg(test)]
impl ValidatedOrgSensingRegistration {
pub(crate) fn capability_for_test(
spec: InterestSpec,
consumer: u64,
requested_sample_interval: Duration,
soft_state_ttl: Duration,
subscriber: EntityId,
org_id: OrgId,
) -> Self {
Self(ValidatedInner::Capability {
spec,
consumer,
requested_sample_interval,
soft_state_ttl,
subscriber,
org_id,
gate_proof: GateProof(()),
})
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum OrgSensingRejection {
NotOrgRegistration,
Semantic(FrameSpecError),
ConsumerBindingMismatch,
SenderMemberMismatch,
MissingAuthority,
ForeignOrg,
CertInvalid(OrgError),
BelowFloor,
AudienceMismatch,
}
#[allow(clippy::too_many_arguments)]
pub fn verify_org_sensing_registration(
frame: &SensingInterestFrame,
from_node: u64,
sender_entity: &EntityId,
node_authority: Option<OrgAuthorityView>,
revocation: &OrgRevocationState,
now_secs: u64,
counters: &SensingCounters,
) -> Result<ValidatedOrgSensingRegistration, OrgSensingRejection> {
let result = verify_org_sensing_registration_inner(
frame,
from_node,
sender_entity,
node_authority,
revocation,
now_secs,
counters,
);
if let Err(rejection) = &result {
let counter = match rejection {
OrgSensingRejection::CertInvalid(_) => Some(&counters.org_cert_invalid),
OrgSensingRejection::BelowFloor => Some(&counters.org_below_floor),
OrgSensingRejection::ForeignOrg => Some(&counters.org_foreign_org),
OrgSensingRejection::SenderMemberMismatch => Some(&counters.org_sender_member_mismatch),
OrgSensingRejection::AudienceMismatch => Some(&counters.org_audience_mismatch),
OrgSensingRejection::MissingAuthority => Some(&counters.org_authority_unavailable),
OrgSensingRejection::ConsumerBindingMismatch
| OrgSensingRejection::NotOrgRegistration => Some(&counters.protocol_invalid),
OrgSensingRejection::Semantic(_) => None,
};
if let Some(counter) = counter {
counter.fetch_add(1, Ordering::Relaxed);
}
}
result
}
#[allow(clippy::too_many_arguments)]
fn verify_org_sensing_registration_inner(
frame: &SensingInterestFrame,
from_node: u64,
sender_entity: &EntityId,
node_authority: Option<OrgAuthorityView>,
revocation: &OrgRevocationState,
now_secs: u64,
counters: &SensingCounters,
) -> Result<ValidatedOrgSensingRegistration, OrgSensingRejection> {
let (membership, leg) = match frame {
SensingInterestFrame::OrgCapabilityRegistration {
subscriber_membership,
consumer,
requested_sample_interval,
soft_state_ttl,
..
} => (
subscriber_membership,
Leg::Capability {
consumer: *consumer,
requested_sample_interval: *requested_sample_interval,
soft_state_ttl: *soft_state_ttl,
},
),
SensingInterestFrame::OrgProviderRegistration {
subscriber_membership,
target,
requested_sample_interval,
soft_state_ttl,
..
} => (
subscriber_membership,
Leg::Provider {
target: *target,
requested_sample_interval: *requested_sample_interval,
soft_state_ttl: *soft_state_ttl,
},
),
_ => return Err(OrgSensingRejection::NotOrgRegistration),
};
let spec = frame
.validated_spec(counters)
.map_err(OrgSensingRejection::Semantic)?;
if *sender_entity != membership.member {
return Err(OrgSensingRejection::SenderMemberMismatch);
}
if let Leg::Capability { consumer, .. } = &leg {
if *consumer != from_node {
return Err(OrgSensingRejection::ConsumerBindingMismatch);
}
}
let authority = node_authority.ok_or(OrgSensingRejection::MissingAuthority)?;
if membership.org_id != authority.owner_org {
return Err(OrgSensingRejection::ForeignOrg);
}
membership
.is_valid_at_with_skew(now_secs, authority.verification_skew_secs)
.map_err(OrgSensingRejection::CertInvalid)?;
if membership.generation < revocation.floor_for(&membership.org_id, &membership.member) {
return Err(OrgSensingRejection::BelowFloor);
}
if spec.audience != canonical_org_sensing_commitment(&membership.org_id) {
return Err(OrgSensingRejection::AudienceMismatch);
}
let subscriber = membership.member.clone();
let org_id = membership.org_id;
Ok(ValidatedOrgSensingRegistration(match leg {
Leg::Capability {
consumer,
requested_sample_interval,
soft_state_ttl,
} => ValidatedInner::Capability {
spec,
consumer,
requested_sample_interval,
soft_state_ttl,
subscriber,
org_id,
gate_proof: GateProof(()),
},
Leg::Provider {
target,
requested_sample_interval,
soft_state_ttl,
} => ValidatedInner::Provider {
spec,
target,
requested_sample_interval,
soft_state_ttl,
subscriber,
org_id,
gate_proof: GateProof(()),
},
}))
}
enum Leg {
Capability {
consumer: u64,
requested_sample_interval: Duration,
soft_state_ttl: Duration,
},
Provider {
target: u64,
requested_sample_interval: Duration,
soft_state_ttl: Duration,
},
}
#[allow(dead_code)]
#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct AdmittedSensingRegistration {
spec: InterestSpec,
leg: RegistrationLeg,
authority: RegistrationAuthority,
}
#[allow(dead_code)]
#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) enum RegistrationAuthority {
Legacy {
proven_root: AudienceScopeCommitment,
},
Org {
org_id: OrgId,
},
}
#[allow(dead_code)]
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) enum RegistrationLeg {
Capability {
consumer: u64,
requested_sample_interval: Duration,
soft_state_ttl: Duration,
},
Provider {
target: u64,
requested_sample_interval: Duration,
soft_state_ttl: Duration,
},
}
#[allow(dead_code)]
impl AdmittedSensingRegistration {
pub(crate) fn from_validated_legacy(spec: InterestSpec, leg: RegistrationLeg) -> Self {
let proven_root = spec.audience;
Self {
spec,
leg,
authority: RegistrationAuthority::Legacy { proven_root },
}
}
pub(crate) fn from_validated_org(value: ValidatedOrgSensingRegistration) -> Self {
match value.0 {
ValidatedInner::Capability {
spec,
consumer,
requested_sample_interval,
soft_state_ttl,
org_id,
..
} => Self {
spec,
leg: RegistrationLeg::Capability {
consumer,
requested_sample_interval,
soft_state_ttl,
},
authority: RegistrationAuthority::Org { org_id },
},
ValidatedInner::Provider {
spec,
target,
requested_sample_interval,
soft_state_ttl,
org_id,
..
} => Self {
spec,
leg: RegistrationLeg::Provider {
target,
requested_sample_interval,
soft_state_ttl,
},
authority: RegistrationAuthority::Org { org_id },
},
}
}
pub(crate) fn spec(&self) -> &InterestSpec {
&self.spec
}
pub(crate) fn leg(&self) -> RegistrationLeg {
self.leg
}
pub(crate) fn proven_root(&self) -> AudienceScopeCommitment {
match &self.authority {
RegistrationAuthority::Legacy { proven_root } => *proven_root,
RegistrationAuthority::Org { org_id } => canonical_org_sensing_commitment(org_id),
}
}
pub(crate) fn authority(&self) -> &RegistrationAuthority {
&self.authority
}
pub(crate) fn provider_continuation(
&self,
target: u64,
requested_sample_interval: Duration,
soft_state_ttl: Duration,
) -> Self {
Self {
spec: self.spec.clone(),
leg: RegistrationLeg::Provider {
target,
requested_sample_interval,
soft_state_ttl,
},
authority: self.authority.clone(),
}
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) struct SensingAuthorityStamp {
authority_ptr: usize,
store_ptr: usize,
store_generation: Option<BarrieredGeneration>,
installation_generation: u64,
poisoned: bool,
}
impl SensingAuthorityStamp {
pub(crate) fn is_current(&self, current: &SensingAuthorityStamp) -> bool {
self.store_generation.is_some()
&& current.store_generation.is_some()
&& self.installation_generation != u64::MAX
&& current.installation_generation != u64::MAX
&& self == current
&& !current.poisoned
}
}
pub(crate) struct SensingAuthoritySnapshot {
stamp: SensingAuthorityStamp,
authority_view: OrgAuthorityView,
floors: Arc<OrgRevocationState>,
_authority: Arc<NodeAuthority>,
_store: Arc<OrgRevocationStore>,
}
#[allow(dead_code)]
impl SensingAuthoritySnapshot {
pub(crate) fn stamp(&self) -> &SensingAuthorityStamp {
&self.stamp
}
pub(crate) fn authority_view(&self) -> OrgAuthorityView {
self.authority_view
}
pub(crate) fn floors(&self) -> &OrgRevocationState {
&self.floors
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) enum SensingAuthorityUnavailable {
NoAuthority,
NoStore,
Poisoned,
GenerationExhausted,
}
pub(crate) fn capture_sensing_authority_snapshot(
org_install: &Mutex<()>,
node_authority: &ArcSwapOption<NodeAuthority>,
org_revocation: &ArcSwapOption<OrgRevocationStore>,
org_install_generation: &AtomicU64,
) -> Result<SensingAuthoritySnapshot, SensingAuthorityUnavailable> {
let _install = org_install.lock();
let authority = node_authority
.load_full()
.ok_or(SensingAuthorityUnavailable::NoAuthority)?;
let store = org_revocation
.load_full()
.ok_or(SensingAuthorityUnavailable::NoStore)?;
if store.is_poisoned() {
return Err(SensingAuthorityUnavailable::Poisoned);
}
let (floors, store_generation) = store
.snapshot_with_generation()
.map_err(|_| SensingAuthorityUnavailable::GenerationExhausted)?;
let store_generation = Some(store_generation);
let stamp = SensingAuthorityStamp {
authority_ptr: Arc::as_ptr(&authority) as *const () as usize,
store_ptr: Arc::as_ptr(&store) as *const () as usize,
store_generation,
installation_generation: org_install_generation.load(Ordering::Acquire),
poisoned: false,
};
let authority_view = OrgAuthorityView {
owner_org: authority.owner_org(),
verification_skew_secs: authority.config.verification_skew_secs,
};
Ok(SensingAuthoritySnapshot {
stamp,
authority_view,
floors,
_authority: authority,
_store: store,
})
}
pub(crate) fn capture_current_sensing_stamp(
org_install: &Mutex<()>,
node_authority: &ArcSwapOption<NodeAuthority>,
org_revocation: &ArcSwapOption<OrgRevocationStore>,
org_install_generation: &AtomicU64,
) -> Option<SensingAuthorityStamp> {
let _install = org_install.lock();
let authority = node_authority.load_full()?;
let store = org_revocation.load_full()?;
Some(SensingAuthorityStamp {
authority_ptr: Arc::as_ptr(&authority) as *const () as usize,
store_ptr: Arc::as_ptr(&store) as *const () as usize,
store_generation: store.barriered_generation().ok(),
installation_generation: org_install_generation.load(Ordering::Acquire),
poisoned: store.is_poisoned(),
})
}
#[allow(dead_code)]
pub(crate) struct LiveOrgRelayMembership {
owner_cert: OrgMembershipCert,
org_id: OrgId,
_authority: Arc<NodeAuthority>,
_store: Arc<OrgRevocationStore>,
}
#[allow(dead_code)]
impl LiveOrgRelayMembership {
pub(crate) fn owner_cert(&self) -> &OrgMembershipCert {
&self.owner_cert
}
pub(crate) fn org_id(&self) -> OrgId {
self.org_id
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) enum RelayMembershipUnavailable {
NoAuthority,
NoStore,
Poisoned,
GenerationExhausted,
ForeignOrg,
NotForThisNode,
CertInvalid,
BelowFloor,
ViewChanged,
}
pub(crate) fn capture_live_org_relay_membership(
org_install: &Mutex<()>,
node_authority: &ArcSwapOption<NodeAuthority>,
org_revocation: &ArcSwapOption<OrgRevocationStore>,
local_entity: &EntityId,
expected_org: OrgId,
now_secs: u64,
) -> Result<LiveOrgRelayMembership, RelayMembershipUnavailable> {
capture_live_org_relay_membership_seamed(
org_install,
node_authority,
org_revocation,
local_entity,
expected_org,
now_secs,
|| {},
)
}
fn capture_live_org_relay_membership_seamed(
org_install: &Mutex<()>,
node_authority: &ArcSwapOption<NodeAuthority>,
org_revocation: &ArcSwapOption<OrgRevocationStore>,
local_entity: &EntityId,
expected_org: OrgId,
now_secs: u64,
after_floor_snapshot: impl FnOnce(),
) -> Result<LiveOrgRelayMembership, RelayMembershipUnavailable> {
let _install = org_install.lock();
let authority = node_authority
.load_full()
.ok_or(RelayMembershipUnavailable::NoAuthority)?;
let store = org_revocation
.load_full()
.ok_or(RelayMembershipUnavailable::NoStore)?;
if store.is_poisoned() {
return Err(RelayMembershipUnavailable::Poisoned);
}
if authority.owner_org() != expected_org {
return Err(RelayMembershipUnavailable::ForeignOrg);
}
let (floors, captured_generation) = store
.snapshot_with_generation()
.map_err(|_| RelayMembershipUnavailable::GenerationExhausted)?;
after_floor_snapshot();
let verification = authority
.config
.self_verify_at(local_entity, &floors, now_secs);
let current_generation = store
.barriered_generation()
.map_err(|_| RelayMembershipUnavailable::GenerationExhausted)?;
if store.is_poisoned() {
return Err(RelayMembershipUnavailable::Poisoned);
}
if current_generation != captured_generation {
return Err(RelayMembershipUnavailable::ViewChanged);
}
verification.map_err(map_self_verify_error)?;
Ok(LiveOrgRelayMembership {
owner_cert: authority.config.owner_cert.clone(),
org_id: authority.owner_org(),
_authority: authority,
_store: store,
})
}
fn map_self_verify_error(e: OrgAuthorityError) -> RelayMembershipUnavailable {
match e {
OrgAuthorityError::CertBelowFloor { .. } => RelayMembershipUnavailable::BelowFloor,
OrgAuthorityError::CertNotForThisNode { .. } => RelayMembershipUnavailable::NotForThisNode,
_ => RelayMembershipUnavailable::CertInvalid,
}
}
pub(crate) fn plan_provider_continuation(
admitted: &AdmittedSensingRegistration,
capture_membership: impl FnOnce(OrgId) -> Option<LiveOrgRelayMembership>,
) -> Option<SensingInterestFrame> {
let RegistrationLeg::Provider {
target,
requested_sample_interval: strictest,
soft_state_ttl: ttl,
} = admitted.leg()
else {
return None;
};
match admitted.authority() {
RegistrationAuthority::Legacy { .. } => Some(SensingInterestFrame::provider_registration(
admitted.spec(),
target,
strictest,
ttl,
)),
RegistrationAuthority::Org { org_id } => {
let membership = capture_membership(*org_id)?;
Some(SensingInterestFrame::org_provider_registration(
admitted.spec(),
target,
strictest,
ttl,
membership.owner_cert().clone(),
))
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::adapter::net::behavior::org::{
OrgKeypair, OrgMembershipCert, ORG_CERT_TTL_SECS_RECOMMENDED,
};
use crate::adapter::net::behavior::sensing::identity::{
CanonicalConstraints, CapabilityId, DisclosureClass, ProviderSelector, ResultMode,
WorkLatencyEnvelope,
};
use std::collections::BTreeMap;
const FROM_NODE: u64 = 0xA11CE;
const D: Duration = Duration::from_millis(100);
const TTL: Duration = Duration::from_secs(30);
fn now_secs() -> u64 {
crate::adapter::net::behavior::org::current_timestamp()
}
fn org_kp() -> OrgKeypair {
OrgKeypair::from_bytes([0x42u8; 32])
}
fn member() -> EntityId {
EntityId::from_bytes([0x24u8; 32])
}
fn cert_gen(generation: u32) -> OrgMembershipCert {
OrgMembershipCert::try_issue(
&org_kp(),
member(),
generation,
ORG_CERT_TTL_SECS_RECOMMENDED,
)
.expect("issue cert")
}
fn authority() -> OrgAuthorityView {
OrgAuthorityView {
owner_org: org_kp().org_id(),
verification_skew_secs: 60,
}
}
fn spec_with(audience: AudienceScopeCommitment) -> InterestSpec {
InterestSpec {
capability_id: CapabilityId::new("gpu.infer"),
constraints: CanonicalConstraints::from_entries([("model", "llama-70b")]).unwrap(),
work_latency: WorkLatencyEnvelope::start_within(Duration::from_secs(2)),
providers: ProviderSelector::Node(0x77),
result_mode: ResultMode::Any,
disclosure_class: DisclosureClass::Owner,
audience,
}
}
fn org_commit() -> AudienceScopeCommitment {
canonical_org_sensing_commitment(&org_kp().org_id())
}
fn cap_frame_with(
cert: OrgMembershipCert,
audience: AudienceScopeCommitment,
) -> SensingInterestFrame {
SensingInterestFrame::org_capability_registration(
&spec_with(audience),
D,
TTL,
FROM_NODE,
cert,
)
}
fn empty_floors() -> OrgRevocationState {
OrgRevocationState::default()
}
fn floors_at(org: OrgId, member: EntityId, floor: u32) -> OrgRevocationState {
let mut map = BTreeMap::new();
map.insert((org, member), floor);
OrgRevocationState::from_floors_for_test(map)
}
fn run(
frame: &SensingInterestFrame,
sender: &EntityId,
authority: Option<OrgAuthorityView>,
floors: &OrgRevocationState,
) -> Result<ValidatedOrgSensingRegistration, OrgSensingRejection> {
verify_org_sensing_registration(
frame,
FROM_NODE,
sender,
authority,
floors,
now_secs(),
&SensingCounters::default(),
)
}
#[test]
fn distinct_orgs_derive_distinct_commitments() {
let a = canonical_org_sensing_commitment(&OrgId([1u8; 32]));
let b = canonical_org_sensing_commitment(&OrgId([2u8; 32]));
assert_ne!(a, b);
assert_ne!(a.as_bytes(), &[1u8; 32]);
}
#[test]
fn a_valid_org_capability_registration_is_admitted() {
let frame = cap_frame_with(cert_gen(5), org_commit());
let validated =
run(&frame, &member(), Some(authority()), &empty_floors()).expect("admitted");
match validated.0 {
ValidatedInner::Capability {
consumer,
subscriber,
org_id,
..
} => {
assert_eq!(consumer, FROM_NODE);
assert_eq!(subscriber, member());
assert_eq!(org_id, org_kp().org_id());
}
other => panic!("expected Capability, got {other:?}"),
}
}
#[test]
fn a_valid_org_provider_registration_is_admitted() {
let frame = SensingInterestFrame::org_provider_registration(
&spec_with(org_commit()),
0x77,
D,
TTL,
cert_gen(5),
);
let validated =
run(&frame, &member(), Some(authority()), &empty_floors()).expect("admitted");
assert!(matches!(
validated.0,
ValidatedInner::Provider { target: 0x77, .. }
));
}
#[test]
fn a_non_org_frame_is_refused_as_not_org_registration() {
let frame =
SensingInterestFrame::provider_registration(&spec_with(org_commit()), 0x77, D, TTL);
let err = run(&frame, &member(), Some(authority()), &empty_floors()).unwrap_err();
assert!(
matches!(err, OrgSensingRejection::NotOrgRegistration),
"got {err:?}"
);
}
#[test]
fn a_digest_inconsistent_audience_on_an_org_frame_fails_at_step_2() {
let mut frame = SensingInterestFrame::org_provider_registration(
&spec_with(org_commit()),
0x77,
D,
TTL,
cert_gen(5),
);
if let SensingInterestFrame::OrgProviderRegistration { audience_scope, .. } = &mut frame {
*audience_scope = AudienceScopeCommitment::from_bytes([0xABu8; 32]);
}
let err = run(&frame, &member(), Some(authority()), &empty_floors()).unwrap_err();
assert!(
matches!(err, OrgSensingRejection::Semantic(_)),
"a digest-inconsistent org audience must fail at step 2, got {err:?}"
);
}
#[test]
fn a_foreign_org_refusal_bumps_its_counter() {
let counters = SensingCounters::default();
let foreign_kp = OrgKeypair::from_bytes([0x99u8; 32]);
let foreign_cert =
OrgMembershipCert::try_issue(&foreign_kp, member(), 5, ORG_CERT_TTL_SECS_RECOMMENDED)
.expect("foreign cert");
let frame = SensingInterestFrame::org_provider_registration(
&spec_with(canonical_org_sensing_commitment(&foreign_kp.org_id())),
0x77,
D,
TTL,
foreign_cert,
);
let err = verify_org_sensing_registration(
&frame,
FROM_NODE,
&member(),
Some(authority()),
&empty_floors(),
now_secs(),
&counters,
)
.unwrap_err();
assert!(
matches!(err, OrgSensingRejection::ForeignOrg),
"got {err:?}"
);
assert_eq!(
SensingCounters::get(&counters.org_foreign_org),
1,
"the foreign-org refusal is counted (previously silent)"
);
}
#[test]
fn foreign_organization_is_refused() {
let foreign_kp = OrgKeypair::from_bytes([0x99u8; 32]);
let foreign_cert =
OrgMembershipCert::try_issue(&foreign_kp, member(), 5, ORG_CERT_TTL_SECS_RECOMMENDED)
.expect("foreign cert");
let foreign_audience = canonical_org_sensing_commitment(&foreign_kp.org_id());
let frame = cap_frame_with(foreign_cert, foreign_audience);
assert_eq!(
run(&frame, &member(), Some(authority()), &empty_floors()),
Err(OrgSensingRejection::ForeignOrg)
);
}
#[test]
fn sender_member_mismatch_is_refused() {
let frame = cap_frame_with(cert_gen(5), org_commit());
let other_sender = EntityId::from_bytes([0xEEu8; 32]);
assert_eq!(
run(&frame, &other_sender, Some(authority()), &empty_floors()),
Err(OrgSensingRejection::SenderMemberMismatch)
);
}
#[test]
fn consumer_from_node_mismatch_is_refused() {
let frame = SensingInterestFrame::org_capability_registration(
&spec_with(org_commit()),
D,
TTL,
FROM_NODE ^ 0x1, cert_gen(5),
);
assert_eq!(
run(&frame, &member(), Some(authority()), &empty_floors()),
Err(OrgSensingRejection::ConsumerBindingMismatch)
);
}
#[test]
fn missing_authority_is_refused() {
let frame = cap_frame_with(cert_gen(5), org_commit());
assert_eq!(
run(&frame, &member(), None, &empty_floors()),
Err(OrgSensingRejection::MissingAuthority)
);
}
#[test]
fn generation_below_floor_is_refused() {
let frame = cap_frame_with(cert_gen(5), org_commit());
let floors = floors_at(org_kp().org_id(), member(), 6);
assert_eq!(
run(&frame, &member(), Some(authority()), &floors),
Err(OrgSensingRejection::BelowFloor)
);
}
#[test]
fn generation_at_floor_is_admitted() {
let frame = cap_frame_with(cert_gen(6), org_commit());
let floors = floors_at(org_kp().org_id(), member(), 6);
assert!(run(&frame, &member(), Some(authority()), &floors).is_ok());
}
#[test]
fn audience_org_mismatch_is_refused() {
let frame = cap_frame_with(
cert_gen(5),
AudienceScopeCommitment::from_bytes([0x55u8; 32]),
);
assert_eq!(
run(&frame, &member(), Some(authority()), &empty_floors()),
Err(OrgSensingRejection::AudienceMismatch)
);
}
#[test]
fn a_forged_signature_is_refused() {
let mut cert = cert_gen(5);
cert.signature[0] ^= 0xFF;
let frame = cap_frame_with(cert, org_commit());
assert!(matches!(
run(&frame, &member(), Some(authority()), &empty_floors()),
Err(OrgSensingRejection::CertInvalid(_))
));
}
#[test]
fn an_expired_certificate_is_refused() {
let now = now_secs();
let expired = OrgMembershipCert::issue_at(&org_kp(), member(), 5, now - 200, now - 100, 1);
let frame = cap_frame_with(expired, org_commit());
assert!(matches!(
run(&frame, &member(), Some(authority()), &empty_floors()),
Err(OrgSensingRejection::CertInvalid(_))
));
}
#[test]
fn a_not_yet_valid_certificate_is_refused() {
let now = now_secs();
let future = OrgMembershipCert::issue_at(&org_kp(), member(), 5, now + 100, now + 200, 1);
let frame = cap_frame_with(future, org_commit());
assert!(matches!(
run(&frame, &member(), Some(authority()), &empty_floors()),
Err(OrgSensingRejection::CertInvalid(_))
));
}
#[test]
fn org_admitted_wrapper_derives_proven_root_from_org_id() {
let validated = ValidatedOrgSensingRegistration::capability_for_test(
spec_with(org_commit()),
FROM_NODE,
D,
TTL,
member(),
org_kp().org_id(),
);
let admitted = AdmittedSensingRegistration::from_validated_org(validated);
assert_eq!(
admitted.proven_root(),
canonical_org_sensing_commitment(&org_kp().org_id())
);
assert_ne!(
admitted.proven_root(),
AudienceScopeCommitment::owner_root(&member())
);
assert!(matches!(
admitted.authority(),
RegistrationAuthority::Org { .. }
));
}
#[test]
fn legacy_admitted_wrapper_derives_the_proven_root_from_its_own_spec() {
let root = AudienceScopeCommitment::owner_root(&member());
let admitted = AdmittedSensingRegistration::from_validated_legacy(
spec_with(root),
RegistrationLeg::Capability {
consumer: FROM_NODE,
requested_sample_interval: D,
soft_state_ttl: TTL,
},
);
assert_eq!(
admitted.proven_root(),
root,
"the admitted root is the admitted spec's audience, by construction"
);
assert_eq!(
admitted.proven_root(),
admitted.spec().audience,
"and they cannot diverge — there is no second parameter to disagree with"
);
assert!(matches!(
admitted.authority(),
RegistrationAuthority::Legacy { .. }
));
}
fn stamp(
authority_ptr: usize,
store_ptr: usize,
store_generation: u64,
installation_generation: u64,
poisoned: bool,
) -> SensingAuthorityStamp {
SensingAuthorityStamp {
authority_ptr,
store_ptr,
store_generation: Some(BarrieredGeneration::from_raw_for_test(store_generation)),
installation_generation,
poisoned,
}
}
#[test]
fn identical_stamp_is_current() {
let captured = stamp(1, 2, 3, 4, false);
assert!(captured.is_current(&stamp(1, 2, 3, 4, false)));
}
#[test]
fn a_b_a_rotation_is_stale_via_installation_generation() {
let captured = stamp(1, 2, 3, 4, false);
let after_a_b_a = stamp(1, 2, 3, 6, false);
assert!(!captured.is_current(&after_a_b_a));
}
#[test]
fn floor_raise_on_same_store_is_stale() {
let captured = stamp(1, 2, 3, 4, false);
assert!(!captured.is_current(&stamp(1, 2, 4, 4, false)));
}
#[test]
fn authority_or_store_pointer_change_is_stale() {
let captured = stamp(1, 2, 3, 4, false);
assert!(!captured.is_current(&stamp(9, 2, 3, 4, false)));
assert!(!captured.is_current(&stamp(1, 9, 3, 4, false)));
}
#[test]
fn poison_transition_is_stale_even_with_equal_numeric_fields() {
let captured = stamp(1, 2, 3, 4, false);
assert!(!captured.is_current(&stamp(1, 2, 3, 4, true)));
}
#[test]
fn an_exhausted_generation_is_never_current_in_either_position() {
let base = stamp(1, 2, 3, 4, false);
let mut exhausted = base;
exhausted.store_generation = None;
assert!(
!base.is_current(&exhausted),
"a live view whose generation is exhausted cannot be shown current"
);
assert!(
!exhausted.is_current(&base),
"a captured view sampled at exhaustion cannot be shown still live"
);
assert!(
!exhausted.is_current(&exhausted),
"and two exhausted samples must NOT compare equal-and-current"
);
}
use crate::adapter::net::behavior::org::OrgRevocationBundle;
use crate::adapter::net::behavior::org_revocation::ProvisioningExpectation;
use crate::adapter::net::identity::EntityKeypair;
use std::sync::atomic::AtomicUsize;
fn scratch(tag: &str) -> std::path::PathBuf {
static SEQ: AtomicUsize = AtomicUsize::new(0);
let seq = SEQ.fetch_add(1, Ordering::Relaxed);
let dir = std::env::temp_dir().join(format!(
"net-relay-membership-{tag}-{}-{seq}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&dir);
dir
}
fn foreign_org() -> OrgId {
OrgKeypair::from_bytes([0x99u8; 32]).org_id()
}
fn adopt_relay(
tag: &str,
generation: u32,
) -> (EntityId, Arc<NodeAuthority>, Arc<OrgRevocationStore>) {
let kp = EntityKeypair::generate();
let entity = kp.entity_id().clone();
let cert = OrgMembershipCert::try_issue(
&org_kp(),
entity.clone(),
generation,
ORG_CERT_TTL_SECS_RECOMMENDED,
)
.expect("issue relay cert");
let authority = Arc::new(
NodeAuthority::adopt(&scratch(tag), cert, &entity, 60, None)
.expect("adopt relay authority"),
);
let store = authority.revocation.clone();
(entity, authority, store)
}
fn capture_relay(
authority: Option<Arc<NodeAuthority>>,
store: Option<Arc<OrgRevocationStore>>,
local_entity: &EntityId,
expected_org: OrgId,
now: u64,
) -> Result<LiveOrgRelayMembership, RelayMembershipUnavailable> {
let na = ArcSwapOption::from(authority);
let rev = ArcSwapOption::from(store);
let lock = Mutex::new(());
capture_live_org_relay_membership(&lock, &na, &rev, local_entity, expected_org, now)
}
#[test]
fn live_relay_membership_returns_this_nodes_own_cert() {
let (entity, authority, store) = adopt_relay("ok", 3);
let membership = capture_relay(
Some(authority),
Some(store),
&entity,
org_kp().org_id(),
now_secs(),
)
.expect("relay membership");
assert_eq!(membership.owner_cert().member, entity);
assert_eq!(membership.owner_cert().org_id, org_kp().org_id());
assert_eq!(membership.owner_cert().generation, 3);
assert_eq!(membership.org_id(), org_kp().org_id());
}
#[test]
fn no_authority_installed_is_refused() {
let (entity, _authority, store) = adopt_relay("noauth", 1);
assert_eq!(
capture_relay(None, Some(store), &entity, org_kp().org_id(), now_secs()).err(),
Some(RelayMembershipUnavailable::NoAuthority)
);
}
#[test]
fn no_store_installed_is_refused() {
let (entity, authority, _store) = adopt_relay("nostore", 1);
assert_eq!(
capture_relay(
Some(authority),
None,
&entity,
org_kp().org_id(),
now_secs()
)
.err(),
Some(RelayMembershipUnavailable::NoStore)
);
}
#[test]
fn poisoned_store_is_refused() {
let (entity, authority, store) = adopt_relay("poison", 1);
store.mark_poisoned_for_test();
assert_eq!(
capture_relay(
Some(authority),
Some(store),
&entity,
org_kp().org_id(),
now_secs()
)
.err(),
Some(RelayMembershipUnavailable::Poisoned)
);
}
#[test]
fn foreign_org_registration_is_refused() {
let (entity, authority, store) = adopt_relay("foreign", 1);
assert_eq!(
capture_relay(
Some(authority),
Some(store),
&entity,
foreign_org(),
now_secs()
)
.err(),
Some(RelayMembershipUnavailable::ForeignOrg)
);
}
#[test]
fn wrong_local_entity_is_not_for_this_node() {
let (_entity, authority, store) = adopt_relay("wrongentity", 1);
let other = EntityId::from_bytes([0xEEu8; 32]);
assert_eq!(
capture_relay(
Some(authority),
Some(store),
&other,
org_kp().org_id(),
now_secs()
)
.err(),
Some(RelayMembershipUnavailable::NotForThisNode)
);
}
#[test]
fn relay_expired_at_future_now_is_cert_invalid() {
let (entity, authority, store) = adopt_relay("expired", 1);
assert!(capture_relay(
Some(authority.clone()),
Some(store.clone()),
&entity,
org_kp().org_id(),
now_secs(),
)
.is_ok());
let far_future = now_secs() + ORG_CERT_TTL_SECS_RECOMMENDED + 1_000;
assert_eq!(
capture_relay(
Some(authority),
Some(store),
&entity,
org_kp().org_id(),
far_future
)
.err(),
Some(RelayMembershipUnavailable::CertInvalid)
);
}
#[test]
fn relay_below_live_floor_is_refused() {
let (entity, authority, _embedded) = adopt_relay("floored", 1);
let live = Arc::new(
OrgRevocationStore::init(scratch("floored-live"), ProvisioningExpectation::MayBeFresh)
.expect("init live store"),
);
let mut floors = BTreeMap::new();
floors.insert(entity.clone(), 2u32);
let bundle = OrgRevocationBundle::try_issue(&org_kp(), &floors).expect("bundle");
live.apply_bundle(&bundle).expect("apply floor raise");
assert_eq!(
capture_relay(
Some(authority),
Some(live),
&entity,
org_kp().org_id(),
now_secs()
)
.err(),
Some(RelayMembershipUnavailable::BelowFloor)
);
}
#[test]
fn relay_at_live_floor_is_admitted() {
let (entity, authority, store) = adopt_relay("atfloor", 2);
let mut floors = BTreeMap::new();
floors.insert(entity.clone(), 2u32);
let bundle = OrgRevocationBundle::try_issue(&org_kp(), &floors).expect("bundle");
store.apply_bundle(&bundle).expect("apply floor raise");
assert!(capture_relay(
Some(authority),
Some(store),
&entity,
org_kp().org_id(),
now_secs(),
)
.is_ok());
}
#[test]
fn foreign_org_precedes_cert_validity() {
let (entity, authority, store) = adopt_relay("orderfirst", 1);
let far_future = now_secs() + ORG_CERT_TTL_SECS_RECOMMENDED + 1_000;
assert_eq!(
capture_relay(
Some(authority),
Some(store),
&entity,
foreign_org(),
far_future
)
.err(),
Some(RelayMembershipUnavailable::ForeignOrg)
);
}
#[test]
fn floor_raise_during_capture_is_view_changed() {
use std::sync::Barrier;
let (entity, authority, store) = adopt_relay("viewchange", 1);
let na = ArcSwapOption::from(Some(authority));
let rev = ArcSwapOption::from(Some(store.clone()));
let lock = Mutex::new(());
let paused = Barrier::new(2);
let release = Barrier::new(2);
std::thread::scope(|s| {
let gate = s.spawn(|| {
capture_live_org_relay_membership_seamed(
&lock,
&na,
&rev,
&entity,
org_kp().org_id(),
now_secs(),
|| {
paused.wait();
release.wait();
},
)
});
paused.wait(); let mut floors = BTreeMap::new();
floors.insert(entity.clone(), 2u32);
let bundle = OrgRevocationBundle::try_issue(&org_kp(), &floors).expect("bundle");
store.apply_bundle(&bundle).expect("apply floor raise");
release.wait(); let result = gate.join().expect("gate thread");
assert_eq!(result.err(), Some(RelayMembershipUnavailable::ViewChanged));
});
}
#[test]
fn poison_during_capture_is_refused() {
use std::sync::Barrier;
let (entity, authority, store) = adopt_relay("poison-mid", 1);
let na = ArcSwapOption::from(Some(authority));
let rev = ArcSwapOption::from(Some(store.clone()));
let lock = Mutex::new(());
let paused = Barrier::new(2);
let release = Barrier::new(2);
std::thread::scope(|s| {
let gate = s.spawn(|| {
capture_live_org_relay_membership_seamed(
&lock,
&na,
&rev,
&entity,
org_kp().org_id(),
now_secs(),
|| {
paused.wait();
release.wait();
},
)
});
paused.wait(); store.mark_poisoned_for_test();
release.wait();
let result = gate.join().expect("gate thread");
assert_eq!(result.err(), Some(RelayMembershipUnavailable::Poisoned));
});
}
use crate::adapter::net::behavior::sensing::identity::ProviderInterestKey;
fn legacy_admitted(audience: AudienceScopeCommitment) -> AdmittedSensingRegistration {
AdmittedSensingRegistration::from_validated_legacy(
spec_with(audience),
RegistrationLeg::Provider {
target: 0x77,
requested_sample_interval: D,
soft_state_ttl: TTL,
},
)
}
fn org_admitted(consumer: EntityId) -> AdmittedSensingRegistration {
AdmittedSensingRegistration::from_validated_org(
ValidatedOrgSensingRegistration::capability_for_test(
spec_with(org_commit()),
FROM_NODE,
D,
TTL,
consumer,
org_kp().org_id(),
),
)
}
#[test]
fn legacy_provider_continuation_emits_legacy_frame() {
let admitted = legacy_admitted(AudienceScopeCommitment::owner_root(&member()));
let frame = plan_provider_continuation(&admitted, |_| {
panic!("capture must not run for a legacy admission")
})
.expect("legacy continuation frame");
assert_eq!(
frame,
SensingInterestFrame::provider_registration(admitted.spec(), 0x77, D, TTL)
);
}
#[test]
fn org_provider_continuation_carries_the_relays_own_cert() {
let consumer = EntityId::from_bytes([0xCCu8; 32]);
let seed = org_admitted(consumer.clone());
let admitted = seed.provider_continuation(0x77, D, TTL);
let (relay_entity, authority, store) = adopt_relay("continuation", 1);
let membership = capture_relay(
Some(authority),
Some(store),
&relay_entity,
org_kp().org_id(),
now_secs(),
)
.expect("relay membership");
let frame = plan_provider_continuation(&admitted, |org| {
assert_eq!(org, org_kp().org_id(), "captured for the admitted org");
Some(membership)
})
.expect("org continuation frame");
match frame {
SensingInterestFrame::OrgProviderRegistration {
subscriber_membership,
target,
..
} => {
assert_eq!(target, 0x77);
assert_eq!(
subscriber_membership.member, relay_entity,
"the continuation carries the relay's own cert"
);
assert_ne!(
subscriber_membership.member, consumer,
"never the downstream consumer's cert"
);
}
_ => panic!("expected OrgProviderRegistration"),
}
}
#[test]
fn org_continuation_without_membership_emits_nothing_and_no_fallback() {
let admitted =
org_admitted(EntityId::from_bytes([0xCCu8; 32])).provider_continuation(0x77, D, TTL);
let planned = plan_provider_continuation(&admitted, |_| None);
assert!(
planned.is_none(),
"an org admission with no live membership emits nothing (no legacy fallback)"
);
}
#[test]
fn capability_leg_cannot_enter_provider_planner() {
let org_cap = org_admitted(EntityId::from_bytes([0xCCu8; 32])); assert!(
plan_provider_continuation(&org_cap, |_| {
panic!("capture must not run for a non-provider leg")
})
.is_none(),
"an org capability seed is not a provider continuation"
);
let legacy_cap = AdmittedSensingRegistration::from_validated_legacy(
spec_with(AudienceScopeCommitment::owner_root(&member())),
RegistrationLeg::Capability {
consumer: FROM_NODE,
requested_sample_interval: D,
soft_state_ttl: TTL,
},
);
assert!(
plan_provider_continuation(&legacy_cap, |_| None).is_none(),
"a legacy capability seed is not a provider continuation"
);
}
#[test]
fn entity_legacy_root_and_org_commitment_are_cryptographically_separated() {
let org_spec = spec_with(org_commit());
let legacy_spec = spec_with(AudienceScopeCommitment::owner_root(&member()));
let org_key = ProviderInterestKey::new(org_spec.key(), 0x77);
let legacy_key = ProviderInterestKey::new(legacy_spec.key(), 0x77);
assert_ne!(
org_key, legacy_key,
"an entity owner-root audience and an org commitment must not collide on the key"
);
}
#[test]
fn provider_continuation_preserves_spec_and_authority() {
let legacy = legacy_admitted(AudienceScopeCommitment::owner_root(&member()));
let legacy_cont = legacy.provider_continuation(0x99, D, TTL);
assert_eq!(legacy_cont.spec(), legacy.spec(), "spec preserved");
assert_eq!(
legacy_cont.proven_root(),
legacy.proven_root(),
"legacy authority preserved"
);
assert!(matches!(
legacy_cont.authority(),
RegistrationAuthority::Legacy { .. }
));
assert_eq!(
legacy_cont.leg(),
RegistrationLeg::Provider {
target: 0x99,
requested_sample_interval: D,
soft_state_ttl: TTL,
},
"leg re-targeted to the provider"
);
let org = org_admitted(EntityId::from_bytes([0xCCu8; 32]));
let org_cont = org.provider_continuation(0x99, D, TTL);
assert_eq!(org_cont.spec(), org.spec(), "spec preserved");
assert!(matches!(
org_cont.authority(),
RegistrationAuthority::Org { .. }
));
assert_eq!(
org_cont.proven_root(),
canonical_org_sensing_commitment(&org_kp().org_id()),
"org authority preserved"
);
}
}