use std::collections::{BTreeMap, BTreeSet};
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::Arc;
use super::org::OrgId;
use super::org_grant::CapabilityAuthorityId;
use super::org_revocation::OrgRevocationState;
use super::org_scoped_ingest::{
CapabilityAudienceScope, PreparedScopedCapability, VerifiedScopedCapability,
};
use crate::adapter::net::identity::EntityId;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PrivateCapabilityProvider {
pub provider: EntityId,
pub owner_org: OrgId,
pub expires_at: u64,
pub generation: u64,
}
impl PrivateCapabilityProvider {
pub(crate) fn from_verified(c: &VerifiedScopedCapability) -> Self {
Self {
provider: c.provider().clone(),
owner_org: *c.owner_org(),
expires_at: c.expires_at(),
generation: c.generation(),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ScopedStoreOutcome {
Inserted,
Updated,
Stale,
RejectedPublic,
AtCapacity,
TooManyDeclarations,
}
const MAX_DECLARATIONS_PER_RECORD: usize = 64;
type ScopedKey = (CapabilityAudienceScope, EntityId);
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ScopedIngestReport {
pub outcome: ScopedStoreOutcome,
pub swept_live: Vec<ScopedKey>,
}
const MAX_DIRTY_CAPABILITIES: usize = 256;
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub(crate) enum DirtyCapabilities {
#[default]
Clean,
Caps(BTreeSet<CapabilityAuthorityId>),
RebuildAll,
}
impl DirtyCapabilities {
fn mark(&mut self, caps: &BTreeSet<CapabilityAuthorityId>) {
if caps.is_empty() {
return;
}
match self {
DirtyCapabilities::RebuildAll => {}
DirtyCapabilities::Clean => {
*self = if caps.len() > MAX_DIRTY_CAPABILITIES {
DirtyCapabilities::RebuildAll
} else {
DirtyCapabilities::Caps(caps.clone())
};
}
DirtyCapabilities::Caps(existing) => {
existing.extend(caps.iter().copied());
if existing.len() > MAX_DIRTY_CAPABILITIES {
*self = DirtyCapabilities::RebuildAll;
}
}
}
}
fn take(&mut self) -> DirtyCapabilities {
std::mem::take(self)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct PrivateDiscoveryChangeBatch {
pub generation: u64,
pub dirty: DirtyCapabilities,
}
struct StoredEntry {
generation: u64,
expires_at: u64,
tombstone_until: u64,
capability: Option<VerifiedScopedCapability>,
}
#[derive(Default)]
pub struct ScopedDiscoveryStore {
entries: BTreeMap<(CapabilityAudienceScope, EntityId), StoredEntry>,
scope_counts: BTreeMap<CapabilityAudienceScope, usize>,
}
fn release_scope_slot(
scope_counts: &mut BTreeMap<CapabilityAudienceScope, usize>,
scope: &CapabilityAudienceScope,
) {
if let Some(count) = scope_counts.get_mut(scope) {
*count -= 1;
if *count == 0 {
scope_counts.remove(scope);
}
}
}
impl ScopedDiscoveryStore {
pub fn new() -> Self {
Self::default()
}
const MAX_ENTRIES: usize = 8192;
const MAX_ENTRIES_PER_SCOPE: usize = 1024;
const OWNER_RESERVED_ENTRIES: usize = Self::MAX_ENTRIES_PER_SCOPE;
fn admission_ceiling(scope: &CapabilityAudienceScope) -> usize {
match scope {
CapabilityAudienceScope::Owner { .. } => Self::MAX_ENTRIES,
_ => Self::MAX_ENTRIES - Self::OWNER_RESERVED_ENTRIES,
}
}
fn entries_in_scope(&self, scope: &CapabilityAudienceScope) -> usize {
self.scope_counts.get(scope).copied().unwrap_or(0)
}
pub fn ingest(
&mut self,
capability: VerifiedScopedCapability,
now_secs: u64,
installed: &dyn InstalledConsumerGrants,
) -> ScopedIngestReport {
if matches!(capability.scope(), CapabilityAudienceScope::Public) {
return ScopedIngestReport {
outcome: ScopedStoreOutcome::RejectedPublic,
swept_live: Vec::new(),
};
}
let key = (capability.scope().clone(), capability.provider().clone());
let generation = capability.generation();
let expires_at = capability.expires_at();
match self.entries.get_mut(&key) {
Some(existing) if generation <= existing.generation => ScopedIngestReport {
outcome: ScopedStoreOutcome::Stale,
swept_live: Vec::new(),
},
Some(existing) => {
existing.generation = generation;
existing.expires_at = expires_at;
existing.tombstone_until = existing.tombstone_until.max(expires_at);
existing.capability = Some(capability);
ScopedIngestReport {
outcome: ScopedStoreOutcome::Updated,
swept_live: Vec::new(),
}
}
None => {
let mut swept_live = Vec::new();
let ceiling = Self::admission_ceiling(capability.scope());
if !self.reclaim_for_admission(ceiling, now_secs, installed, &mut swept_live) {
return ScopedIngestReport {
outcome: ScopedStoreOutcome::AtCapacity,
swept_live,
};
}
if self.entries_in_scope(capability.scope()) >= Self::MAX_ENTRIES_PER_SCOPE {
swept_live.append(&mut self.sweep_expired(now_secs));
if self.entries_in_scope(capability.scope()) >= Self::MAX_ENTRIES_PER_SCOPE {
return ScopedIngestReport {
outcome: ScopedStoreOutcome::AtCapacity,
swept_live,
};
}
}
*self.scope_counts.entry(key.0.clone()).or_insert(0) += 1;
self.entries.insert(
key,
StoredEntry {
generation,
expires_at,
tombstone_until: expires_at,
capability: Some(capability),
},
);
ScopedIngestReport {
outcome: ScopedStoreOutcome::Inserted,
swept_live,
}
}
}
}
fn reclaim_for_admission(
&mut self,
ceiling: usize,
now_secs: u64,
installed: &dyn InstalledConsumerGrants,
swept_live: &mut Vec<ScopedKey>,
) -> bool {
if self.entries.len() < ceiling {
return true;
}
swept_live.append(&mut self.sweep_expired(now_secs));
if self.entries.len() < ceiling {
return true;
}
for reclaim_live in [false, true] {
let victims: Vec<ScopedKey> = self
.entries
.iter()
.filter(|(_, entry)| entry.capability.is_some() == reclaim_live)
.filter(|(key, _)| dormant_grant_scope(&key.0, installed))
.map(|(key, _)| key.clone())
.collect();
for key in victims {
if self.entries.len() < ceiling {
return true;
}
if self.entries.remove(&key).is_some() {
release_scope_slot(&mut self.scope_counts, &key.0);
if reclaim_live {
swept_live.push(key);
}
}
}
}
self.entries.len() < ceiling
}
pub fn find_capabilities_for_grant<F>(
&self,
grant_id: &[u8; 32],
now_secs: u64,
floors: &OrgRevocationState,
mut predicate: F,
) -> Vec<&VerifiedScopedCapability>
where
F: FnMut(&VerifiedScopedCapability) -> bool,
{
self.entries
.values()
.filter(|e| now_secs < e.expires_at)
.filter_map(|e| e.capability.as_ref())
.filter(|c| {
matches!(
c.scope(),
CapabilityAudienceScope::Grant { grant_id: g, .. } if g == grant_id
)
})
.filter(|c| is_current(c, floors))
.filter(|c| predicate(c))
.collect()
}
pub fn find_owner_private_capabilities<F>(
&self,
now_secs: u64,
floors: &OrgRevocationState,
mut predicate: F,
) -> Vec<&VerifiedScopedCapability>
where
F: FnMut(&VerifiedScopedCapability) -> bool,
{
self.entries
.values()
.filter(|e| now_secs < e.expires_at)
.filter_map(|e| e.capability.as_ref())
.filter(|c| matches!(c.scope(), CapabilityAudienceScope::Owner { .. }))
.filter(|c| is_current(c, floors))
.filter(|c| predicate(c))
.collect()
}
pub fn sweep_expired(&mut self, now_secs: u64) -> Vec<ScopedKey> {
let mut swept = Vec::new();
let Self {
entries,
scope_counts,
} = self;
entries.retain(|key, e| {
if e.capability.is_some() && now_secs >= e.expires_at {
e.capability = None; swept.push(key.clone());
}
if now_secs < e.tombstone_until {
return true;
}
release_scope_slot(scope_counts, &key.0);
false
});
swept
}
fn live_record(&self, key: &ScopedKey) -> Option<&VerifiedScopedCapability> {
self.entries.get(key).and_then(|e| e.capability.as_ref())
}
pub fn len(&self) -> usize {
self.entries
.values()
.filter(|e| e.capability.is_some())
.count()
}
pub fn is_empty(&self) -> bool {
!self.entries.values().any(|e| e.capability.is_some())
}
}
pub trait InstalledConsumerGrants {
fn is_installed(&self, grant_id: &[u8; 32]) -> bool;
}
impl InstalledConsumerGrants for super::org_grant_registry::ConsumerGrantSnapshot {
fn is_installed(&self, grant_id: &[u8; 32]) -> bool {
self.get(grant_id).is_some()
}
}
pub struct NoConsumerGrants;
impl InstalledConsumerGrants for NoConsumerGrants {
fn is_installed(&self, _grant_id: &[u8; 32]) -> bool {
false
}
}
fn dormant_grant_scope(
scope: &CapabilityAudienceScope,
installed: &dyn InstalledConsumerGrants,
) -> bool {
match scope {
CapabilityAudienceScope::Grant { grant_id, .. } => !installed.is_installed(grant_id),
_ => false,
}
}
#[derive(Default)]
pub struct ScopedIngestCounters {
inserted: AtomicU64,
updated: AtomicU64,
stale: AtomicU64,
rejected_public: AtomicU64,
at_capacity: AtomicU64,
too_many_declarations: AtomicU64,
verify_refused: AtomicU64,
race_refused: AtomicU64,
}
impl ScopedIngestCounters {
pub fn note_outcome(&self, outcome: ScopedStoreOutcome) -> bool {
let counter = match outcome {
ScopedStoreOutcome::Inserted => &self.inserted,
ScopedStoreOutcome::Updated => &self.updated,
ScopedStoreOutcome::Stale => &self.stale,
ScopedStoreOutcome::RejectedPublic => &self.rejected_public,
ScopedStoreOutcome::TooManyDeclarations => &self.too_many_declarations,
ScopedStoreOutcome::AtCapacity => &self.at_capacity,
};
let previous = counter.fetch_add(1, Ordering::AcqRel);
matches!(outcome, ScopedStoreOutcome::AtCapacity) && previous % 1024 == 0
}
pub fn note_verify_refused(&self) {
self.verify_refused.fetch_add(1, Ordering::AcqRel);
}
pub fn note_race_refused(&self) {
self.race_refused.fetch_add(1, Ordering::AcqRel);
}
pub fn snapshot(&self) -> [u64; 8] {
[
self.inserted.load(Ordering::Acquire),
self.updated.load(Ordering::Acquire),
self.stale.load(Ordering::Acquire),
self.rejected_public.load(Ordering::Acquire),
self.at_capacity.load(Ordering::Acquire),
self.too_many_declarations.load(Ordering::Acquire),
self.verify_refused.load(Ordering::Acquire),
self.race_refused.load(Ordering::Acquire),
]
}
}
#[derive(Default)]
struct ScopedCapabilityIndex {
owner_by_capability: BTreeMap<CapabilityAuthorityId, BTreeSet<ScopedKey>>,
declarations_by_record: BTreeMap<ScopedKey, Arc<[CapabilityAuthorityId]>>,
floor_visible_by_provider: BTreeMap<EntityId, BTreeSet<ScopedKey>>,
}
impl ScopedCapabilityIndex {
fn insert_record(&mut self, key: ScopedKey, cap_ids: Arc<[CapabilityAuthorityId]>) {
if matches!(key.0, CapabilityAudienceScope::Owner { .. }) {
for cap in cap_ids.iter() {
self.owner_by_capability
.entry(*cap)
.or_default()
.insert(key.clone());
}
}
self.floor_visible_by_provider
.entry(key.1.clone())
.or_default()
.insert(key.clone());
self.declarations_by_record.insert(key, cap_ids);
}
fn remove_record(&mut self, key: &ScopedKey) {
let Some(cap_ids) = self.declarations_by_record.remove(key) else {
return;
};
if matches!(key.0, CapabilityAudienceScope::Owner { .. }) {
for cap in cap_ids.iter() {
if let Some(bucket) = self.owner_by_capability.get_mut(cap) {
bucket.remove(key);
if bucket.is_empty() {
self.owner_by_capability.remove(cap);
}
}
}
}
self.retract_floor_visibility(key);
}
fn retract_floor_visibility(&mut self, key: &ScopedKey) -> bool {
let Some(bucket) = self.floor_visible_by_provider.get_mut(&key.1) else {
return false;
};
let was_visible = bucket.remove(key);
if bucket.is_empty() {
self.floor_visible_by_provider.remove(&key.1);
}
was_visible
}
fn is_floor_visible(&self, key: &ScopedKey) -> bool {
self.floor_visible_by_provider
.get(&key.1)
.is_some_and(|bucket| bucket.contains(key))
}
fn replace_record(&mut self, key: ScopedKey, cap_ids: Arc<[CapabilityAuthorityId]>) {
self.remove_record(&key);
self.insert_record(key, cap_ids);
}
}
#[derive(Default)]
pub struct ScopedDiscoveryState {
store: ScopedDiscoveryStore,
index: ScopedCapabilityIndex,
revision: u64,
owner_revision: u64,
generations_exhausted: bool,
pending_global: DirtyCapabilities,
pending_owner: DirtyCapabilities,
live_expiries: BTreeMap<u64, u32>,
expiry_by_key: BTreeMap<ScopedKey, u64>,
drain_leases: PrivateDiscoveryLeaseState,
}
#[derive(Default)]
struct PrivateDiscoveryLeaseState {
global: Arc<AtomicBool>,
owner: Arc<AtomicBool>,
}
fn note_removed_record(
index: &ScopedCapabilityIndex,
key: &ScopedKey,
global: &mut BTreeSet<CapabilityAuthorityId>,
owner: &mut BTreeSet<CapabilityAuthorityId>,
) {
if !index.is_floor_visible(key) {
return;
}
if let Some(caps) = index.declarations_by_record.get(key) {
note_caps(&key.0, caps, global, owner);
}
}
fn note_caps(
scope: &CapabilityAudienceScope,
caps: &[CapabilityAuthorityId],
global: &mut BTreeSet<CapabilityAuthorityId>,
owner: &mut BTreeSet<CapabilityAuthorityId>,
) {
let is_owner = matches!(scope, CapabilityAudienceScope::Owner { .. });
for c in caps {
global.insert(*c);
if is_owner {
owner.insert(*c);
}
}
}
impl ScopedDiscoveryState {
pub fn new() -> Self {
Self::default()
}
pub fn ingest(
&mut self,
prepared: PreparedScopedCapability,
now_secs: u64,
installed: &dyn InstalledConsumerGrants,
) -> ScopedStoreOutcome {
let (capability, cap_ids) = prepared.into_parts();
if cap_ids.len() > MAX_DECLARATIONS_PER_RECORD {
return ScopedStoreOutcome::TooManyDeclarations;
}
let scope = capability.scope().clone();
let key = (scope.clone(), capability.provider().clone());
let expires_at = capability.expires_at();
let report = self.store.ingest(capability, now_secs, installed);
let mut global = BTreeSet::new();
let mut owner = BTreeSet::new();
for swept in &report.swept_live {
note_removed_record(&self.index, swept, &mut global, &mut owner);
self.index.remove_record(swept);
self.forget_live_expiry(swept);
}
match report.outcome {
ScopedStoreOutcome::Inserted => {
note_caps(&scope, &cap_ids, &mut global, &mut owner);
let declares = !cap_ids.is_empty();
self.index.insert_record(key.clone(), cap_ids);
if declares {
self.track_live_expiry(&key, expires_at);
}
}
ScopedStoreOutcome::Updated => {
note_removed_record(&self.index, &key, &mut global, &mut owner);
note_caps(&scope, &cap_ids, &mut global, &mut owner);
let declares = !cap_ids.is_empty();
self.index.replace_record(key.clone(), cap_ids);
if declares {
self.track_live_expiry(&key, expires_at);
} else {
self.forget_live_expiry(&key);
}
}
ScopedStoreOutcome::Stale
| ScopedStoreOutcome::RejectedPublic
| ScopedStoreOutcome::AtCapacity
| ScopedStoreOutcome::TooManyDeclarations => {}
}
self.record_change(&global, &owner);
report.outcome
}
pub fn sweep_expired(&mut self, now_secs: u64) -> usize {
let removed = self.store.sweep_expired(now_secs);
let dropped = removed.len();
let mut global = BTreeSet::new();
let mut owner = BTreeSet::new();
for key in &removed {
note_removed_record(&self.index, key, &mut global, &mut owner);
}
self.record_change(&global, &owner);
for key in &removed {
self.index.remove_record(key);
self.forget_live_expiry(key);
}
dropped
}
pub(crate) fn note_floors_raised(&mut self, raised: &[(OrgId, EntityId, u32)]) -> usize {
let mut retract: BTreeSet<ScopedKey> = BTreeSet::new();
for (org, provider, floor) in raised {
let Some(keys) = self.index.floor_visible_by_provider.get(provider) else {
continue;
};
for key in keys {
let Some(record) = self.store.live_record(key) else {
continue;
};
if record.owner_org() != org || record.provider_cert_generation() >= *floor {
continue;
}
retract.insert(key.clone());
}
}
let mut global = BTreeSet::new();
let mut owner = BTreeSet::new();
for key in &retract {
note_removed_record(&self.index, key, &mut global, &mut owner);
self.index.retract_floor_visibility(key);
self.forget_live_expiry(key);
}
self.record_change(&global, &owner);
retract.len()
}
#[cfg(test)]
pub(crate) fn advance_query_visible_generation_for_test(
&mut self,
capability: CapabilityAuthorityId,
) {
let mut global = BTreeSet::new();
global.insert(capability);
self.record_change(&global, &BTreeSet::new());
}
fn record_change(
&mut self,
global: &BTreeSet<CapabilityAuthorityId>,
owner: &BTreeSet<CapabilityAuthorityId>,
) {
if !global.is_empty() {
self.revision = self.advance_revision(self.revision);
self.pending_global.mark(global);
}
if !owner.is_empty() {
self.owner_revision = self.advance_revision(self.owner_revision);
self.pending_owner.mark(owner);
}
}
fn advance_revision(&mut self, current: u64) -> u64 {
match current.checked_add(1) {
Some(next) if next != u64::MAX => next,
_ => {
if !self.generations_exhausted {
self.generations_exhausted = true;
tracing::error!(
"org scoped discovery: change-generation space exhausted; private \
discovery is fenced rather than reusing a generation identity"
);
}
u64::MAX
}
}
}
pub fn generations_exhausted(&self) -> bool {
self.generations_exhausted
}
#[cfg(test)]
pub(crate) fn park_revisions_at_ceiling_for_test(&mut self) {
self.revision = u64::MAX - 1;
self.owner_revision = u64::MAX - 1;
}
fn track_live_expiry(&mut self, key: &ScopedKey, expires_at: u64) {
if let Some(previous) = self.expiry_by_key.insert(key.clone(), expires_at) {
self.release_expiry_slot(previous);
}
*self.live_expiries.entry(expires_at).or_insert(0) += 1;
}
fn forget_live_expiry(&mut self, key: &ScopedKey) {
if let Some(previous) = self.expiry_by_key.remove(key) {
self.release_expiry_slot(previous);
}
}
fn release_expiry_slot(&mut self, deadline: u64) {
if let Some(count) = self.live_expiries.get_mut(&deadline) {
*count -= 1;
if *count == 0 {
self.live_expiries.remove(&deadline);
}
}
}
pub fn next_visible_expiry(&self) -> Option<u64> {
self.live_expiries.keys().next().copied()
}
pub fn revision(&self) -> u64 {
self.revision
}
pub fn owner_revision(&self) -> u64 {
self.owner_revision
}
fn take_global_change_batch(&mut self) -> PrivateDiscoveryChangeBatch {
PrivateDiscoveryChangeBatch {
generation: self.revision,
dirty: self.pending_global.take(),
}
}
fn take_owner_change_batch(&mut self) -> PrivateDiscoveryChangeBatch {
PrivateDiscoveryChangeBatch {
generation: self.owner_revision,
dirty: self.pending_owner.take(),
}
}
fn mark_rebuild_all(&mut self, stream: PrivateDiscoveryStream) {
match stream {
PrivateDiscoveryStream::Global => self.pending_global = DirtyCapabilities::RebuildAll,
PrivateDiscoveryStream::Owner => self.pending_owner = DirtyCapabilities::RebuildAll,
}
}
pub fn find_owner_private_providers(
&self,
capability: Option<&CapabilityAuthorityId>,
now_secs: u64,
floors: &OrgRevocationState,
) -> Vec<(PrivateCapabilityProvider, EntityId)> {
let Some(cap) = capability else {
return self
.store
.find_owner_private_capabilities(now_secs, floors, |_| true)
.into_iter()
.map(|c| {
(
PrivateCapabilityProvider::from_verified(c),
c.provider().clone(),
)
})
.collect();
};
let mut out = Vec::new();
let Some(keys) = self.index.owner_by_capability.get(cap) else {
return out;
};
for key in keys {
let Some(rec) = self.store.live_record(key) else {
continue;
};
if now_secs < rec.expires_at() && is_current(rec, floors) {
out.push((
PrivateCapabilityProvider::from_verified(rec),
rec.provider().clone(),
));
}
}
out
}
pub(crate) fn find_scope_exact_private_providers(
&self,
scope: &CapabilityAudienceScope,
capability: &CapabilityAuthorityId,
now_secs: u64,
floors: &OrgRevocationState,
) -> Vec<PrivateCapabilityProvider> {
if !matches!(scope, CapabilityAudienceScope::Owner { .. }) {
return Vec::new();
}
let Some(keys) = self.index.owner_by_capability.get(capability) else {
return Vec::new();
};
keys.iter()
.filter(|(key_scope, _)| key_scope == scope)
.filter_map(|key| self.store.live_record(key))
.filter(|rec| now_secs < rec.expires_at() && is_current(rec, floors))
.map(PrivateCapabilityProvider::from_verified)
.collect()
}
pub fn find_capabilities_for_grant<F>(
&self,
grant_id: &[u8; 32],
now_secs: u64,
floors: &OrgRevocationState,
predicate: F,
) -> Vec<&VerifiedScopedCapability>
where
F: FnMut(&VerifiedScopedCapability) -> bool,
{
self.store
.find_capabilities_for_grant(grant_id, now_secs, floors, predicate)
}
pub(crate) fn find_grant_exact_private_providers<F>(
&self,
grant_id: &[u8; 32],
capability: &CapabilityAuthorityId,
now_secs: u64,
floors: &OrgRevocationState,
predicate: F,
) -> Vec<PrivateCapabilityProvider>
where
F: FnMut(&VerifiedScopedCapability) -> bool,
{
let declarations = &self.index.declarations_by_record;
self.store
.find_capabilities_for_grant(grant_id, now_secs, floors, predicate)
.into_iter()
.filter(|rec| {
declarations
.get(&(rec.scope().clone(), rec.provider().clone()))
.is_some_and(|caps| caps.contains(capability))
})
.map(PrivateCapabilityProvider::from_verified)
.collect()
}
pub fn len(&self) -> usize {
self.store.len()
}
pub fn is_empty(&self) -> bool {
self.store.is_empty()
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) enum PrivateDiscoveryStream {
Global,
#[allow(dead_code)]
Owner,
}
pub(crate) struct PrivateDiscoveryDrain {
state: Arc<parking_lot::Mutex<ScopedDiscoveryState>>,
stream: PrivateDiscoveryStream,
lease: Arc<AtomicBool>,
}
impl PrivateDiscoveryDrain {
pub(crate) fn drain(&mut self) -> PrivateDiscoveryChangeBatch {
let mut state = self.state.lock();
match self.stream {
PrivateDiscoveryStream::Global => state.take_global_change_batch(),
PrivateDiscoveryStream::Owner => state.take_owner_change_batch(),
}
}
}
impl Drop for PrivateDiscoveryDrain {
fn drop(&mut self) {
self.lease.store(false, Ordering::Release);
}
}
struct LeaseRollback<'a> {
lease: &'a AtomicBool,
armed: bool,
}
impl LeaseRollback<'_> {
fn disarm(mut self) {
self.armed = false;
}
}
impl Drop for LeaseRollback<'_> {
fn drop(&mut self) {
if self.armed {
self.lease.store(false, Ordering::Release);
}
}
}
pub(crate) struct PrivateDiscoveryDrains {
state: Arc<parking_lot::Mutex<ScopedDiscoveryState>>,
global_lease: Arc<AtomicBool>,
owner_lease: Arc<AtomicBool>,
}
impl PrivateDiscoveryDrains {
pub(crate) fn new(state: Arc<parking_lot::Mutex<ScopedDiscoveryState>>) -> Self {
let (global_lease, owner_lease) = {
let held = state.lock();
(
held.drain_leases.global.clone(),
held.drain_leases.owner.clone(),
)
};
Self {
state,
global_lease,
owner_lease,
}
}
fn lease_for(&self, stream: PrivateDiscoveryStream) -> &Arc<AtomicBool> {
match stream {
PrivateDiscoveryStream::Global => &self.global_lease,
PrivateDiscoveryStream::Owner => &self.owner_lease,
}
}
pub(crate) fn mint(&self, stream: PrivateDiscoveryStream) -> Option<PrivateDiscoveryDrain> {
let lease = self.lease_for(stream);
lease
.compare_exchange(false, true, Ordering::Acquire, Ordering::Relaxed)
.ok()?;
let rollback = LeaseRollback { lease, armed: true };
self.state.lock().mark_rebuild_all(stream);
let drain = PrivateDiscoveryDrain {
state: self.state.clone(),
stream,
lease: lease.clone(),
};
rollback.disarm();
Some(drain)
}
}
fn is_current(cap: &VerifiedScopedCapability, floors: &OrgRevocationState) -> bool {
floors.floor_for(cap.owner_org(), cap.provider()) <= cap.provider_cert_generation()
}
#[cfg(test)]
mod tests {
use super::*;
use crate::adapter::net::behavior::org::{OrgId, OrgKeypair, OrgRevocationBundle};
use std::collections::BTreeMap;
const FIXTURE_CERT_GEN: u32 = 5;
fn no_floors() -> OrgRevocationState {
OrgRevocationState::empty()
}
fn provider(seed: u8) -> EntityId {
EntityId::from_bytes([seed; 32])
}
fn org(seed: u8) -> OrgId {
OrgId::from_bytes([seed; 32])
}
fn owner_cap(provider_seed: u8, generation: u64, expires_at: u64) -> VerifiedScopedCapability {
VerifiedScopedCapability::for_test(
CapabilityAudienceScope::Owner {
org_id: org(1),
audience_handle: [0x11; 32],
},
provider(provider_seed),
org(1),
generation,
expires_at,
FIXTURE_CERT_GEN,
None,
b"owner-descriptor".to_vec(),
)
}
fn grant_cap(
grant_id: [u8; 32],
provider_seed: u8,
generation: u64,
expires_at: u64,
) -> VerifiedScopedCapability {
VerifiedScopedCapability::for_test(
CapabilityAudienceScope::Grant {
grant_id,
audience_handle: [0x22; 32],
},
provider(provider_seed),
org(2),
generation,
expires_at,
FIXTURE_CERT_GEN,
Some([0x5A; 64]),
b"grant-descriptor".to_vec(),
)
}
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 owner_cap_n(
provider_index: u64,
generation: u64,
expires_at: u64,
) -> VerifiedScopedCapability {
VerifiedScopedCapability::for_test(
CapabilityAudienceScope::Owner {
org_id: org(1),
audience_handle: [0x11; 32],
},
provider_n(provider_index),
org(1),
generation,
expires_at,
FIXTURE_CERT_GEN,
None,
b"owner-descriptor".to_vec(),
)
}
fn scoped_cap_in(
scope: CapabilityAudienceScope,
provider_index: u64,
generation: u64,
expires_at: u64,
) -> VerifiedScopedCapability {
VerifiedScopedCapability::for_test(
scope,
provider_n(provider_index),
org(1),
generation,
expires_at,
FIXTURE_CERT_GEN,
Some([0x5Au8; 64]),
b"granted-descriptor".to_vec(),
)
}
#[test]
fn ingest_bounds_cardinality_fail_closed_under_a_distinct_provider_flood() {
let mut store = ScopedDiscoveryStore::new();
let cap = ScopedDiscoveryStore::MAX_ENTRIES_PER_SCOPE;
for index in 0..cap as u64 {
assert_eq!(
store
.ingest(owner_cap_n(index, 1, 10_000), 1, &NoConsumerGrants)
.outcome,
ScopedStoreOutcome::Inserted
);
}
assert_eq!(store.len(), cap);
assert_eq!(
store
.ingest(owner_cap_n(u64::MAX, 1, 10_000), 1, &NoConsumerGrants)
.outcome,
ScopedStoreOutcome::AtCapacity
);
assert_eq!(store.len(), cap);
assert_eq!(
store
.ingest(owner_cap_n(0, 2, 10_000), 1, &NoConsumerGrants)
.outcome,
ScopedStoreOutcome::Updated
);
assert_eq!(store.len(), cap);
}
#[test]
fn one_exhausted_scope_never_denies_another() {
let mut store = ScopedDiscoveryStore::new();
let hostile = CapabilityAudienceScope::Grant {
grant_id: [0x7Au8; 32],
audience_handle: [0x7Bu8; 32],
};
for index in 0..ScopedDiscoveryStore::MAX_ENTRIES_PER_SCOPE as u64 {
assert_eq!(
store
.ingest(
scoped_cap_in(hostile.clone(), index, 1, 10_000),
1,
&AllInstalled
)
.outcome,
ScopedStoreOutcome::Inserted
);
}
assert_eq!(
store
.ingest(
scoped_cap_in(hostile.clone(), u64::MAX, 1, 10_000),
1,
&AllInstalled
)
.outcome,
ScopedStoreOutcome::AtCapacity,
"the hostile scope must be capped at its own share",
);
assert_eq!(
store
.ingest(owner_cap_n(0, 1, 10_000), 1, &NoConsumerGrants)
.outcome,
ScopedStoreOutcome::Inserted,
"a flooded grant scope must not deny owner-scoped discovery",
);
let other = CapabilityAudienceScope::Grant {
grant_id: [0x0Cu8; 32],
audience_handle: [0x0Du8; 32],
};
assert_eq!(
store
.ingest(scoped_cap_in(other, 0, 1, 10_000), 1, &NoConsumerGrants)
.outcome,
ScopedStoreOutcome::Inserted,
"a flooded grant scope must not deny an unrelated grant",
);
assert!(store.len() < ScopedDiscoveryStore::MAX_ENTRIES);
}
#[test]
fn intake_counters_separate_the_outcomes_and_announce_the_first_wedge() {
let counters = ScopedIngestCounters::default();
assert!(
counters.note_outcome(ScopedStoreOutcome::AtCapacity),
"the first capacity refusal must warn — a wedge announced late is a \
wedge an operator learns about from its consequences"
);
for _ in 1..1024 {
assert!(!counters.note_outcome(ScopedStoreOutcome::AtCapacity));
}
assert!(
counters.note_outcome(ScopedStoreOutcome::AtCapacity),
"and a SUSTAINED wedge re-announces rather than going quiet"
);
counters.note_outcome(ScopedStoreOutcome::Inserted);
counters.note_outcome(ScopedStoreOutcome::TooManyDeclarations);
counters.note_verify_refused();
counters.note_verify_refused();
counters.note_race_refused();
let counts = counters.snapshot();
assert_eq!(counts[0], 1, "inserted");
assert_eq!(counts[4], 1025, "at_capacity");
assert_eq!(counts[5], 1, "too_many_declarations");
assert_eq!(
counts[6], 2,
"verify_refused — a forged-envelope storm reads as a rate here"
);
assert_eq!(
counts[7], 1,
"race_refused — a persistent one means valid announcements never land"
);
assert_eq!(counts[1] + counts[2] + counts[3], 0, "and nothing bled");
}
#[test]
fn composed_grant_floods_never_deny_a_new_owner_key() {
let mut store = ScopedDiscoveryStore::new();
let scopes =
ScopedDiscoveryStore::MAX_ENTRIES / ScopedDiscoveryStore::MAX_ENTRIES_PER_SCOPE;
let mut provider_index = 0u64;
let mut admitted = 0usize;
for scope_index in 0..scopes as u8 {
let scope = CapabilityAudienceScope::Grant {
grant_id: [0xA0 ^ scope_index; 32],
audience_handle: [0xB0 ^ scope_index; 32],
};
for _ in 0..ScopedDiscoveryStore::MAX_ENTRIES_PER_SCOPE {
let outcome = store
.ingest(
scoped_cap_in(scope.clone(), provider_index, 1, 10_000),
1,
&AllInstalled,
)
.outcome;
provider_index += 1;
if matches!(outcome, ScopedStoreOutcome::Inserted) {
admitted += 1;
}
}
}
assert_eq!(
admitted,
ScopedDiscoveryStore::MAX_ENTRIES - ScopedDiscoveryStore::OWNER_RESERVED_ENTRIES,
"grant scopes COLLECTIVELY stop at the reservation boundary, not the global cap",
);
assert_eq!(
store
.ingest(owner_cap_n(0, 1, 10_000), 1, &NoConsumerGrants)
.outcome,
ScopedStoreOutcome::Inserted,
"a composed grant flood must not deny a NEW owner key",
);
let latecomer = CapabilityAudienceScope::Grant {
grant_id: [0xCC; 32],
audience_handle: [0xDD; 32],
};
assert_eq!(
store
.ingest(
scoped_cap_in(latecomer, u64::MAX, 1, 10_000),
1,
&AllInstalled
)
.outcome,
ScopedStoreOutcome::AtCapacity,
"an unrelated grant is refused at the reserve, with its own share empty",
);
}
struct AllInstalled;
impl InstalledConsumerGrants for AllInstalled {
fn is_installed(&self, _grant_id: &[u8; 32]) -> bool {
true
}
}
struct Installed(Vec<[u8; 32]>);
impl InstalledConsumerGrants for Installed {
fn is_installed(&self, grant_id: &[u8; 32]) -> bool {
self.0.contains(grant_id)
}
}
fn grant_scope(id: u8, handle: u8) -> CapabilityAudienceScope {
CapabilityAudienceScope::Grant {
grant_id: [id; 32],
audience_handle: [handle; 32],
}
}
#[test]
fn a_dormant_grants_rows_are_reclaimed_before_at_capacity() {
let mut store = ScopedDiscoveryStore::new();
let hostile_id = [0x7A; 32];
let live_id = [0x0C; 32];
let live = grant_scope(0x0C, 0x0D);
let non_owner_ceiling =
ScopedDiscoveryStore::MAX_ENTRIES - ScopedDiscoveryStore::OWNER_RESERVED_ENTRIES;
let mut provider_index = 0u64;
let mut handle = 0u8;
while store.entries.len() < non_owner_ceiling {
let scope = CapabilityAudienceScope::Grant {
grant_id: hostile_id,
audience_handle: [handle; 32],
};
for _ in 0..ScopedDiscoveryStore::MAX_ENTRIES_PER_SCOPE {
if store.entries.len() >= non_owner_ceiling {
break;
}
store.ingest(
scoped_cap_in(scope.clone(), provider_index, 1, u64::MAX),
1,
&AllInstalled,
);
provider_index += 1;
}
handle += 1;
}
let occupied = store.entries.len();
assert_eq!(occupied, non_owner_ceiling, "the non-owner budget is full");
assert_eq!(
store
.ingest(
scoped_cap_in(live.clone(), 9_000, 1, u64::MAX),
1,
&AllInstalled
)
.outcome,
ScopedStoreOutcome::AtCapacity,
"rows of an INSTALLED grant are never reclaimed for another key",
);
assert_eq!(store.entries.len(), occupied, "and nothing was taken");
assert_eq!(
store.entries.len(),
occupied,
"removal is a read-time filter — the rows are still stored",
);
let report = store.ingest(
scoped_cap_in(live, 9_000, 1, u64::MAX),
1,
&Installed(vec![live_id]),
);
assert_eq!(
report.outcome,
ScopedStoreOutcome::Inserted,
"a dormant grant's rows are reclaimed before AtCapacity",
);
assert_eq!(
report.swept_live.len(),
1,
"MINIMAL: exactly one dormant row was freed to admit one key; got {:?}",
report.swept_live.len(),
);
assert!(
report.swept_live.iter().all(|key| matches!(
&key.0,
CapabilityAudienceScope::Grant { grant_id, .. } if grant_id == &hostile_id
)),
"and only the dormant grant's rows were taken",
);
}
#[test]
fn owner_rows_are_never_reclaimed_for_a_grant_admission() {
let mut store = ScopedDiscoveryStore::new();
let dormant = grant_scope(0x11, 0x12);
for index in 0..ScopedDiscoveryStore::MAX_ENTRIES_PER_SCOPE as u64 {
store.ingest(owner_cap_n(index, 1, u64::MAX), 1, &NoConsumerGrants);
}
let owner_occupancy = store.len();
assert_eq!(owner_occupancy, ScopedDiscoveryStore::MAX_ENTRIES_PER_SCOPE);
store.ingest(
scoped_cap_in(dormant.clone(), 1, 1, u64::MAX),
1,
&NoConsumerGrants,
);
assert_eq!(
store.len(),
owner_occupancy + 1,
"owner rows are not candidates for reclamation at any pressure",
);
}
#[test]
fn tombstone_only_reclamation_reports_no_provider_set_change() {
let mut store = ScopedDiscoveryStore::new();
let dormant = grant_scope(0x21, 0x22);
let fresh = grant_scope(0x23, 0x24);
store.ingest(
scoped_cap_in(dormant.clone(), 1, 1, 10),
1,
&Installed(vec![[0x21; 32]]),
);
let report = store.ingest(scoped_cap_in(fresh, 2, 1, u64::MAX), 50, &NoConsumerGrants);
assert_eq!(report.outcome, ScopedStoreOutcome::Inserted);
assert!(
report.swept_live.is_empty(),
"an already-tombstoned row carries no live capability, so reclaiming \
it fabricates no provider-set change; got {:?}",
report.swept_live,
);
}
#[test]
fn capacity_pressure_never_rolls_back_a_known_high_water() {
let mut store = ScopedDiscoveryStore::new();
assert_eq!(
store
.ingest(owner_cap_n(0, 2, 10_000), 1, &NoConsumerGrants)
.outcome,
ScopedStoreOutcome::Inserted
);
for index in 1..ScopedDiscoveryStore::MAX_ENTRIES_PER_SCOPE as u64 {
store.ingest(owner_cap_n(index, 1, 10_000), 1, &NoConsumerGrants);
}
assert_eq!(store.len(), ScopedDiscoveryStore::MAX_ENTRIES_PER_SCOPE);
assert_eq!(
store
.ingest(owner_cap_n(u64::MAX, 1, 10_000), 1, &NoConsumerGrants)
.outcome,
ScopedStoreOutcome::AtCapacity
);
assert_eq!(
store
.ingest(owner_cap_n(0, 1, 10_000), 1, &NoConsumerGrants)
.outcome,
ScopedStoreOutcome::Stale
);
}
#[test]
fn ingest_reports_insert_update_and_stale() {
let mut store = ScopedDiscoveryStore::new();
assert_eq!(
store
.ingest(owner_cap(3, 1, 1000), 0, &NoConsumerGrants)
.outcome,
ScopedStoreOutcome::Inserted
);
assert_eq!(
store
.ingest(owner_cap(3, 2, 1000), 0, &NoConsumerGrants)
.outcome,
ScopedStoreOutcome::Updated
);
assert_eq!(
store
.ingest(owner_cap(3, 2, 1000), 0, &NoConsumerGrants)
.outcome,
ScopedStoreOutcome::Stale
);
assert_eq!(
store
.ingest(owner_cap(3, 1, 1000), 0, &NoConsumerGrants)
.outcome,
ScopedStoreOutcome::Stale
);
assert_eq!(store.len(), 1);
}
#[test]
fn public_scope_is_refused() {
let mut store = ScopedDiscoveryStore::new();
let public = VerifiedScopedCapability::for_test(
CapabilityAudienceScope::Public,
provider(3),
org(1),
1,
1000,
FIXTURE_CERT_GEN,
None,
b"x".to_vec(),
);
assert_eq!(
store.ingest(public, 0, &NoConsumerGrants).outcome,
ScopedStoreOutcome::RejectedPublic
);
assert!(store.is_empty());
}
#[test]
fn owner_and_grant_partitions_are_mutually_invisible() {
let mut store = ScopedDiscoveryStore::new();
let grant_x = [0xAA; 32];
let grant_y = [0xBB; 32];
store.ingest(owner_cap(3, 1, 1000), 0, &NoConsumerGrants);
store.ingest(grant_cap(grant_x, 4, 1, 1000), 0, &NoConsumerGrants);
store.ingest(grant_cap(grant_y, 5, 1, 1000), 0, &NoConsumerGrants);
assert_eq!(store.len(), 3);
let x = store.find_capabilities_for_grant(&grant_x, 0, &no_floors(), |_| true);
assert_eq!(x.len(), 1);
assert_eq!(x[0].provider(), &provider(4));
let y = store.find_capabilities_for_grant(&grant_y, 0, &no_floors(), |_| true);
assert_eq!(y.len(), 1);
assert_eq!(y[0].provider(), &provider(5));
let owner = store.find_owner_private_capabilities(0, &no_floors(), |_| true);
assert_eq!(owner.len(), 1);
assert_eq!(owner[0].provider(), &provider(3));
assert!(store
.find_capabilities_for_grant(&[0xCC; 32], 0, &no_floors(), |_| true)
.is_empty());
}
#[test]
fn predicate_filters_within_a_partition() {
let mut store = ScopedDiscoveryStore::new();
let grant = [0xAA; 32];
store.ingest(grant_cap(grant, 4, 1, 1000), 0, &NoConsumerGrants);
store.ingest(grant_cap(grant, 5, 1, 1000), 0, &NoConsumerGrants);
let hits = store
.find_capabilities_for_grant(&grant, 0, &no_floors(), |c| c.provider() == &provider(5));
assert_eq!(hits.len(), 1);
assert_eq!(hits[0].provider(), &provider(5));
}
#[test]
fn distinct_providers_under_one_grant_coexist() {
let mut store = ScopedDiscoveryStore::new();
let grant = [0xAA; 32];
store.ingest(grant_cap(grant, 4, 1, 1000), 0, &NoConsumerGrants);
store.ingest(grant_cap(grant, 5, 1, 1000), 0, &NoConsumerGrants);
assert_eq!(
store
.find_capabilities_for_grant(&grant, 0, &no_floors(), |_| true)
.len(),
2
);
}
#[test]
fn sweep_removes_only_expired_entries() {
let mut store = ScopedDiscoveryStore::new();
store.ingest(owner_cap(3, 1, 1000), 0, &NoConsumerGrants); store.ingest(grant_cap([0xAA; 32], 4, 1, 5000), 0, &NoConsumerGrants); assert_eq!(store.sweep_expired(2000).len(), 1);
assert_eq!(store.len(), 1);
assert!(store
.find_owner_private_capabilities(2000, &no_floors(), |_| true)
.is_empty());
assert_eq!(
store
.find_capabilities_for_grant(&[0xAA; 32], 2000, &no_floors(), |_| true)
.len(),
1
);
}
#[test]
fn queries_exclude_expired_entries_before_any_sweep() {
let mut store = ScopedDiscoveryStore::new();
let grant = [0xAA; 32];
store.ingest(grant_cap(grant, 4, 1, 1000), 0, &NoConsumerGrants); assert_eq!(
store
.find_capabilities_for_grant(&grant, 500, &no_floors(), |_| true)
.len(),
1,
"visible before expiry"
);
assert!(
store
.find_capabilities_for_grant(&grant, 2000, &no_floors(), |_| true)
.is_empty(),
"excluded past expiry even with no sweep",
);
}
#[test]
fn a_swept_newer_generation_cannot_be_revived_by_an_older_one() {
let mut store = ScopedDiscoveryStore::new();
let grant = [0xAA; 32];
store.ingest(grant_cap(grant, 4, 1, 5000), 0, &NoConsumerGrants); assert_eq!(
store
.ingest(grant_cap(grant, 4, 2, 2000), 0, &NoConsumerGrants)
.outcome, ScopedStoreOutcome::Updated
);
store.sweep_expired(3000);
assert!(store
.find_capabilities_for_grant(&grant, 3000, &no_floors(), |_| true)
.is_empty());
assert_eq!(
store
.ingest(grant_cap(grant, 4, 1, 5000), 0, &NoConsumerGrants)
.outcome,
ScopedStoreOutcome::Stale
);
assert!(store
.find_capabilities_for_grant(&grant, 3000, &no_floors(), |_| true)
.is_empty());
}
fn floor_state(org_kp: &OrgKeypair, member: &EntityId, floor: u32) -> OrgRevocationState {
let mut floors_map = BTreeMap::new();
floors_map.insert(member.clone(), floor);
let bundle = OrgRevocationBundle::try_issue(org_kp, &floors_map).expect("issue bundle");
let mut state = OrgRevocationState::empty();
state.merge_bundle(&bundle);
state
}
#[test]
fn a_raised_provider_floor_retracts_a_stored_record_at_query_time() {
let org_kp = OrgKeypair::from_bytes([7u8; 32]);
let org_id = org_kp.org_id();
let member = EntityId::from_bytes([9u8; 32]);
let mut store = ScopedDiscoveryStore::new();
store.ingest(
VerifiedScopedCapability::for_test(
CapabilityAudienceScope::Owner {
org_id,
audience_handle: [0x11; 32],
},
member.clone(),
org_id,
1,
10_000,
FIXTURE_CERT_GEN,
None,
b"owner-descriptor".to_vec(),
),
0,
&NoConsumerGrants,
);
assert_eq!(
store
.find_owner_private_capabilities(0, &no_floors(), |_| true)
.len(),
1
);
let floor_at = floor_state(&org_kp, &member, FIXTURE_CERT_GEN);
assert_eq!(
store
.find_owner_private_capabilities(0, &floor_at, |_| true)
.len(),
1,
"a floor equal to the admitted generation keeps the record"
);
let floor_above = floor_state(&org_kp, &member, FIXTURE_CERT_GEN + 1);
assert!(
store
.find_owner_private_capabilities(0, &floor_above, |_| true)
.is_empty(),
"a floor above the admitted generation retracts the record at query time"
);
assert_eq!(store.len(), 1);
}
use crate::adapter::net::behavior::capability::CapabilitySet;
use crate::adapter::net::behavior::org_scoped_ingest::PreparedScopedCapability;
fn owner_scope() -> CapabilityAudienceScope {
CapabilityAudienceScope::Owner {
org_id: org(1),
audience_handle: [0x11; 32],
}
}
fn cap_id(tag: &str) -> CapabilityAuthorityId {
CapabilityAuthorityId::for_tag(tag)
}
fn descriptor(tags: &[&str]) -> Vec<u8> {
let mut caps = CapabilitySet::new();
for t in tags {
caps = caps.add_tag(*t);
}
caps.to_bytes_compact()
}
fn owner_cap_declaring(
provider_seed: u8,
generation: u64,
expires_at: u64,
tags: &[&str],
) -> VerifiedScopedCapability {
VerifiedScopedCapability::for_test(
owner_scope(),
provider(provider_seed),
org(1),
generation,
expires_at,
FIXTURE_CERT_GEN,
None,
descriptor(tags),
)
}
fn owner_cap_declaring_n(
provider_index: u64,
generation: u64,
expires_at: u64,
tags: &[&str],
) -> VerifiedScopedCapability {
VerifiedScopedCapability::for_test(
owner_scope(),
provider_n(provider_index),
org(1),
generation,
expires_at,
FIXTURE_CERT_GEN,
None,
descriptor(tags),
)
}
fn ingest_indexed(
state: &mut ScopedDiscoveryState,
cap: VerifiedScopedCapability,
now: u64,
) -> ScopedStoreOutcome {
state.ingest(
PreparedScopedCapability::prepare(cap),
now,
&NoConsumerGrants,
)
}
#[test]
fn indexed_owner_query_matches_only_the_declared_capability() {
let mut state = ScopedDiscoveryState::new();
ingest_indexed(
&mut state,
owner_cap_declaring(3, 1, 10_000, &["nrpc:a"]),
0,
);
ingest_indexed(
&mut state,
owner_cap_declaring(4, 1, 10_000, &["nrpc:b"]),
0,
);
let a = state.find_owner_private_providers(Some(&cap_id("nrpc:a")), 0, &no_floors());
assert_eq!(a.len(), 1);
assert_eq!(a[0].0.provider, provider(3));
let b = state.find_owner_private_providers(Some(&cap_id("nrpc:b")), 0, &no_floors());
assert_eq!(b.len(), 1);
assert_eq!(b[0].0.provider, provider(4));
assert!(state
.find_owner_private_providers(Some(&cap_id("nrpc:none")), 0, &no_floors())
.is_empty());
}
#[test]
fn a_multi_tag_owner_record_is_indexed_under_every_capability() {
let mut state = ScopedDiscoveryState::new();
ingest_indexed(
&mut state,
owner_cap_declaring(3, 1, 10_000, &["nrpc:a", "nrpc:b"]),
0,
);
assert_eq!(
state
.find_owner_private_providers(Some(&cap_id("nrpc:a")), 0, &no_floors())
.len(),
1
);
assert_eq!(
state
.find_owner_private_providers(Some(&cap_id("nrpc:b")), 0, &no_floors())
.len(),
1
);
}
#[test]
fn an_updated_descriptor_reindexes_the_record() {
let mut state = ScopedDiscoveryState::new();
ingest_indexed(
&mut state,
owner_cap_declaring(3, 1, 10_000, &["nrpc:a"]),
0,
);
assert_eq!(
state
.find_owner_private_providers(Some(&cap_id("nrpc:a")), 0, &no_floors())
.len(),
1
);
assert_eq!(
ingest_indexed(
&mut state,
owner_cap_declaring(3, 2, 10_000, &["nrpc:b"]),
0
),
ScopedStoreOutcome::Updated
);
assert!(
state
.find_owner_private_providers(Some(&cap_id("nrpc:a")), 0, &no_floors())
.is_empty(),
"the old capability is re-indexed away"
);
assert_eq!(
state
.find_owner_private_providers(Some(&cap_id("nrpc:b")), 0, &no_floors())
.len(),
1,
"the new capability is indexed"
);
}
#[test]
fn the_indexed_owner_query_excludes_an_expired_record_before_sweep() {
let mut state = ScopedDiscoveryState::new();
ingest_indexed(&mut state, owner_cap_declaring(3, 1, 1000, &["nrpc:a"]), 0);
let a = cap_id("nrpc:a");
assert_eq!(
state
.find_owner_private_providers(Some(&a), 500, &no_floors())
.len(),
1,
"visible before expiry"
);
assert!(
state
.find_owner_private_providers(Some(&a), 2000, &no_floors())
.is_empty(),
"excluded past expiry with no sweep"
);
}
#[test]
fn the_indexed_owner_query_applies_floor_currentness() {
let org_kp = OrgKeypair::from_bytes([7u8; 32]);
let org_id = org_kp.org_id();
let member = EntityId::from_bytes([9u8; 32]);
let mut state = ScopedDiscoveryState::new();
ingest_indexed(
&mut state,
VerifiedScopedCapability::for_test(
CapabilityAudienceScope::Owner {
org_id,
audience_handle: [0x11; 32],
},
member.clone(),
org_id,
1,
10_000,
FIXTURE_CERT_GEN,
None,
descriptor(&["nrpc:a"]),
),
0,
);
let a = cap_id("nrpc:a");
assert_eq!(
state
.find_owner_private_providers(Some(&a), 0, &no_floors())
.len(),
1
);
let floor_above = floor_state(&org_kp, &member, FIXTURE_CERT_GEN + 1);
assert!(
state
.find_owner_private_providers(Some(&a), 0, &floor_above)
.is_empty(),
"a raised floor retracts the indexed record at query time"
);
}
#[test]
fn sweep_expired_reports_the_demoted_live_keys() {
let mut store = ScopedDiscoveryStore::new();
store.ingest(owner_cap(3, 1, 1000), 0, &NoConsumerGrants); store.ingest(grant_cap([0xAA; 32], 4, 1, 5000), 0, &NoConsumerGrants); let removed = store.sweep_expired(2000);
assert_eq!(removed, vec![(owner_scope(), provider(3))]);
}
#[test]
fn the_internal_capacity_sweep_reports_its_demotions() {
let mut store = ScopedDiscoveryStore::new();
for index in 0..ScopedDiscoveryStore::MAX_ENTRIES_PER_SCOPE as u64 {
store.ingest(owner_cap_n(index, 1, 1000), 0, &NoConsumerGrants);
}
let report = store.ingest(owner_cap_n(u64::MAX, 1, 5000), 2000, &NoConsumerGrants);
assert_eq!(report.outcome, ScopedStoreOutcome::Inserted);
assert_eq!(
report.swept_live.len(),
ScopedDiscoveryStore::MAX_ENTRIES_PER_SCOPE
);
}
fn grant_cap_declaring(
grant_id: [u8; 32],
provider_seed: u8,
generation: u64,
expires_at: u64,
tag: &str,
) -> VerifiedScopedCapability {
VerifiedScopedCapability::for_test(
CapabilityAudienceScope::Grant {
grant_id,
audience_handle: [0x22; 32],
},
provider(provider_seed),
org(2),
generation,
expires_at,
FIXTURE_CERT_GEN,
Some([0x5A; 64]),
descriptor(&[tag]),
)
}
fn one_cap(tag: &str) -> DirtyCapabilities {
DirtyCapabilities::Caps([cap_id(tag)].into_iter().collect())
}
#[test]
fn an_owner_ingest_advances_both_generations_and_dirties_its_capability() {
let mut state = ScopedDiscoveryState::new();
assert_eq!(state.revision(), 0);
assert_eq!(state.owner_revision(), 0);
ingest_indexed(
&mut state,
owner_cap_declaring(3, 1, 10_000, &["nrpc:a"]),
0,
);
assert_eq!(state.revision(), 1);
assert_eq!(state.owner_revision(), 1);
assert_eq!(state.take_global_change_batch().dirty, one_cap("nrpc:a"));
assert_eq!(state.take_owner_change_batch().dirty, one_cap("nrpc:a"));
assert_eq!(
state.take_global_change_batch().dirty,
DirtyCapabilities::Clean
);
assert_eq!(
state.take_owner_change_batch().dirty,
DirtyCapabilities::Clean
);
}
#[test]
fn grant_churn_never_advances_the_owner_stream() {
let mut state = ScopedDiscoveryState::new();
ingest_indexed(
&mut state,
grant_cap_declaring([0xAA; 32], 4, 1, 10_000, "nrpc:g"),
0,
);
assert_eq!(state.revision(), 1, "global advances");
assert_eq!(state.owner_revision(), 0, "owner does not");
assert_eq!(state.take_global_change_batch().dirty, one_cap("nrpc:g"));
assert_eq!(
state.take_owner_change_batch().dirty,
DirtyCapabilities::Clean
);
}
#[test]
fn a_stale_ingest_advances_no_generation() {
let mut state = ScopedDiscoveryState::new();
ingest_indexed(
&mut state,
owner_cap_declaring(3, 2, 10_000, &["nrpc:a"]),
0,
);
let (rev, owner_rev) = (state.revision(), state.owner_revision());
let _ = state.take_global_change_batch().dirty;
let _ = state.take_owner_change_batch().dirty;
assert_eq!(
ingest_indexed(
&mut state,
owner_cap_declaring(3, 1, 10_000, &["nrpc:a"]),
0
),
ScopedStoreOutcome::Stale
);
assert_eq!(state.revision(), rev, "stale advances nothing");
assert_eq!(state.owner_revision(), owner_rev);
assert_eq!(
state.take_global_change_batch().dirty,
DirtyCapabilities::Clean
);
}
#[test]
fn an_update_dirties_the_old_and_new_capabilities() {
let mut state = ScopedDiscoveryState::new();
ingest_indexed(
&mut state,
owner_cap_declaring(3, 1, 10_000, &["nrpc:a"]),
0,
);
let _ = state.take_global_change_batch().dirty;
let _ = state.take_owner_change_batch().dirty;
assert_eq!(
ingest_indexed(
&mut state,
owner_cap_declaring(3, 2, 10_000, &["nrpc:b"]),
0
),
ScopedStoreOutcome::Updated
);
let expected: std::collections::BTreeSet<_> =
[cap_id("nrpc:a"), cap_id("nrpc:b")].into_iter().collect();
assert_eq!(
state.take_global_change_batch().dirty,
DirtyCapabilities::Caps(expected)
);
}
#[test]
fn a_sweep_dirties_the_expired_capability() {
let mut state = ScopedDiscoveryState::new();
ingest_indexed(&mut state, owner_cap_declaring(3, 1, 1000, &["nrpc:a"]), 0);
let rev = state.revision();
let _ = state.take_global_change_batch().dirty;
let _ = state.take_owner_change_batch().dirty;
assert_eq!(state.sweep_expired(2000), 1);
assert_eq!(state.revision(), rev + 1);
assert_eq!(state.take_global_change_batch().dirty, one_cap("nrpc:a"));
}
#[test]
fn the_delta_collapses_to_rebuild_all_past_the_bound() {
let mut state = ScopedDiscoveryState::new();
for i in 0..=MAX_DIRTY_CAPABILITIES as u64 {
let tag = format!("nrpc:svc{i}");
ingest_indexed(&mut state, owner_cap_declaring_n(i, 1, 10_000, &[&tag]), 0);
}
assert_eq!(
state.take_global_change_batch().dirty,
DirtyCapabilities::RebuildAll
);
}
#[test]
fn a_record_declaring_no_capability_advances_nothing() {
let mut state = ScopedDiscoveryState::new();
ingest_indexed(
&mut state,
VerifiedScopedCapability::for_test(
owner_scope(),
provider(3),
org(1),
1,
10_000,
FIXTURE_CERT_GEN,
None,
b"not-a-capability-set".to_vec(),
),
0,
);
assert_eq!(state.revision(), 0);
assert_eq!(
state.take_global_change_batch().dirty,
DirtyCapabilities::Clean
);
}
#[test]
fn the_change_batch_captures_generation_and_delta_atomically() {
let mut state = ScopedDiscoveryState::new();
ingest_indexed(
&mut state,
owner_cap_declaring(3, 1, 10_000, &["nrpc:a"]),
0,
);
let batch = state.take_global_change_batch();
assert_eq!(batch.generation, state.revision());
assert_eq!(batch.generation, 1);
assert_eq!(batch.dirty, one_cap("nrpc:a"));
let drained = state.take_global_change_batch();
assert_eq!(drained.generation, 1, "generation is not reset by a drain");
assert_eq!(drained.dirty, DirtyCapabilities::Clean);
let owner = state.take_owner_change_batch();
assert_eq!(owner.generation, state.owner_revision());
assert_eq!(owner.dirty, one_cap("nrpc:a"));
}
#[test]
fn a_record_declaring_too_many_capabilities_is_refused_fail_closed() {
let mut state = ScopedDiscoveryState::new();
let tags: Vec<String> = (0..=MAX_DECLARATIONS_PER_RECORD)
.map(|i| format!("nrpc:svc{i}"))
.collect();
let tag_refs: Vec<&str> = tags.iter().map(String::as_str).collect();
let outcome = ingest_indexed(&mut state, owner_cap_declaring(3, 1, 10_000, &tag_refs), 0);
assert_eq!(outcome, ScopedStoreOutcome::TooManyDeclarations);
assert_eq!(state.len(), 0, "no row stored");
assert!(
state
.find_owner_private_providers(Some(&cap_id("nrpc:svc0")), 0, &no_floors())
.is_empty(),
"no index association"
);
assert_eq!(state.revision(), 0, "no generation advance");
assert_eq!(
state.take_global_change_batch().dirty,
DirtyCapabilities::Clean,
"no dirty"
);
}
#[test]
fn an_internal_capacity_demotion_is_visible_through_the_state_on_at_capacity() {
let mut state = ScopedDiscoveryState::new();
for index in 0..ScopedDiscoveryStore::MAX_ENTRIES_PER_SCOPE as u64 {
ingest_indexed(
&mut state,
owner_cap_declaring_n(index, 1, 10_000, &["nrpc:a"]),
0,
);
ingest_indexed(
&mut state,
owner_cap_declaring_n(index, 2, 1000, &["nrpc:a"]),
0,
);
}
let _ = state.take_global_change_batch();
let _ = state.take_owner_change_batch();
let (rev, owner_rev) = (state.revision(), state.owner_revision());
assert_eq!(
state
.find_owner_private_providers(Some(&cap_id("nrpc:a")), 500, &no_floors())
.len(),
ScopedDiscoveryStore::MAX_ENTRIES_PER_SCOPE,
"all live before the sweep"
);
let outcome = ingest_indexed(
&mut state,
owner_cap_declaring_n(u64::MAX, 1, 20_000, &["nrpc:b"]),
2000,
);
assert_eq!(outcome, ScopedStoreOutcome::AtCapacity);
assert!(
state
.find_owner_private_providers(Some(&cap_id("nrpc:a")), 2000, &no_floors())
.is_empty(),
"the demoted records left the owner index"
);
assert_eq!(state.revision(), rev + 1, "global generation advanced");
assert_eq!(
state.owner_revision(),
owner_rev + 1,
"owner generation advanced"
);
assert_eq!(state.take_global_change_batch().dirty, one_cap("nrpc:a"));
assert_eq!(state.take_owner_change_batch().dirty, one_cap("nrpc:a"));
}
#[test]
fn next_visible_expiry_tracks_the_earliest_live_deadline() {
let mut state = ScopedDiscoveryState::new();
assert_eq!(
state.next_visible_expiry(),
None,
"empty live set has no deadline"
);
ingest_indexed(&mut state, owner_cap_declaring(3, 1, 5000, &["nrpc:a"]), 0);
assert_eq!(state.next_visible_expiry(), Some(5000));
ingest_indexed(&mut state, owner_cap_declaring(4, 1, 9000, &["nrpc:a"]), 0);
assert_eq!(state.next_visible_expiry(), Some(5000));
ingest_indexed(&mut state, owner_cap_declaring(5, 1, 1000, &["nrpc:a"]), 0);
assert_eq!(state.next_visible_expiry(), Some(1000));
}
#[test]
fn an_update_moves_the_records_expiry_slot() {
let mut state = ScopedDiscoveryState::new();
ingest_indexed(&mut state, owner_cap_declaring(3, 1, 1000, &["nrpc:a"]), 0);
assert_eq!(state.next_visible_expiry(), Some(1000));
assert_eq!(
ingest_indexed(&mut state, owner_cap_declaring(3, 2, 5000, &["nrpc:a"]), 0),
ScopedStoreOutcome::Updated
);
assert_eq!(
state.next_visible_expiry(),
Some(5000),
"the update released the vacated 1000 slot"
);
}
#[test]
fn records_sharing_a_deadline_are_reference_counted() {
let mut state = ScopedDiscoveryState::new();
ingest_indexed(&mut state, owner_cap_declaring(3, 1, 1000, &["nrpc:a"]), 0);
ingest_indexed(&mut state, owner_cap_declaring(4, 1, 1000, &["nrpc:a"]), 0);
assert_eq!(state.next_visible_expiry(), Some(1000));
ingest_indexed(&mut state, owner_cap_declaring(3, 2, 5000, &["nrpc:a"]), 0);
assert_eq!(
state.next_visible_expiry(),
Some(1000),
"provider 4 still holds the shared deadline"
);
ingest_indexed(&mut state, owner_cap_declaring(4, 2, 5000, &["nrpc:a"]), 0);
assert_eq!(state.next_visible_expiry(), Some(5000));
}
#[test]
fn a_sweep_advances_next_visible_expiry_to_the_survivor() {
let mut state = ScopedDiscoveryState::new();
ingest_indexed(&mut state, owner_cap_declaring(3, 1, 1000, &["nrpc:a"]), 0);
ingest_indexed(&mut state, owner_cap_declaring(4, 1, 5000, &["nrpc:a"]), 0);
assert_eq!(state.next_visible_expiry(), Some(1000));
assert_eq!(state.sweep_expired(2000), 1);
assert_eq!(
state.next_visible_expiry(),
Some(5000),
"the swept 1000 slot was released with its live record"
);
assert_eq!(state.sweep_expired(6000), 1);
assert_eq!(state.next_visible_expiry(), None);
}
fn owner_cap_declaring_nothing(
provider_seed: u8,
generation: u64,
expires_at: u64,
) -> VerifiedScopedCapability {
VerifiedScopedCapability::for_test(
owner_scope(),
provider(provider_seed),
org(1),
generation,
expires_at,
FIXTURE_CERT_GEN,
None,
b"not-a-capability-set".to_vec(),
)
}
#[test]
fn a_declaration_empty_insert_does_not_gate_the_next_expiry() {
let mut state = ScopedDiscoveryState::new();
ingest_indexed(&mut state, owner_cap_declaring_nothing(3, 1, 500), 0);
assert_eq!(
state.next_visible_expiry(),
None,
"an inert record reports no deadline"
);
ingest_indexed(&mut state, owner_cap_declaring(4, 1, 5000, &["nrpc:a"]), 0);
assert_eq!(
state.next_visible_expiry(),
Some(5000),
"only the declaring record's deadline is reported"
);
}
#[test]
fn an_update_to_no_declarations_releases_the_expiry_slot() {
let mut state = ScopedDiscoveryState::new();
ingest_indexed(&mut state, owner_cap_declaring(3, 1, 1000, &["nrpc:a"]), 0);
ingest_indexed(&mut state, owner_cap_declaring(4, 1, 5000, &["nrpc:a"]), 0);
assert_eq!(state.next_visible_expiry(), Some(1000));
assert_eq!(
ingest_indexed(&mut state, owner_cap_declaring_nothing(3, 2, 1000), 0),
ScopedStoreOutcome::Updated
);
assert_eq!(
state.next_visible_expiry(),
Some(5000),
"the now-inert record released its deadline"
);
}
#[test]
fn an_update_to_declared_installs_the_expiry_slot() {
let mut state = ScopedDiscoveryState::new();
ingest_indexed(&mut state, owner_cap_declaring_nothing(3, 1, 1000), 0);
ingest_indexed(&mut state, owner_cap_declaring(4, 1, 5000, &["nrpc:a"]), 0);
assert_eq!(
state.next_visible_expiry(),
Some(5000),
"only the declaring record gates it"
);
assert_eq!(
ingest_indexed(&mut state, owner_cap_declaring(3, 2, 1000, &["nrpc:a"]), 0),
ScopedStoreOutcome::Updated
);
assert_eq!(
state.next_visible_expiry(),
Some(1000),
"becoming query-visible installed the earlier deadline"
);
}
fn leased_state() -> Arc<parking_lot::Mutex<ScopedDiscoveryState>> {
let state = Arc::new(parking_lot::Mutex::new(ScopedDiscoveryState::new()));
{
let mut s = state.lock();
let prepared =
PreparedScopedCapability::prepare(owner_cap_declaring(3, 1, 10_000, &["nrpc:a"]));
s.ingest(prepared, 0, &NoConsumerGrants);
}
state
}
#[test]
fn a_stream_leases_to_one_holder_at_a_time() {
let drains = PrivateDiscoveryDrains::new(leased_state());
let global = drains
.mint(PrivateDiscoveryStream::Global)
.expect("first global claim");
assert!(
drains.mint(PrivateDiscoveryStream::Global).is_none(),
"a second live claim on a held stream is refused"
);
let owner = drains
.mint(PrivateDiscoveryStream::Owner)
.expect("owner claims independently");
assert!(drains.mint(PrivateDiscoveryStream::Owner).is_none());
drop((global, owner));
}
#[test]
fn two_mint_facades_over_one_source_cannot_split_the_global_stream() {
let state = leased_state();
let a = PrivateDiscoveryDrains::new(state.clone());
let b = PrivateDiscoveryDrains::new(state);
let _held = a.mint(PrivateDiscoveryStream::Global).expect("first mint");
assert!(
b.mint(PrivateDiscoveryStream::Global).is_none(),
"lease identity must belong to the source, not the mint façade"
);
}
#[test]
fn mint_facades_share_one_lease_per_stream_including_release() {
let state = leased_state();
let a = PrivateDiscoveryDrains::new(state.clone());
let b = PrivateDiscoveryDrains::new(state);
let held = a.mint(PrivateDiscoveryStream::Owner).expect("first mint");
assert!(
b.mint(PrivateDiscoveryStream::Owner).is_none(),
"the owner stream is exclusive across façades too"
);
drop(held);
assert!(
b.mint(PrivateDiscoveryStream::Owner).is_some(),
"releasing through one façade frees the stream for the other — one \
shared lease word, not two coincidentally-held ones"
);
}
#[test]
fn dropping_the_handle_releases_the_lease() {
let drains = PrivateDiscoveryDrains::new(leased_state());
let first = drains
.mint(PrivateDiscoveryStream::Global)
.expect("first claim");
assert!(drains.mint(PrivateDiscoveryStream::Global).is_none());
drop(first);
assert!(
drains.mint(PrivateDiscoveryStream::Global).is_some(),
"the released lease can be reclaimed by a successor"
);
}
#[test]
fn a_minted_drain_starts_from_rebuild_all() {
let state = leased_state();
let drains = PrivateDiscoveryDrains::new(state.clone());
let mut predecessor = drains
.mint(PrivateDiscoveryStream::Global)
.expect("first claim");
assert_eq!(
predecessor.drain().dirty,
DirtyCapabilities::RebuildAll,
"even the first ever mint starts from a complete recapture"
);
assert_eq!(predecessor.drain().dirty, DirtyCapabilities::Clean);
drop(predecessor);
let mut successor = drains
.mint(PrivateDiscoveryStream::Global)
.expect("successor claims");
assert_eq!(
successor.drain().dirty,
DirtyCapabilities::RebuildAll,
"a successor recaptures completely; no consumed delta is silently lost"
);
}
#[test]
fn the_rollback_guard_releases_an_unpublished_claim() {
let lease = AtomicBool::new(true);
drop(LeaseRollback {
lease: &lease,
armed: true,
});
assert!(
!lease.load(Ordering::Acquire),
"an armed guard releases the claim it was protecting"
);
let lease = AtomicBool::new(true);
LeaseRollback {
lease: &lease,
armed: true,
}
.disarm();
assert!(
lease.load(Ordering::Acquire),
"a disarmed guard leaves the successful claim in place"
);
let lease = Arc::new(AtomicBool::new(true));
let unwound = {
let lease = lease.clone();
std::panic::catch_unwind(std::panic::AssertUnwindSafe(move || {
let _rollback = LeaseRollback {
lease: &lease,
armed: true,
};
panic!("mint failed after claiming");
}))
};
assert!(unwound.is_err(), "the mint unwound");
assert!(
!lease.load(Ordering::Acquire),
"an unwinding mint releases its claim rather than stranding the stream"
);
}
#[test]
fn a_leaked_handle_strands_its_stream_rather_than_double_draining() {
let drains = PrivateDiscoveryDrains::new(leased_state());
let leaked = drains
.mint(PrivateDiscoveryStream::Global)
.expect("first claim");
std::mem::forget(leaked);
assert!(
drains.mint(PrivateDiscoveryStream::Global).is_none(),
"a leaked lease is never silently reclaimed into a second drainer"
);
}
#[test]
fn a_drain_routes_to_its_own_stream() {
let state = leased_state();
let drains = PrivateDiscoveryDrains::new(state);
let expected = one_cap("nrpc:a");
let mut global = drains
.mint(PrivateDiscoveryStream::Global)
.expect("global claim");
let _ = global.drain();
{
let mut s = global.state.lock();
let prepared =
PreparedScopedCapability::prepare(owner_cap_declaring(4, 1, 10_000, &["nrpc:a"]));
s.ingest(prepared, 0, &NoConsumerGrants);
}
assert_eq!(global.drain().dirty, expected, "global reports its delta");
assert_eq!(global.drain().dirty, DirtyCapabilities::Clean);
let mut owner = drains
.mint(PrivateDiscoveryStream::Owner)
.expect("owner claim");
assert_eq!(
owner.drain().dirty,
DirtyCapabilities::RebuildAll,
"draining global did not clean owner"
);
}
fn scan_entries_in_scope(
store: &ScopedDiscoveryStore,
scope: &CapabilityAudienceScope,
) -> usize {
store.entries.keys().filter(|(s, _)| s == scope).count()
}
fn assert_counts_match_scan(store: &ScopedDiscoveryStore) {
let scopes: BTreeSet<CapabilityAudienceScope> =
store.entries.keys().map(|(s, _)| s.clone()).collect();
for scope in &scopes {
assert_eq!(
store.entries_in_scope(scope),
scan_entries_in_scope(store, scope),
"maintained count must equal the scan for {scope:?}"
);
}
assert_eq!(
store.scope_counts.len(),
scopes.len(),
"no scope row may outlive its last entry"
);
assert!(
store.scope_counts.values().all(|c| *c > 0),
"no scope row may sit at zero"
);
}
#[test]
fn the_maintained_scope_count_matches_the_scan_across_transitions() {
let mut store = ScopedDiscoveryStore::new();
let grant = [0xAA; 32];
store.ingest(owner_cap(3, 1, 1000), 0, &NoConsumerGrants);
store.ingest(owner_cap(4, 1, 5000), 0, &NoConsumerGrants);
store.ingest(grant_cap(grant, 5, 1, 1000), 0, &NoConsumerGrants);
assert_counts_match_scan(&store);
assert_eq!(store.entries_in_scope(&owner_scope()), 2);
assert_eq!(
store
.ingest(owner_cap(3, 2, 1000), 0, &NoConsumerGrants)
.outcome,
ScopedStoreOutcome::Updated
);
assert_eq!(
store.entries_in_scope(&owner_scope()),
2,
"an update must not grow the scope's occupancy"
);
assert_counts_match_scan(&store);
store.sweep_expired(2000);
assert_counts_match_scan(&store);
assert_eq!(
store.entries_in_scope(&owner_scope()),
1,
"the forgotten key freed its slot; the survivor keeps its own"
);
store.sweep_expired(6000);
assert_counts_match_scan(&store);
assert!(
store.scope_counts.is_empty(),
"an emptied store carries no scope rows"
);
}
#[test]
fn a_retained_tombstone_still_occupies_its_scope_slot() {
let mut store = ScopedDiscoveryStore::new();
store.ingest(owner_cap(3, 1, 10_000), 0, &NoConsumerGrants);
store.ingest(owner_cap(3, 2, 1000), 0, &NoConsumerGrants);
assert_eq!(store.entries_in_scope(&owner_scope()), 1);
store.sweep_expired(2000);
assert_eq!(store.len(), 0, "no live capability remains");
assert_eq!(
store.entries_in_scope(&owner_scope()),
1,
"the retained tombstone still holds its slot"
);
assert_counts_match_scan(&store);
}
#[test]
fn the_maintained_count_drives_the_fail_closed_scope_guard() {
let mut store = ScopedDiscoveryStore::new();
for index in 0..ScopedDiscoveryStore::MAX_ENTRIES_PER_SCOPE as u64 {
store.ingest(owner_cap_n(index, 1, 10_000), 0, &NoConsumerGrants);
}
assert_eq!(
store.entries_in_scope(&owner_scope()),
ScopedDiscoveryStore::MAX_ENTRIES_PER_SCOPE
);
let refused = store.ingest(owner_cap_n(u64::MAX, 1, 10_000), 0, &NoConsumerGrants);
assert_eq!(refused.outcome, ScopedStoreOutcome::AtCapacity);
assert_eq!(
store.entries_in_scope(&owner_scope()),
ScopedDiscoveryStore::MAX_ENTRIES_PER_SCOPE,
"a refused admission consumes no slot"
);
assert_counts_match_scan(&store);
}
#[test]
fn a_floor_raise_dirties_only_the_retracted_providers_capabilities() {
let mut state = ScopedDiscoveryState::new();
ingest_indexed(
&mut state,
owner_cap_declaring(3, 1, 10_000, &["nrpc:a"]),
0,
);
ingest_indexed(
&mut state,
owner_cap_declaring(4, 1, 10_000, &["nrpc:b"]),
0,
);
let _ = state.take_global_change_batch();
let _ = state.take_owner_change_batch();
let (rev, owner_rev) = (state.revision(), state.owner_revision());
let retracted = state.note_floors_raised(&[(org(1), provider(3), FIXTURE_CERT_GEN + 1)]);
assert_eq!(retracted, 1, "exactly provider 3's record was retracted");
assert_eq!(state.revision(), rev + 1, "global generation advanced");
assert_eq!(
state.owner_revision(),
owner_rev + 1,
"owner generation advanced"
);
assert_eq!(
state.take_global_change_batch().dirty,
one_cap("nrpc:a"),
"only the retracted provider's capability is dirty"
);
assert_eq!(state.take_owner_change_batch().dirty, one_cap("nrpc:a"));
}
#[test]
fn a_floor_at_the_admitted_generation_dirties_nothing() {
let mut state = ScopedDiscoveryState::new();
ingest_indexed(
&mut state,
owner_cap_declaring(3, 1, 10_000, &["nrpc:a"]),
0,
);
let _ = state.take_global_change_batch();
let _ = state.take_owner_change_batch();
let (rev, owner_rev) = (state.revision(), state.owner_revision());
let retracted = state.note_floors_raised(&[(org(1), provider(3), FIXTURE_CERT_GEN)]);
assert_eq!(
retracted, 0,
"a floor at the admitted generation is current"
);
assert_eq!(state.revision(), rev, "no generation advance");
assert_eq!(state.owner_revision(), owner_rev);
assert_eq!(
state.take_global_change_batch().dirty,
DirtyCapabilities::Clean
);
}
#[test]
fn a_floor_raise_for_another_org_retracts_nothing() {
let mut state = ScopedDiscoveryState::new();
ingest_indexed(
&mut state,
owner_cap_declaring(3, 1, 10_000, &["nrpc:a"]),
0,
);
let _ = state.take_global_change_batch();
let _ = state.take_owner_change_batch();
let rev = state.revision();
let retracted = state.note_floors_raised(&[(org(9), provider(3), FIXTURE_CERT_GEN + 1)]);
assert_eq!(retracted, 0, "a foreign org's raise retracts nothing");
assert_eq!(state.revision(), rev, "no generation advance");
assert_eq!(
state.take_global_change_batch().dirty,
DirtyCapabilities::Clean
);
}
#[test]
fn a_floor_raise_on_a_grant_record_never_advances_the_owner_stream() {
let mut state = ScopedDiscoveryState::new();
ingest_indexed(
&mut state,
grant_cap_declaring([0xAA; 32], 4, 1, 10_000, "nrpc:g"),
0,
);
let _ = state.take_global_change_batch();
let _ = state.take_owner_change_batch();
let (rev, owner_rev) = (state.revision(), state.owner_revision());
let retracted = state.note_floors_raised(&[(org(2), provider(4), FIXTURE_CERT_GEN + 1)]);
assert_eq!(retracted, 1);
assert_eq!(state.revision(), rev + 1, "global advances");
assert_eq!(state.owner_revision(), owner_rev, "owner does not");
assert_eq!(state.take_global_change_batch().dirty, one_cap("nrpc:g"));
assert_eq!(
state.take_owner_change_batch().dirty,
DirtyCapabilities::Clean
);
}
#[test]
fn the_reverse_provider_index_drops_swept_records() {
let mut state = ScopedDiscoveryState::new();
ingest_indexed(&mut state, owner_cap_declaring(3, 1, 1000, &["nrpc:a"]), 0);
ingest_indexed(&mut state, owner_cap_declaring(4, 1, 5000, &["nrpc:b"]), 0);
assert_eq!(state.index.floor_visible_by_provider.len(), 2);
state.sweep_expired(2000);
assert_eq!(
state
.index
.floor_visible_by_provider
.keys()
.collect::<Vec<_>>(),
vec![&provider(4)],
"the swept provider left the reverse index"
);
state.sweep_expired(6000);
assert!(
state.index.floor_visible_by_provider.is_empty(),
"no provider entry outlives its last live record"
);
}
#[test]
fn an_incremental_raise_over_a_hidden_row_dirties_nothing() {
let mut state = ScopedDiscoveryState::new();
ingest_indexed(
&mut state,
owner_cap_declaring(3, 1, 10_000, &["nrpc:a"]),
0,
);
let _ = state.take_global_change_batch();
let _ = state.take_owner_change_batch();
assert_eq!(
state.note_floors_raised(&[(org(1), provider(3), FIXTURE_CERT_GEN + 1)]),
1
);
let (rev, owner_rev) = (state.revision(), state.owner_revision());
assert_eq!(state.take_global_change_batch().dirty, one_cap("nrpc:a"));
assert_eq!(state.take_owner_change_batch().dirty, one_cap("nrpc:a"));
assert_eq!(
state.note_floors_raised(&[(org(1), provider(3), FIXTURE_CERT_GEN + 2)]),
0,
"the record was already invisible below the previous floor"
);
assert_eq!(state.revision(), rev, "no second generation advance");
assert_eq!(state.owner_revision(), owner_rev);
assert_eq!(
state.take_global_change_batch().dirty,
DirtyCapabilities::Clean,
"and no second invalidation"
);
}
#[test]
fn replaying_the_same_raise_is_a_complete_no_op() {
let mut state = ScopedDiscoveryState::new();
ingest_indexed(
&mut state,
owner_cap_declaring(3, 1, 10_000, &["nrpc:a"]),
0,
);
let _ = state.take_global_change_batch();
let _ = state.take_owner_change_batch();
let raise = [(org(1), provider(3), FIXTURE_CERT_GEN + 1)];
assert_eq!(state.note_floors_raised(&raise), 1, "the callback's pass");
let rev = state.revision();
let _ = state.take_global_change_batch();
let _ = state.take_owner_change_batch();
assert_eq!(state.note_floors_raised(&raise), 0, "the snapshot's pass");
assert_eq!(state.revision(), rev);
assert_eq!(
state.take_global_change_batch().dirty,
DirtyCapabilities::Clean
);
let mut fresh = ScopedDiscoveryState::new();
ingest_indexed(
&mut fresh,
owner_cap_declaring(3, 1, 10_000, &["nrpc:a"]),
0,
);
let _ = fresh.take_global_change_batch();
let before = fresh.revision();
assert_eq!(
fresh.note_floors_raised(&[
(org(1), provider(3), FIXTURE_CERT_GEN + 1),
(org(1), provider(3), FIXTURE_CERT_GEN + 2),
]),
1,
"a duplicated provider in one batch retracts once"
);
assert_eq!(
fresh.revision(),
before + 1,
"exactly one generation advance"
);
}
#[test]
fn the_expiry_of_a_floor_hidden_row_dirties_nothing() {
let mut state = ScopedDiscoveryState::new();
ingest_indexed(&mut state, owner_cap_declaring(3, 1, 1000, &["nrpc:a"]), 0);
let _ = state.take_global_change_batch();
let _ = state.take_owner_change_batch();
assert_eq!(
state.note_floors_raised(&[(org(1), provider(3), FIXTURE_CERT_GEN + 1)]),
1
);
let rev = state.revision();
let _ = state.take_global_change_batch();
let _ = state.take_owner_change_batch();
assert_eq!(
state.next_visible_expiry(),
None,
"a hidden row no longer gates the exact-expiry timer"
);
assert_eq!(state.sweep_expired(2000), 1, "the row is still reclaimed");
assert_eq!(state.revision(), rev, "no second generation advance");
assert_eq!(
state.take_global_change_batch().dirty,
DirtyCapabilities::Clean
);
}
#[test]
fn the_capacity_demotion_of_a_floor_hidden_row_dirties_nothing() {
let mut state = ScopedDiscoveryState::new();
for index in 0..ScopedDiscoveryStore::MAX_ENTRIES_PER_SCOPE as u64 {
ingest_indexed(
&mut state,
owner_cap_declaring_n(index, 1, 10_000, &["nrpc:a"]),
0,
);
ingest_indexed(
&mut state,
owner_cap_declaring_n(index, 2, 1000, &["nrpc:a"]),
0,
);
}
let raises: Vec<_> = (0..ScopedDiscoveryStore::MAX_ENTRIES_PER_SCOPE as u64)
.map(|i| (org(1), provider_n(i), FIXTURE_CERT_GEN + 1))
.collect();
assert_eq!(
state.note_floors_raised(&raises),
ScopedDiscoveryStore::MAX_ENTRIES_PER_SCOPE,
"every filler row is hidden"
);
let rev = state.revision();
let _ = state.take_global_change_batch();
let _ = state.take_owner_change_batch();
ingest_indexed(
&mut state,
owner_cap_declaring_n(u64::MAX, 1, 20_000, &["nrpc:b"]),
2000,
);
assert_eq!(
state.revision(),
rev,
"reclaiming already-hidden rows is not a visible transition"
);
assert_eq!(
state.take_global_change_batch().dirty,
DirtyCapabilities::Clean
);
}
#[test]
fn a_current_reannouncement_restores_a_floor_hidden_row() {
let mut state = ScopedDiscoveryState::new();
ingest_indexed(
&mut state,
owner_cap_declaring(3, 1, 10_000, &["nrpc:a"]),
0,
);
assert_eq!(
state.note_floors_raised(&[(org(1), provider(3), FIXTURE_CERT_GEN + 1)]),
1
);
let _ = state.take_global_change_batch();
let _ = state.take_owner_change_batch();
assert_eq!(state.next_visible_expiry(), None, "hidden: no deadline");
let restored = VerifiedScopedCapability::for_test(
owner_scope(),
provider(3),
org(1),
2,
10_000,
FIXTURE_CERT_GEN + 1,
None,
descriptor(&["nrpc:a"]),
);
assert_eq!(
ingest_indexed(&mut state, restored, 0),
ScopedStoreOutcome::Updated
);
assert_eq!(
state.next_visible_expiry(),
Some(10_000),
"the restored row gates the timer again"
);
assert_eq!(state.take_global_change_batch().dirty, one_cap("nrpc:a"));
assert_eq!(
state.note_floors_raised(&[(org(1), provider(3), FIXTURE_CERT_GEN + 1)]),
0,
"the restored certificate is current against the old floor"
);
assert_eq!(
state.note_floors_raised(&[(org(1), provider(3), FIXTURE_CERT_GEN + 2)]),
1,
"a higher floor retracts the restored row exactly once"
);
assert_eq!(
state.note_floors_raised(&[(org(1), provider(3), FIXTURE_CERT_GEN + 3)]),
0,
"and not again"
);
}
#[test]
fn floor_raise_dirtying_agrees_with_the_query_time_filter() {
let org_kp = OrgKeypair::from_bytes([7u8; 32]);
let org_id = org_kp.org_id();
let member = EntityId::from_bytes([9u8; 32]);
let mut state = ScopedDiscoveryState::new();
ingest_indexed(
&mut state,
VerifiedScopedCapability::for_test(
CapabilityAudienceScope::Owner {
org_id,
audience_handle: [0x11; 32],
},
member.clone(),
org_id,
1,
10_000,
FIXTURE_CERT_GEN,
None,
descriptor(&["nrpc:a"]),
),
0,
);
let a = cap_id("nrpc:a");
let _ = state.take_global_change_batch();
let floor_at = floor_state(&org_kp, &member, FIXTURE_CERT_GEN);
assert_eq!(
state
.find_owner_private_providers(Some(&a), 0, &floor_at)
.len(),
1,
"still visible at the boundary"
);
assert_eq!(
state.note_floors_raised(&[(org_id, member.clone(), FIXTURE_CERT_GEN)]),
0,
"and not reported retracted"
);
let floor_above = floor_state(&org_kp, &member, FIXTURE_CERT_GEN + 1);
assert!(
state
.find_owner_private_providers(Some(&a), 0, &floor_above)
.is_empty(),
"the query retracts it"
);
assert_eq!(
state.note_floors_raised(&[(org_id, member, FIXTURE_CERT_GEN + 1)]),
1,
"and the source dirties it"
);
assert_eq!(state.take_global_change_batch().dirty, one_cap("nrpc:a"));
}
#[test]
fn a_grant_records_deadline_also_gates_the_next_expiry() {
let mut state = ScopedDiscoveryState::new();
ingest_indexed(&mut state, owner_cap_declaring(3, 1, 2000, &["nrpc:a"]), 0);
ingest_indexed(
&mut state,
grant_cap_declaring([0xAA; 32], 4, 1, 800, "nrpc:g"),
0,
);
assert_eq!(
state.next_visible_expiry(),
Some(800),
"the earlier grant deadline gates the timer"
);
}
}