#![cfg(all(feature = "net", feature = "cortex"))]
use std::net::SocketAddr;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::Arc;
use std::time::{Duration, Instant};
use bytes::Bytes;
use net::adapter::net::behavior::capability::CapabilitySet;
use net::adapter::net::behavior::org::{OrgId, OrgKeypair, OrgMembershipCert, OrgRevocationBundle};
use net::adapter::net::behavior::org_admission::OrgAdmission;
use net::adapter::net::behavior::org_authority::NodeAuthority;
use net::adapter::net::behavior::org_call::MAX_ORG_PROOF_TTL_SECS;
use net::adapter::net::behavior::org_grant::OrgAudienceSecret;
use net::adapter::net::behavior::org_grant::{
CapabilityAuthorityId, DispatcherScope, GrantRights, GrantTargetScope, OrgCapabilityGrant,
OrgDispatcherGrant,
};
use net::adapter::net::behavior::org_grant_registry::{
GrantAudienceInstallError, GrantAudienceInstalled,
};
use net::adapter::net::behavior::CapabilityAnnouncement;
use net::adapter::net::cortex::{
RpcContext, RpcHandler, RpcHandlerError, RpcResponsePayload, RpcStatus,
};
use net::adapter::net::identity::EntityId;
use net::adapter::net::mesh_rpc::{
CallOptions, CodecDirection, OrgProofIntent, RpcError, ServeError, ServeHandle,
};
use net::adapter::net::{EntityKeypair, MeshNode, MeshNodeConfig, SocketBufferConfig};
const PSK: [u8; 32] = [0x42u8; 32];
const TEST_BUFFER_SIZE: usize = 256 * 1024;
const ORG_ADMISSION_HEADER: &str = "net-org-admission";
fn test_config() -> MeshNodeConfig {
let addr: SocketAddr = "127.0.0.1:0".parse().unwrap();
let mut cfg = MeshNodeConfig::new(addr, PSK)
.with_heartbeat_interval(Duration::from_millis(200))
.with_session_timeout(Duration::from_secs(5))
.with_handshake(3, Duration::from_secs(2))
.with_capability_gc_interval(Duration::from_millis(250));
cfg.socket_buffers = SocketBufferConfig {
send_buffer_size: TEST_BUFFER_SIZE,
recv_buffer_size: TEST_BUFFER_SIZE,
};
cfg
}
async fn build_node_with(keypair: EntityKeypair) -> Arc<MeshNode> {
Arc::new(
MeshNode::new(keypair, test_config())
.await
.expect("MeshNode::new"),
)
}
async fn handshake_pair(a: &Arc<MeshNode>, b: &Arc<MeshNode>) {
let a_id = a.node_id();
let b_id = b.node_id();
let b_pub = *b.public_key();
let b_addr = b.local_addr();
let b_clone = b.clone();
let accept = tokio::spawn(async move { b_clone.accept(a_id).await });
a.connect(b_addr, &b_pub, b_id)
.await
.expect("connect failed");
accept
.await
.expect("accept task panicked")
.expect("accept failed");
a.start();
b.start();
}
async fn build_node_fast_announce(keypair: EntityKeypair) -> Arc<MeshNode> {
let mut cfg = test_config();
cfg.min_announce_interval = Duration::from_millis(50);
Arc::new(MeshNode::new(keypair, cfg).await.expect("MeshNode::new"))
}
async fn connect_no_start(initiator: &Arc<MeshNode>, responder: &Arc<MeshNode>) {
let r_id = responder.node_id();
let r_pub = *responder.public_key();
let r_addr = responder.local_addr();
let i_id = initiator.node_id();
let responder_c = responder.clone();
let accept = tokio::spawn(async move { responder_c.accept(i_id).await });
initiator
.connect(r_addr, &r_pub, r_id)
.await
.expect("connect failed");
accept
.await
.expect("accept task panicked")
.expect("accept failed");
}
async fn wait_until<F: Fn() -> bool>(limit: Duration, cond: F) -> bool {
let start = Instant::now();
while start.elapsed() < limit {
if cond() {
return true;
}
tokio::time::sleep(Duration::from_millis(25)).await;
}
cond()
}
async fn bring_up(caller: &Arc<MeshNode>, server: &Arc<MeshNode>) {
handshake_pair(caller, server).await;
server
.announce_capabilities(CapabilitySet::new())
.await
.expect("server announce");
caller
.announce_capabilities(CapabilitySet::new())
.await
.expect("caller announce");
let caller_id = caller.node_id();
let server_id = server.node_id();
assert!(
wait_until(Duration::from_secs(5), || {
caller.peer_entity_id(server_id).is_some() && server.peer_entity_id(caller_id).is_some()
})
.await,
"entity pins established in both directions",
);
}
fn install_authority(server: &Arc<MeshNode>, tag: &str) -> (OrgKeypair, std::path::PathBuf) {
let node_entity = server.entity_id().clone();
let org_b = OrgKeypair::from_bytes([0x42u8; 32]);
let node_cert =
OrgMembershipCert::try_issue(&org_b, node_entity.clone(), 1, 3600).expect("node cert");
let dir = std::env::temp_dir().join(format!("net-oa2-live-{tag}-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let authority =
NodeAuthority::adopt(&dir, node_cert, &node_entity, 0, None).expect("adopt authority");
server
.install_node_authority(Arc::new(authority))
.expect("install authority");
(org_b, dir)
}
fn fold_restrictive_announcement(
nodes: &[&Arc<MeshNode>],
target: &Arc<MeshNode>,
version: u64,
tag: &str,
allowed_nodes: Vec<u64>,
) {
let caps = CapabilitySet::new().add_tag(tag);
let mut ann =
CapabilityAnnouncement::new(target.node_id(), target.entity_id().clone(), version, caps);
ann.allowed_nodes = allowed_nodes;
for n in nodes {
n.test_inject_capability_announcement(ann.clone());
}
}
fn owner_delegated_intent(
caller_kp: EntityKeypair,
org_b: &OrgKeypair,
provider: EntityId,
service: &str,
) -> OrgProofIntent {
owner_delegated_intent_gen(caller_kp, org_b, provider, service, 1)
}
fn owner_delegated_intent_gen(
caller_kp: EntityKeypair,
org_b: &OrgKeypair,
provider: EntityId,
service: &str,
generation: u32,
) -> OrgProofIntent {
let caller_entity = caller_kp.entity_id().clone();
let cap = CapabilityAuthorityId::for_tag(&format!("nrpc:{service}"));
let membership = OrgMembershipCert::try_issue(org_b, caller_entity.clone(), generation, 3600)
.expect("membership");
let dispatcher =
OrgDispatcherGrant::try_issue(org_b, caller_entity, DispatcherScope::Exact(cap), 3600)
.expect("dispatcher");
OrgProofIntent {
caller: Arc::new(caller_kp),
membership,
dispatcher,
capability_grant: None,
acting_org: org_b.org_id(),
provider_owner_org: org_b.org_id(),
provider,
capability: cap,
proof_ttl_secs: 30,
}
}
struct AdmitHandler {
calls: Arc<AtomicUsize>,
saw_admission: Arc<AtomicBool>,
attribution_ok: Arc<AtomicBool>,
proof_stripped: Arc<AtomicBool>,
expected_caller: EntityId,
expected_acting_org: OrgId,
expected_provider_org: OrgId,
expected_provider: EntityId,
expected_capability: CapabilityAuthorityId,
}
#[async_trait::async_trait]
impl RpcHandler for AdmitHandler {
async fn call(&self, ctx: RpcContext) -> Result<RpcResponsePayload, RpcHandlerError> {
self.calls.fetch_add(1, Ordering::SeqCst);
if let Some(admitted) = ctx.org_admission.as_ref() {
self.saw_admission.store(true, Ordering::SeqCst);
if admitted.caller == self.expected_caller
&& admitted.acting_org == self.expected_acting_org
&& admitted.provider_org == self.expected_provider_org
&& admitted.provider == self.expected_provider
&& admitted.capability == self.expected_capability
{
self.attribution_ok.store(true, Ordering::SeqCst);
}
}
let stripped = !ctx
.payload
.headers
.iter()
.any(|(name, _)| name == ORG_ADMISSION_HEADER);
self.proof_stripped.store(stripped, Ordering::SeqCst);
Ok(RpcResponsePayload {
status: RpcStatus::Ok,
headers: vec![],
body: Bytes::from_static(b"pong"),
})
}
}
struct DarkHandler {
calls: Arc<AtomicUsize>,
}
#[async_trait::async_trait]
impl RpcHandler for DarkHandler {
async fn call(&self, _ctx: RpcContext) -> Result<RpcResponsePayload, RpcHandlerError> {
self.calls.fetch_add(1, Ordering::SeqCst);
Ok(RpcResponsePayload {
status: RpcStatus::Ok,
headers: vec![],
body: Bytes::new(),
})
}
}
async fn assert_handler_stays_dark(calls: &Arc<AtomicUsize>, what: &str) {
const SETTLE: Duration = Duration::from_millis(200);
const STEP: Duration = Duration::from_millis(10);
let deadline = Instant::now() + SETTLE;
loop {
let observed = calls.load(Ordering::SeqCst);
assert_eq!(
observed, 0,
"{what}: the handler RAN ({observed} call(s)) despite the denial — \
the request reached the fold even though the caller was denied",
);
if Instant::now() >= deadline {
return;
}
tokio::time::sleep(STEP).await;
}
}
#[tokio::test]
#[should_panic(expected = "the handler RAN")]
async fn the_darkness_window_catches_a_handler_that_runs_late() {
let calls = Arc::new(AtomicUsize::new(0));
assert_eq!(calls.load(Ordering::SeqCst), 0);
let late = Arc::clone(&calls);
tokio::spawn(async move {
tokio::time::sleep(Duration::from_millis(50)).await;
late.fetch_add(1, Ordering::SeqCst);
});
assert_handler_stays_dark(&calls, "a handler scheduled after the denial").await;
}
#[tokio::test]
async fn live_two_node_owner_delegated_admit() {
const CALLER_SEED: [u8; 32] = [0x07u8; 32];
let server = build_node_with(EntityKeypair::generate()).await;
let caller = build_node_with(EntityKeypair::from_bytes(CALLER_SEED)).await;
bring_up(&caller, &server).await;
let (org_b, _dir) = install_authority(&server, "admit");
let provider = server.entity_id().clone();
let caller_entity = caller.entity_id().clone();
let calls = Arc::new(AtomicUsize::new(0));
let saw = Arc::new(AtomicBool::new(false));
let attribution_ok = Arc::new(AtomicBool::new(false));
let stripped = Arc::new(AtomicBool::new(false));
let _serve = server
.serve_rpc_protected(
"svc",
Arc::new(AdmitHandler {
calls: calls.clone(),
saw_admission: saw.clone(),
attribution_ok: attribution_ok.clone(),
proof_stripped: stripped.clone(),
expected_caller: caller_entity,
expected_acting_org: org_b.org_id(),
expected_provider_org: org_b.org_id(),
expected_provider: provider.clone(),
expected_capability: CapabilityAuthorityId::for_tag("nrpc:svc"),
}),
OrgAdmission::OwnerDelegated,
Arc::new(|_| true),
)
.expect("serve protected");
let intent = owner_delegated_intent(
EntityKeypair::from_bytes(CALLER_SEED),
&org_b,
provider,
"svc",
);
let opts = CallOptions {
org_proof_intent: Some(intent),
deadline: Some(Instant::now() + Duration::from_secs(5)),
..Default::default()
};
let reply = caller
.call(server.node_id(), "svc", Bytes::from_static(b"ping"), opts)
.await
.expect("admitted call returns Ok");
assert_eq!(reply.body.as_ref(), b"pong", "handler reply body");
assert_eq!(calls.load(Ordering::SeqCst), 1, "handler ran exactly once");
assert!(
saw.load(Ordering::SeqCst),
"handler observed org_admission attribution",
);
assert!(
attribution_ok.load(Ordering::SeqCst),
"all four attribution parties (caller, acting org, provider org, provider) plus the \
exact nrpc:svc capability match",
);
assert!(
stripped.load(Ordering::SeqCst),
"the raw net-org-admission proof header was stripped from the handler view",
);
}
#[tokio::test]
async fn live_two_node_missing_proof_denied() {
let server = build_node_with(EntityKeypair::generate()).await;
let caller = build_node_with(EntityKeypair::from_bytes([0x08u8; 32])).await;
bring_up(&caller, &server).await;
let (_org_b, _dir) = install_authority(&server, "deny");
let calls = Arc::new(AtomicUsize::new(0));
let _serve = server
.serve_rpc_protected(
"svc",
Arc::new(DarkHandler {
calls: calls.clone(),
}),
OrgAdmission::OwnerDelegated,
Arc::new(|_| true),
)
.expect("serve protected");
let opts = CallOptions {
deadline: Some(Instant::now() + Duration::from_secs(5)),
..Default::default()
};
let err = caller
.call(server.node_id(), "svc", Bytes::from_static(b"ping"), opts)
.await
.expect_err("a public call to a protected service must be denied");
match err {
RpcError::ServerError {
status, message, ..
} => {
assert_eq!(status, 0x0009, "status is exactly AdmissionDenied (0x0009)");
assert_eq!(
message.len(),
1,
"the deny body carries exactly one coarse reason byte",
);
assert_eq!(
message.as_bytes(),
&[0u8],
"a missing proof must report exactly Denied (0); a different \
coarse reason leaks provider state to an uncredentialed caller",
);
}
other => panic!(
"expected an AdmissionDenied ServerError, got {other:?} \
(a Timeout here would be a denial masquerading as a timeout)"
),
}
assert_handler_stays_dark(&calls, "the handler stayed dark for the denied call").await;
}
#[cfg(feature = "fixtures")]
#[tokio::test]
async fn live_two_node_provider_store_poison_denies() {
const CALLER_SEED: [u8; 32] = [0x09u8; 32];
let server = build_node_with(EntityKeypair::generate()).await;
let caller = build_node_with(EntityKeypair::from_bytes(CALLER_SEED)).await;
bring_up(&caller, &server).await;
let (org_b, _dir) = install_authority(&server, "poison");
let provider = server.entity_id().clone();
let calls = Arc::new(AtomicUsize::new(0));
let _serve = server
.serve_rpc_protected(
"svc",
Arc::new(DarkHandler {
calls: calls.clone(),
}),
OrgAdmission::OwnerDelegated,
Arc::new(|_| true),
)
.expect("serve protected");
server
.org_revocation_store()
.expect("a revocation store is installed")
.mark_poisoned_for_test();
let intent = owner_delegated_intent(
EntityKeypair::from_bytes(CALLER_SEED),
&org_b,
provider,
"svc",
);
let opts = CallOptions {
org_proof_intent: Some(intent),
deadline: Some(Instant::now() + Duration::from_secs(5)),
..Default::default()
};
let err = caller
.call(server.node_id(), "svc", Bytes::from_static(b"ping"), opts)
.await
.expect_err("a poisoned provider store must deny even a valid proof");
match err {
RpcError::ServerError {
status, message, ..
} => {
assert_eq!(status, 0x0009, "status is AdmissionDenied (0x0009)");
assert_eq!(
message.as_bytes(),
&[2u8],
"coarse reason is exactly Unavailable (provider authority unavailable)",
);
}
other => panic!("expected an AdmissionDenied ServerError, got {other:?}"),
}
assert_handler_stays_dark(&calls, "the handler stayed dark under the poisoned store").await;
}
#[tokio::test]
async fn live_two_node_policy_veto_denies() {
const CALLER_SEED: [u8; 32] = [0x0au8; 32];
let server = build_node_with(EntityKeypair::generate()).await;
let caller = build_node_with(EntityKeypair::from_bytes(CALLER_SEED)).await;
bring_up(&caller, &server).await;
let (org_b, _dir) = install_authority(&server, "veto");
let provider = server.entity_id().clone();
let calls = Arc::new(AtomicUsize::new(0));
let _serve = server
.serve_rpc_protected(
"svc",
Arc::new(DarkHandler {
calls: calls.clone(),
}),
OrgAdmission::OwnerDelegated,
Arc::new(|_| false),
)
.expect("serve protected");
let intent = owner_delegated_intent(
EntityKeypair::from_bytes(CALLER_SEED),
&org_b,
provider,
"svc",
);
let opts = CallOptions {
org_proof_intent: Some(intent),
deadline: Some(Instant::now() + Duration::from_secs(5)),
..Default::default()
};
let err = caller
.call(server.node_id(), "svc", Bytes::from_static(b"ping"), opts)
.await
.expect_err("a vetoing provider policy must deny a valid proof");
match err {
RpcError::ServerError {
status, message, ..
} => {
assert_eq!(status, 0x0009, "status is AdmissionDenied (0x0009)");
assert_eq!(
message.len(),
1,
"the deny body carries exactly one coarse reason byte",
);
}
other => {
panic!("expected an AdmissionDenied ServerError, got {other:?} (no timeout masquerade)")
}
}
assert_handler_stays_dark(&calls, "the handler stayed dark under the policy veto").await;
}
struct HeaderSpyHandler {
calls: Arc<AtomicUsize>,
saw_proof: Arc<AtomicBool>,
}
#[async_trait::async_trait]
impl RpcHandler for HeaderSpyHandler {
async fn call(&self, ctx: RpcContext) -> Result<RpcResponsePayload, RpcHandlerError> {
self.calls.fetch_add(1, Ordering::SeqCst);
if ctx
.payload
.headers
.iter()
.any(|(name, _)| name == ORG_ADMISSION_HEADER)
{
self.saw_proof.store(true, Ordering::SeqCst);
}
Ok(RpcResponsePayload {
status: RpcStatus::Ok,
headers: vec![],
body: Bytes::from_static(b"pong"),
})
}
}
#[tokio::test]
async fn live_two_node_public_handler_never_sees_proof_header() {
const CALLER_SEED: [u8; 32] = [0x0bu8; 32];
let server = build_node_with(EntityKeypair::generate()).await;
let caller = build_node_with(EntityKeypair::from_bytes(CALLER_SEED)).await;
bring_up(&caller, &server).await;
let org_b = OrgKeypair::from_bytes([0x42u8; 32]);
let provider = server.entity_id().clone();
let calls = Arc::new(AtomicUsize::new(0));
let saw_proof = Arc::new(AtomicBool::new(false));
let _serve = server
.serve_rpc(
"pub",
Arc::new(HeaderSpyHandler {
calls: calls.clone(),
saw_proof: saw_proof.clone(),
}),
)
.expect("serve public");
let intent = owner_delegated_intent(
EntityKeypair::from_bytes(CALLER_SEED),
&org_b,
provider,
"pub",
);
let opts = CallOptions {
org_proof_intent: Some(intent),
deadline: Some(Instant::now() + Duration::from_secs(5)),
..Default::default()
};
let reply = caller
.call(server.node_id(), "pub", Bytes::from_static(b"ping"), opts)
.await
.expect("a public call carrying a stray proof still succeeds");
assert_eq!(reply.body.as_ref(), b"pong");
assert_eq!(
calls.load(Ordering::SeqCst),
1,
"the public handler ran once"
);
assert!(
!saw_proof.load(Ordering::SeqCst),
"the public handler never saw the org-admission proof header (stripped by the bridge)",
);
}
#[tokio::test]
async fn live_two_node_protected_call_service_bypasses_legacy_gate() {
const CALLER_SEED: [u8; 32] = [0x0cu8; 32];
let server = build_node_with(EntityKeypair::generate()).await;
let caller = build_node_with(EntityKeypair::from_bytes(CALLER_SEED)).await;
bring_up(&caller, &server).await;
let (org_b, _dir) = install_authority(&server, "callservice");
let provider = server.entity_id().clone();
let calls = Arc::new(AtomicUsize::new(0));
let _serve = server
.serve_rpc_protected(
"svc",
Arc::new(DarkHandler {
calls: calls.clone(),
}),
OrgAdmission::OwnerDelegated,
Arc::new(|_| true),
)
.expect("serve protected");
fold_restrictive_announcement(&[&caller], &server, 100, "nrpc:svc", vec![server.node_id()]);
let deny = caller
.call_service(
"svc",
Bytes::from_static(b"ping"),
CallOptions {
deadline: Some(Instant::now() + Duration::from_secs(5)),
..Default::default()
},
)
.await
.expect_err("a public call_service must be denied by the legacy allow-list");
assert!(
matches!(deny, RpcError::CapabilityDenied { .. }),
"public call_service is denied by the legacy allow-list, got {deny:?}",
);
assert_handler_stays_dark(&calls, "the denied public call never reached the handler").await;
let intent = owner_delegated_intent(
EntityKeypair::from_bytes(CALLER_SEED),
&org_b,
provider,
"svc",
);
caller
.call_service(
"svc",
Bytes::from_static(b"ping"),
CallOptions {
org_proof_intent: Some(intent),
deadline: Some(Instant::now() + Duration::from_secs(5)),
..Default::default()
},
)
.await
.expect("protected call_service must bypass the legacy gate and admit");
assert_eq!(
calls.load(Ordering::SeqCst),
1,
"the protected call bypassed the legacy gate and reached the handler exactly once",
);
}
struct TrivialHandler;
#[async_trait::async_trait]
impl RpcHandler for TrivialHandler {
async fn call(&self, _ctx: RpcContext) -> Result<RpcResponsePayload, RpcHandlerError> {
Ok(RpcResponsePayload {
status: RpcStatus::Ok,
headers: vec![],
body: Bytes::from_static(b"ok"),
})
}
}
#[tokio::test]
async fn owner_scoped_residue_is_stripped_from_the_plaintext_announcement() {
let server = build_node_with(EntityKeypair::from_bytes([0x51u8; 32])).await;
let (_org_b, _dir) = install_authority(&server, "residue-strip");
let _secret = server
.serve_rpc_owner_scoped("secret", Arc::new(TrivialHandler), Arc::new(|_| true))
.expect("owner-scoped serve");
let baseline = CapabilitySet::new()
.add_tag("nrpc:secret")
.add_tag("region:eu-west");
server
.announce_capabilities(baseline)
.await
.expect("announce");
assert!(
wait_until(Duration::from_secs(5), || {
server
.local_announcement_for_test()
.map(|a| {
!a.capabilities.has_tag("nrpc:secret")
&& a.capabilities.has_tag("region:eu-west")
})
.unwrap_or(false)
})
.await,
"owner-scoped baseline residue stripped from plaintext; unrelated tag kept",
);
tokio::time::sleep(Duration::from_millis(250)).await;
let settled = server
.local_announcement_for_test()
.expect("an announcement is published by now");
assert!(
!settled.capabilities.has_tag("nrpc:secret"),
"the owner-scoped tag reappeared in plaintext after convergence — a \
second re-announce path is republishing the baseline residue",
);
assert!(
settled.capabilities.has_tag("region:eu-west"),
"the unrelated baseline tag must survive the strip",
);
}
#[tokio::test]
async fn owner_scoped_service_ships_only_inside_the_encrypted_owner_envelope() {
use net::adapter::net::behavior::org_revocation::OrgRevocationState;
use net::adapter::net::behavior::org_scoped_ann::ScopedCapabilityAnnouncement;
use net::adapter::net::behavior::org_scoped_ingest::{
verify_scoped_ingest, AudienceAuthority, ScopedIngestContext,
};
let server = build_node_with(EntityKeypair::from_bytes([0x53u8; 32])).await;
let (_org_b, _dir) = install_authority(&server, "scoped-delivery");
server
.set_owner_cert_emission(true)
.expect("enable owner-cert emission");
let _secret = server
.serve_rpc_owner_scoped("secret", Arc::new(TrivialHandler), Arc::new(|_| true))
.expect("owner-scoped serve");
let _public = server
.serve_rpc("open", Arc::new(TrivialHandler))
.expect("public serve");
server
.announce_capabilities(CapabilitySet::new())
.await
.expect("announce");
assert!(
wait_until(Duration::from_secs(5), || {
server.announcement_scoped_for_send_for_test().len() == 1
&& server
.local_announcement_for_test()
.map(|a| {
a.capabilities.has_tag("nrpc:open")
&& !a.capabilities.has_tag("nrpc:secret")
})
.unwrap_or(false)
})
.await,
"one scoped envelope emitted; plaintext keeps nrpc:open, drops nrpc:secret",
);
let scoped = server.announcement_scoped_for_send_for_test();
let envelope =
ScopedCapabilityAnnouncement::from_bytes(&scoped[0]).expect("decode scoped envelope");
let authority = server.node_authority().expect("authority installed");
let audience = AudienceAuthority::owner(authority.owner_org(), &authority.audience);
let floors = OrgRevocationState::empty();
let now_secs = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.expect("clock")
.as_secs();
let ctx = ScopedIngestContext {
local_owner_org: authority.owner_org(),
floors: &floors,
now_secs,
skew_secs: 5,
local_member: None,
};
let verified = verify_scoped_ingest(&envelope, &audience, &ctx).expect("owner ingest opens");
let descriptor = CapabilitySet::from_bytes(verified.descriptor()).expect("descriptor caps");
assert!(
descriptor.has_tag("nrpc:secret"),
"the encrypted descriptor names the owner-scoped service",
);
assert!(
!descriptor.has_tag("nrpc:open"),
"the encrypted descriptor carries only owner-scoped services, never public ones",
);
}
#[tokio::test]
async fn serve_rpc_granted_is_dispatchable_but_undiscoverable_without_a_grant() {
use net::adapter::net::behavior::fold::capability_bridge::has_local_capability;
let bare = build_node_with(EntityKeypair::from_bytes([0x62u8; 32])).await;
assert!(
matches!(
bare.serve_rpc_granted("cross", Arc::new(TrivialHandler), Arc::new(|_| true)),
Err(ServeError::ProtectedAuthorityRequired(_))
),
"a granted registration without authority is refused",
);
let server = build_node_with(EntityKeypair::from_bytes([0x63u8; 32])).await;
let (_org_b, _dir) = install_authority(&server, "granted-seam");
server
.set_owner_cert_emission(true)
.expect("enable owner-cert emission");
let _granted = server
.serve_rpc_granted("cross", Arc::new(TrivialHandler), Arc::new(|_| true))
.expect("granted serve");
let _protected = server
.serve_rpc_protected(
"prot",
Arc::new(TrivialHandler),
OrgAdmission::CrossOrgGranted,
Arc::new(|_| true),
)
.expect("protected serve");
let _public = server
.serve_rpc("open", Arc::new(TrivialHandler))
.expect("public serve");
server
.announce_capabilities(CapabilitySet::new())
.await
.expect("announce");
assert!(
wait_until(Duration::from_secs(5), || {
server
.local_announcement_for_test()
.map(|a| {
a.capabilities.has_tag("nrpc:open")
&& a.capabilities.has_tag("nrpc:prot")
&& !a.capabilities.has_tag("nrpc:cross")
})
.unwrap_or(false)
&& server.announcement_scoped_for_send_for_test().is_empty()
})
.await,
"plaintext keeps public+protected, drops granted; no envelope without a grant",
);
assert!(
has_local_capability(server.capability_fold(), server.node_id(), "nrpc:cross"),
"the granted service is locally dispatchable despite being undiscoverable",
);
}
fn unix_now() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.expect("clock")
.as_secs()
}
fn copy_secret(secret: &OrgAudienceSecret) -> OrgAudienceSecret {
let mut buf = secret.encode_config();
let copy = OrgAudienceSecret::decode_config(&buf).expect("copy secret");
for byte in buf.iter_mut() {
unsafe { std::ptr::write_volatile(byte, 0) };
}
copy
}
async fn granted_provider(
seed: u8,
tag: &str,
svc: &str,
) -> (
Arc<MeshNode>,
ServeHandle,
std::path::PathBuf,
EntityId,
OrgKeypair,
) {
let p = build_node_with(EntityKeypair::from_bytes([seed; 32])).await;
let entity = p.entity_id().clone();
let org_b = OrgKeypair::from_bytes([0x42u8; 32]);
let cert = OrgMembershipCert::try_issue(&org_b, entity.clone(), 1, 3600).expect("cert");
let dir = std::env::temp_dir().join(format!("net-oa34b2-emit-{tag}-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let authority = NodeAuthority::adopt(&dir, cert, &entity, 0, None).expect("adopt");
p.install_node_authority(Arc::new(authority))
.expect("install authority");
p.set_owner_cert_emission(true).expect("enable emission");
let handle = p
.serve_rpc_granted(svc, Arc::new(TrivialHandler), Arc::new(|_| true))
.expect("granted serve");
(p, handle, dir, entity, org_b)
}
fn open_granted_envelope(
scoped_bytes: &[u8],
grant: &OrgCapabilityGrant,
secret: &OrgAudienceSecret,
grantee_org: OrgId,
now_secs: u64,
) -> Option<CapabilitySet> {
use net::adapter::net::behavior::org_revocation::OrgRevocationState;
use net::adapter::net::behavior::org_scoped_ann::ScopedCapabilityAnnouncement;
use net::adapter::net::behavior::org_scoped_ingest::{
verify_scoped_ingest, AudienceAuthority, ScopedIngestContext,
};
let env = ScopedCapabilityAnnouncement::from_bytes(scoped_bytes).ok()?;
let authority = AudienceAuthority::granted(grant, secret);
let floors = OrgRevocationState::empty();
let ctx = ScopedIngestContext {
local_owner_org: grantee_org,
floors: &floors,
now_secs,
skew_secs: 5,
local_member: None,
};
let verified = verify_scoped_ingest(&env, &authority, &ctx).ok()?;
CapabilitySet::from_bytes(verified.descriptor())
}
async fn converge_scoped_count(p: &Arc<MeshNode>, n: usize) -> bool {
for _ in 0..40 {
p.announce_capabilities(CapabilitySet::new()).await.ok();
if p.announcement_scoped_for_send_for_test().len() == n {
return true;
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
p.announcement_scoped_for_send_for_test().len() == n
}
#[tokio::test]
async fn a_granted_service_ships_only_inside_an_encrypted_grant_envelope() {
let (p, _h, _dir, entity, org_b) = granted_provider(0x70, "one", "cross").await;
let org_a = OrgKeypair::from_bytes([0x7Au8; 32]);
let (grant, secret) = OrgCapabilityGrant::try_issue(
&org_b,
org_a.org_id(),
CapabilityAuthorityId::for_tag("nrpc:cross"),
GrantRights::DISCOVER,
GrantTargetScope::ExactNode(entity.clone()),
3600,
)
.expect("issue grant");
let secret = secret.expect("secret");
let opener = copy_secret(&secret);
assert_eq!(
p.install_provider_grant_audience(grant.clone(), secret)
.expect("install"),
GrantAudienceInstalled::Installed
);
assert!(
converge_scoped_count(&p, 1).await,
"P emits exactly one granted envelope",
);
assert!(
p.local_announcement_for_test()
.map(|a| !a.capabilities.has_tag("nrpc:cross"))
.unwrap_or(false),
"the granted tag never appears in the plaintext announcement",
);
let scoped = p.announcement_scoped_for_send_for_test();
let descriptor = open_granted_envelope(&scoped[0], &grant, &opener, org_a.org_id(), unix_now())
.expect("grantee opens the granted envelope");
assert!(descriptor.has_tag("nrpc:cross"));
assert!(!descriptor.has_tag("nrpc:open"));
}
#[tokio::test]
async fn overlapping_grants_emit_two_independently_decryptable_envelopes() {
let (p, _h, _dir, entity, org_b) = granted_provider(0x71, "two", "cross").await;
let org_a = OrgKeypair::from_bytes([0x7Au8; 32]);
let cap = CapabilityAuthorityId::for_tag("nrpc:cross");
let issue = || {
let (g, s) = OrgCapabilityGrant::try_issue(
&org_b,
org_a.org_id(),
cap,
GrantRights::DISCOVER,
GrantTargetScope::ExactNode(entity.clone()),
3600,
)
.expect("issue");
(g, s.expect("secret"))
};
let (g1, s1) = issue();
let (g2, s2) = issue();
assert_ne!(g1.grant_id, g2.grant_id, "distinct grant ids");
let (o1, o2) = (copy_secret(&s1), copy_secret(&s2));
p.install_provider_grant_audience(g1.clone(), s1)
.expect("install g1");
p.install_provider_grant_audience(g2.clone(), s2)
.expect("install g2");
assert!(
converge_scoped_count(&p, 2).await,
"two overlapping grants emit two envelopes",
);
let scoped = p.announcement_scoped_for_send_for_test();
let now = unix_now();
use net::adapter::net::behavior::org_scoped_ann::ScopedCapabilityAnnouncement;
let envs: Vec<ScopedCapabilityAnnouncement> = scoped
.iter()
.map(|b| ScopedCapabilityAnnouncement::from_bytes(b).expect("decode"))
.collect();
let e1 = envs
.iter()
.find(|e| e.grant_id() == &g1.grant_id)
.expect("g1 envelope present");
let e2 = envs
.iter()
.find(|e| e.grant_id() == &g2.grant_id)
.expect("g2 envelope present");
assert!(open_granted_envelope(&e1.to_bytes(), &g1, &o1, org_a.org_id(), now).is_some());
assert!(open_granted_envelope(&e2.to_bytes(), &g2, &o2, org_a.org_id(), now).is_some());
assert!(e1.open_with(o1.discovery_key()).is_ok(), "K1 opens E1",);
assert!(
e1.open_with(o2.discovery_key()).is_err(),
"K2 cannot open E1",
);
assert!(e2.open_with(o2.discovery_key()).is_ok(), "K2 opens E2",);
assert!(
e2.open_with(o1.discovery_key()).is_err(),
"K1 cannot open E2",
);
}
#[tokio::test]
async fn an_unrelated_capability_grant_emits_no_granted_envelope() {
let (p, _h, _dir, entity, org_b) = granted_provider(0x72, "none", "cross").await;
let org_a = OrgKeypair::from_bytes([0x7Au8; 32]);
let (matching, matching_secret) = OrgCapabilityGrant::try_issue(
&org_b,
org_a.org_id(),
CapabilityAuthorityId::for_tag("nrpc:cross"),
GrantRights::DISCOVER,
GrantTargetScope::ExactNode(entity.clone()),
3600,
)
.expect("issue matching");
let matching_id = matching.grant_id;
p.install_provider_grant_audience(matching, matching_secret.expect("secret"))
.expect("install matching");
assert!(
converge_scoped_count(&p, 1).await,
"precondition: a grant that DOES name the local service emits one \
envelope — without this, the zero asserted below is indistinguishable \
from a provider that never emitted anything",
);
assert!(
p.remove_provider_grant_audience(&matching_id),
"the matching grant must actually be removed",
);
assert!(
converge_scoped_count(&p, 0).await,
"removing the matching grant retracts its envelope",
);
let (grant, secret) = OrgCapabilityGrant::try_issue(
&org_b,
org_a.org_id(),
CapabilityAuthorityId::for_tag("nrpc:unrelated"),
GrantRights::DISCOVER,
GrantTargetScope::ExactNode(entity.clone()),
3600,
)
.expect("issue");
p.install_provider_grant_audience(grant, secret.expect("secret"))
.expect("install");
assert!(
converge_scoped_count(&p, 0).await,
"an unrelated-capability grant emits no granted envelope",
);
}
#[tokio::test]
async fn a_granted_envelope_never_outlives_its_grant() {
let (p, _h, _dir, entity, org_b) = granted_provider(0x73, "ttl", "cross").await;
let org_a = OrgKeypair::from_bytes([0x7Au8; 32]);
let (grant, secret) = OrgCapabilityGrant::try_issue(
&org_b,
org_a.org_id(),
CapabilityAuthorityId::for_tag("nrpc:cross"),
GrantRights::DISCOVER,
GrantTargetScope::ExactNode(entity.clone()),
120, )
.expect("issue");
p.install_provider_grant_audience(grant.clone(), secret.expect("secret"))
.expect("install");
assert!(converge_scoped_count(&p, 1).await, "one granted envelope");
let scoped = p.announcement_scoped_for_send_for_test();
let env =
net::adapter::net::behavior::org_scoped_ann::ScopedCapabilityAnnouncement::from_bytes(
&scoped[0],
)
.expect("decode");
assert!(
env.expires_at() <= grant.not_after,
"envelope expiry {} must not outlive the grant not_after {}",
env.expires_at(),
grant.not_after,
);
assert!(
env.expires_at() <= unix_now() + 200,
"the short grant TTL clamped the envelope expiry below the announce TTL",
);
}
#[tokio::test]
async fn removing_a_provider_grant_refuses_the_cached_granted_envelope() {
let (p, _h, _dir, entity, org_b) = granted_provider(0x74, "remove", "cross").await;
let org_a = OrgKeypair::from_bytes([0x7Au8; 32]);
let (grant, secret) = OrgCapabilityGrant::try_issue(
&org_b,
org_a.org_id(),
CapabilityAuthorityId::for_tag("nrpc:cross"),
GrantRights::DISCOVER,
GrantTargetScope::ExactNode(entity.clone()),
3600,
)
.expect("issue");
let grant_id = grant.grant_id;
p.install_provider_grant_audience(grant, secret.expect("secret"))
.expect("install");
assert!(converge_scoped_count(&p, 1).await, "one granted envelope");
assert!(p.remove_provider_grant_audience(&grant_id));
assert!(
p.announcement_scoped_for_send_for_test().is_empty(),
"the cached granted envelope is refused after the grant is removed",
);
p.announce_capabilities(CapabilitySet::new())
.await
.expect("announce");
assert!(
p.announcement_scoped_for_send_for_test().is_empty(),
"the rebuilt emission carries no granted envelope",
);
}
async fn adopted_node(
seed: u8,
org: &OrgKeypair,
tag: &str,
) -> (Arc<MeshNode>, std::path::PathBuf) {
let n = build_node_with(EntityKeypair::from_bytes([seed; 32])).await;
let entity = n.entity_id().clone();
let cert = OrgMembershipCert::try_issue(org, entity.clone(), 1, 3600).expect("cert");
let dir = std::env::temp_dir().join(format!("net-oa34b2-cons-{tag}-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let authority = NodeAuthority::adopt(&dir, cert, &entity, 0, None).expect("adopt");
n.install_node_authority(Arc::new(authority))
.expect("install authority");
(n, dir)
}
fn granted_envelope_bytes(
provider_kp: &EntityKeypair,
org_b: &OrgKeypair,
grant: &OrgCapabilityGrant,
secret: &OrgAudienceSecret,
svc: &str,
now: u64,
) -> Vec<u8> {
use net::adapter::net::behavior::org_scoped_ann::ScopedCapabilityAnnouncement;
let cert = OrgMembershipCert::try_issue(org_b, provider_kp.entity_id().clone(), 1, 3600)
.expect("cert");
let descriptor = CapabilitySet::new()
.add_tag(format!("nrpc:{svc}"))
.to_bytes_compact();
ScopedCapabilityAnnouncement::build_granted(
provider_kp,
org_b.org_id(),
cert,
grant.grant_id,
secret.audience_handle,
secret.discovery_key(),
1,
now + 600,
&descriptor,
)
.expect("build granted envelope")
.to_bytes()
}
fn cross_org_grant(
org_b: &OrgKeypair,
org_a: &OrgKeypair,
p: &EntityId,
svc: &str,
) -> (OrgCapabilityGrant, OrgAudienceSecret) {
let (g, s) = OrgCapabilityGrant::try_issue(
org_b,
org_a.org_id(),
CapabilityAuthorityId::for_tag(&format!("nrpc:{svc}")),
GrantRights::DISCOVER,
GrantTargetScope::ExactNode(p.clone()),
3600,
)
.expect("issue cross-org grant");
(g, s.expect("secret"))
}
#[tokio::test]
async fn an_inbound_granted_announcement_is_verified_and_stored() {
let org_b = OrgKeypair::from_bytes([0x42u8; 32]);
let org_a = OrgKeypair::from_bytes([0x7Au8; 32]);
let provider = EntityKeypair::from_bytes([0x90u8; 32]);
let p_entity = provider.entity_id().clone();
let now = unix_now();
let (grant, secret) = cross_org_grant(&org_b, &org_a, &p_entity, "cross");
let grant_id = grant.grant_id;
let env = granted_envelope_bytes(&provider, &org_b, &grant, &secret, "cross", now);
let (c, _c_dir) = adopted_node(0x91, &org_a, "resolve").await;
c.install_consumer_grant_audience(grant.clone(), copy_secret(&secret))
.expect("install consumer grant");
c.ingest_scoped_announcement_for_test(&env);
assert_eq!(
c.scoped_granted_providers_for_test(&grant_id, now),
vec![p_entity.clone()],
"the grantee opens and resolves P under the grant",
);
let (d, _d_dir) = adopted_node(0x92, &org_a, "nostore").await;
d.ingest_scoped_announcement_for_test(&env);
d.install_consumer_grant_audience(grant.clone(), copy_secret(&secret))
.expect("install after the drop");
assert!(
d.scoped_granted_providers_for_test(&grant_id, now)
.is_empty(),
"a node without the pair at ingest time stored nothing",
);
}
#[tokio::test]
async fn the_ingest_selector_drops_a_grant_id_it_does_not_hold() {
let org_b = OrgKeypair::from_bytes([0x42u8; 32]);
let org_a = OrgKeypair::from_bytes([0x7Au8; 32]);
let provider = EntityKeypair::from_bytes([0x93u8; 32]);
let p_entity = provider.entity_id().clone();
let now = unix_now();
let (g1, s1) = cross_org_grant(&org_b, &org_a, &p_entity, "cross");
let (g2, s2) = cross_org_grant(&org_b, &org_a, &p_entity, "other");
let env1 = granted_envelope_bytes(&provider, &org_b, &g1, &s1, "cross", now);
let (c, _dir) = adopted_node(0x94, &org_a, "wrongid").await;
c.install_consumer_grant_audience(g2, s2)
.expect("install g2");
c.ingest_scoped_announcement_for_test(&env1);
c.install_consumer_grant_audience(g1.clone(), copy_secret(&s1))
.expect("install g1");
assert!(
c.scoped_granted_providers_for_test(&g1.grant_id, now)
.is_empty(),
"an envelope whose grant id the node did not hold was dropped, not stored",
);
}
#[tokio::test]
async fn removing_the_consumer_credential_hides_the_stored_granted_record() {
let org_b = OrgKeypair::from_bytes([0x42u8; 32]);
let org_a = OrgKeypair::from_bytes([0x7Au8; 32]);
let provider = EntityKeypair::from_bytes([0x95u8; 32]);
let p_entity = provider.entity_id().clone();
let now = unix_now();
let (grant, secret) = cross_org_grant(&org_b, &org_a, &p_entity, "cross");
let grant_id = grant.grant_id;
let env = granted_envelope_bytes(&provider, &org_b, &grant, &secret, "cross", now);
let (c, _dir) = adopted_node(0x96, &org_a, "hide").await;
c.install_consumer_grant_audience(grant.clone(), copy_secret(&secret))
.expect("install");
c.ingest_scoped_announcement_for_test(&env);
assert_eq!(
c.scoped_granted_providers_for_test(&grant_id, now),
vec![p_entity.clone()],
"resolves before removal",
);
assert!(c.remove_consumer_grant_audience(&grant_id));
assert!(
c.scoped_granted_providers_for_test(&grant_id, now)
.is_empty(),
"removing the consumer credential retracts the record at query time",
);
c.install_consumer_grant_audience(grant.clone(), copy_secret(&secret))
.expect("re-install");
assert_eq!(
c.scoped_granted_providers_for_test(&grant_id, now),
vec![p_entity],
"the record was hidden, not evicted — re-installing re-exposes it",
);
}
#[tokio::test]
async fn a_consumer_credential_replacement_racing_the_granted_insert_is_refused() {
let org_b = OrgKeypair::from_bytes([0x42u8; 32]);
let org_a = OrgKeypair::from_bytes([0x7Au8; 32]);
let provider = EntityKeypair::from_bytes([0x97u8; 32]);
let p_entity = provider.entity_id().clone();
let now = unix_now();
let (grant, secret) = cross_org_grant(&org_b, &org_a, &p_entity, "cross");
let grant_id = grant.grant_id;
let env = granted_envelope_bytes(&provider, &org_b, &grant, &secret, "cross", now);
let (c, _dir) = adopted_node(0x98, &org_a, "race").await;
c.install_consumer_grant_audience(grant.clone(), copy_secret(&secret))
.expect("install target grant");
let (unrelated, unrelated_secret) = OrgCapabilityGrant::try_issue(
&org_b,
org_a.org_id(),
CapabilityAuthorityId::for_tag("nrpc:unrelated"),
GrantRights::DISCOVER,
GrantTargetScope::AnyNodeOwnedBy(org_b.org_id()),
3600,
)
.expect("issue unrelated");
let unrelated_secret = unrelated_secret.expect("secret");
let pending = parking_lot::Mutex::new(Some((unrelated, unrelated_secret)));
let c_probe = c.clone();
let probe = move || {
if let Some((g, s)) = pending.lock().take() {
c_probe
.install_consumer_grant_audience(g, s)
.expect("probe install");
}
};
c.ingest_scoped_announcement_probed_for_test(&env, &probe);
assert!(
c.scoped_granted_providers_for_test(&grant_id, now)
.is_empty(),
"the raced insert is refused when the consumer snapshot moved during verify",
);
c.ingest_scoped_announcement_for_test(&env);
assert_eq!(
c.scoped_granted_providers_for_test(&grant_id, now),
vec![p_entity],
"a clean re-ingest against the settled snapshot resolves P",
);
}
#[tokio::test]
async fn a_same_org_audience_rotation_refuses_the_stale_scoped_envelope() {
use net::adapter::net::behavior::org::{OrgKeypair, OrgMembershipCert};
use net::adapter::net::behavior::org_authority::NodeAuthority;
use net::adapter::net::behavior::org_revocation::OrgRevocationState;
use net::adapter::net::behavior::org_scoped_ann::ScopedCapabilityAnnouncement;
use net::adapter::net::behavior::org_scoped_ingest::{
verify_scoped_ingest, AudienceAuthority, ScopedIngestContext,
};
let server = build_node_with(EntityKeypair::from_bytes([0x54u8; 32])).await;
let node_entity = server.entity_id().clone();
let org = OrgKeypair::from_bytes([0x77u8; 32]);
let cert = OrgMembershipCert::try_issue(&org, node_entity.clone(), 1, 3600).expect("cert C");
let owner_org = cert.org_id;
let dir_a = std::env::temp_dir().join(format!("net-b2-rot-a-{}", std::process::id()));
let dir_b = std::env::temp_dir().join(format!("net-b2-rot-b-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir_a);
let _ = std::fs::remove_dir_all(&dir_b);
let authority_a = Arc::new(
NodeAuthority::adopt(&dir_a, cert.clone(), &node_entity, 0, None).expect("adopt A"),
);
let authority_b = Arc::new(
NodeAuthority::adopt(&dir_b, cert.clone(), &node_entity, 0, None).expect("adopt B"),
);
let handle_a = authority_a.audience.audience_handle;
let key_a = *authority_a.audience.discovery_key();
let handle_b = authority_b.audience.audience_handle;
let key_b = *authority_b.audience.discovery_key();
assert_ne!(key_a, key_b, "the rotation must change the audience key");
server
.install_node_authority(authority_a)
.expect("install A");
server
.set_owner_cert_emission(true)
.expect("enable owner-cert emission");
let _secret = server
.serve_rpc_owner_scoped("secret", Arc::new(TrivialHandler), Arc::new(|_| true))
.expect("owner-scoped serve");
server
.announce_capabilities(CapabilitySet::new())
.await
.expect("announce under A");
assert!(
wait_until(Duration::from_secs(5), || {
server.announcement_scoped_for_send_for_test().len() == 1
})
.await,
"E1 published under authority A",
);
let e1 = server.announcement_scoped_for_send_for_test();
let env1 = ScopedCapabilityAnnouncement::from_bytes(&e1[0]).expect("decode E1");
server
.install_node_authority(authority_b)
.expect("install B (same-org rotation)");
let after_rotation = server.announcement_scoped_for_send_for_test();
assert!(
after_rotation.is_empty(),
"a rotation must refuse the stale scoped envelope until the emission is rebuilt",
);
server
.announce_capabilities(CapabilitySet::new())
.await
.expect("announce under B");
assert!(
wait_until(Duration::from_secs(5), || {
server.announcement_scoped_for_send_for_test().len() == 1
})
.await,
"E2 published under authority B",
);
let e2 = server.announcement_scoped_for_send_for_test();
let env2 = ScopedCapabilityAnnouncement::from_bytes(&e2[0]).expect("decode E2");
let floors = OrgRevocationState::empty();
let now_secs = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.expect("clock")
.as_secs();
let v1 = verify_scoped_ingest(
&env1,
&AudienceAuthority::Owner {
owner_org,
audience_handle: handle_a,
discovery_key: &key_a,
},
&ScopedIngestContext {
local_owner_org: owner_org,
floors: &floors,
now_secs,
skew_secs: 5,
local_member: None,
},
)
.expect("E1 opens under K1");
assert!(CapabilitySet::from_bytes(v1.descriptor())
.expect("E1 descriptor")
.has_tag("nrpc:secret"));
let v2 = verify_scoped_ingest(
&env2,
&AudienceAuthority::Owner {
owner_org,
audience_handle: handle_b,
discovery_key: &key_b,
},
&ScopedIngestContext {
local_owner_org: owner_org,
floors: &floors,
now_secs,
skew_secs: 5,
local_member: None,
},
)
.expect("E2 opens under K2");
assert!(CapabilitySet::from_bytes(v2.descriptor())
.expect("E2 descriptor")
.has_tag("nrpc:secret"));
assert!(
verify_scoped_ingest(
&env2,
&AudienceAuthority::Owner {
owner_org,
audience_handle: handle_b,
discovery_key: &key_a,
},
&ScopedIngestContext {
local_owner_org: owner_org,
floors: &floors,
now_secs,
skew_secs: 5,
local_member: None,
},
)
.is_err(),
"E2 sealed under the new key must not open under the rotated-out K1",
);
}
#[tokio::test]
async fn an_inbound_owner_scoped_announcement_is_verified_and_stored() {
use net::adapter::net::behavior::org::{OrgKeypair, OrgMembershipCert};
use net::adapter::net::behavior::org_authority::NodeAuthority;
use net::adapter::net::behavior::org_scoped_ann::ScopedCapabilityAnnouncement;
let node = build_node_with(EntityKeypair::from_bytes([0x60u8; 32])).await;
let node_entity = node.entity_id().clone();
let org = OrgKeypair::from_bytes([0x88u8; 32]);
let node_cert =
OrgMembershipCert::try_issue(&org, node_entity.clone(), 1, 3600).expect("node cert");
let dir = std::env::temp_dir().join(format!("net-oa35-ingest-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let authority =
Arc::new(NodeAuthority::adopt(&dir, node_cert, &node_entity, 0, None).expect("adopt"));
let handle = authority.audience.audience_handle;
let key = *authority.audience.discovery_key();
node.install_node_authority(authority).expect("install");
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.expect("clock")
.as_secs();
let descriptor = CapabilitySet::new()
.add_tag("nrpc:peer-secret")
.to_bytes_compact();
let make_envelope = |seed: u8, disc_key: [u8; 32], expires_at: u64| -> (EntityId, Vec<u8>) {
let provider_kp = EntityKeypair::from_bytes([seed; 32]);
let provider_entity = provider_kp.entity_id().clone();
let cert = OrgMembershipCert::try_issue(&org, provider_entity.clone(), 1, 3600)
.expect("provider cert");
let env = ScopedCapabilityAnnouncement::build_owner(
&provider_kp,
org.org_id(),
cert,
handle,
&disc_key,
1,
expires_at,
&descriptor,
)
.expect("build owner envelope");
(provider_entity, env.to_bytes())
};
let (good_provider, good_bytes) = make_envelope(0x61, key, now + 3600);
node.ingest_scoped_announcement_for_test(&good_bytes);
assert!(
node.scoped_owner_providers_for_test(now)
.iter()
.any(|p| p == &good_provider),
"the verified owner-scoped provider is exposed in the private-discovery store",
);
let (bad_provider, bad_bytes) = make_envelope(0x62, [0x99u8; 32], now + 3600);
node.ingest_scoped_announcement_for_test(&bad_bytes);
assert!(
!node
.scoped_owner_providers_for_test(now)
.iter()
.any(|p| p == &bad_provider),
"a wrong-audience envelope is refused and never stored",
);
let (exp_provider, exp_bytes) = make_envelope(0x63, key, now.saturating_sub(10));
node.ingest_scoped_announcement_for_test(&exp_bytes);
assert!(
!node
.scoped_owner_providers_for_test(now)
.iter()
.any(|p| p == &exp_provider),
"an expired envelope is refused at ingest",
);
}
#[tokio::test]
async fn a_floor_publish_racing_the_scoped_insert_is_refused_then_recovers() {
use net::adapter::net::behavior::org::{OrgKeypair, OrgMembershipCert, OrgRevocationBundle};
use net::adapter::net::behavior::org_authority::NodeAuthority;
use net::adapter::net::behavior::org_scoped_ann::ScopedCapabilityAnnouncement;
use std::collections::BTreeMap;
let node = build_node_with(EntityKeypair::from_bytes([0x70u8; 32])).await;
let node_entity = node.entity_id().clone();
let org = OrgKeypair::from_bytes([0x89u8; 32]);
let node_cert =
OrgMembershipCert::try_issue(&org, node_entity.clone(), 1, 3600).expect("node cert");
let dir = std::env::temp_dir().join(format!("net-oa35-race-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let authority =
Arc::new(NodeAuthority::adopt(&dir, node_cert, &node_entity, 0, None).expect("adopt"));
let store = authority.revocation.clone();
let handle = authority.audience.audience_handle;
let key = *authority.audience.discovery_key();
node.install_node_authority(authority)
.expect("install authority");
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.expect("clock")
.as_secs();
let descriptor = CapabilitySet::new()
.add_tag("nrpc:peer-secret")
.to_bytes_compact();
let make_envelope = |seed: u8| -> (EntityId, Vec<u8>) {
let provider_kp = EntityKeypair::from_bytes([seed; 32]);
let provider_entity = provider_kp.entity_id().clone();
let cert = OrgMembershipCert::try_issue(&org, provider_entity.clone(), 1, 3600)
.expect("provider cert");
let env = ScopedCapabilityAnnouncement::build_owner(
&provider_kp,
org.org_id(),
cert,
handle,
&key,
1,
now + 3600,
&descriptor,
)
.expect("build owner envelope");
(provider_entity, env.to_bytes())
};
let (clean_provider, clean_bytes) = make_envelope(0x71);
node.ingest_scoped_announcement_for_test(&clean_bytes);
assert!(
node.scoped_owner_providers_for_test(now)
.iter()
.any(|p| p == &clean_provider),
"a valid owner-scoped envelope lands under an installed revocation store",
);
let (raced_provider, raced_bytes) = make_envelope(0x72);
let unrelated_member = EntityKeypair::from_bytes([0xAAu8; 32]).entity_id().clone();
let publisher: std::sync::OnceLock<std::thread::JoinHandle<()>> = std::sync::OnceLock::new();
let race_probe = || {
let store_for_publish = store.clone();
let member = unrelated_member.clone();
let handle = std::thread::spawn(move || {
let org = OrgKeypair::from_bytes([0x89u8; 32]);
let mut floors_map = BTreeMap::new();
floors_map.insert(member, 5u32);
let bundle =
OrgRevocationBundle::try_issue(&org, &floors_map).expect("issue race bundle");
store_for_publish
.apply_bundle(&bundle)
.expect("apply race floor");
});
let _ = publisher.set(handle);
while store.snapshot().floor_for(&org.org_id(), &unrelated_member) < 5 {
std::thread::yield_now();
}
};
node.ingest_scoped_announcement_probed_for_test(&raced_bytes, &race_probe);
publisher
.into_inner()
.expect("the probe published")
.join()
.expect("the floor publish completes once the gate is released");
assert!(
!node
.scoped_owner_providers_for_test(now)
.iter()
.any(|p| p == &raced_provider),
"an insert racing a floor publish is refused — the raced provider never enters the store",
);
node.ingest_scoped_announcement_for_test(&raced_bytes);
assert!(
node.scoped_owner_providers_for_test(now)
.iter()
.any(|p| p == &raced_provider),
"the identical envelope re-announced against the settled view lands cleanly",
);
}
#[tokio::test]
async fn an_owner_scoped_announcement_floods_opaquely_through_a_relay_to_the_audience() {
use net::adapter::net::behavior::org::{OrgKeypair, OrgMembershipCert};
use net::adapter::net::behavior::org_authority::{NodeAuthority, OwnerAudienceCredential};
let p = build_node_fast_announce(EntityKeypair::from_bytes([0x80u8; 32])).await;
let r = build_node_fast_announce(EntityKeypair::from_bytes([0x81u8; 32])).await;
let c = build_node_fast_announce(EntityKeypair::from_bytes([0x82u8; 32])).await;
let p_entity = p.entity_id().clone();
let c_entity = c.entity_id().clone();
let org = OrgKeypair::from_bytes([0x8Au8; 32]);
let base = std::env::temp_dir().join(format!("net-oa35-relay-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&base);
let p_cert = OrgMembershipCert::try_issue(&org, p_entity.clone(), 1, 3600).expect("P cert");
let p_authority = Arc::new(
NodeAuthority::adopt(&base.join("p"), p_cert, &p_entity, 0, None).expect("adopt P"),
);
let shared_audience = p_authority.audience.encode_config();
p.install_node_authority(p_authority)
.expect("install P authority");
p.set_owner_cert_emission(true).expect("enable P emission");
let _svc = p
.serve_rpc_owner_scoped("secret", Arc::new(TrivialHandler), Arc::new(|_| true))
.expect("P owner-scoped serve");
let c_cert = OrgMembershipCert::try_issue(&org, c_entity.clone(), 1, 3600).expect("C cert");
let mut c_authority =
NodeAuthority::adopt(&base.join("c"), c_cert, &c_entity, 0, None).expect("adopt C");
c_authority.audience =
OwnerAudienceCredential::decode_config(&shared_audience).expect("decode shared audience");
c.install_node_authority(Arc::new(c_authority))
.expect("install C authority");
connect_no_start(&p, &r).await;
connect_no_start(&r, &c).await;
p.start();
r.start();
c.start();
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.expect("clock")
.as_secs();
assert!(
wait_until(Duration::from_secs(5), || {
p.announcement_scoped_for_send_for_test().len() == 1
})
.await,
"P emits exactly one owner-scoped envelope",
);
let mut c_resolved_p = false;
for _ in 0..40 {
p.announce_capabilities(CapabilitySet::new()).await.ok();
if c.scoped_owner_providers_for_test(now)
.iter()
.any(|prov| prov == &p_entity)
{
c_resolved_p = true;
break;
}
tokio::time::sleep(Duration::from_millis(150)).await;
}
assert!(
c_resolved_p,
"C resolves P through the relay despite having no direct session with P",
);
assert!(
r.scoped_relay_gate_len_for_test() >= 1,
"the relay admitted and forwarded the opaque envelope",
);
assert!(
r.scoped_owner_providers_for_test(now).is_empty(),
"the authority-less relay forwards but never decrypts or stores the envelope",
);
}
#[tokio::test]
async fn a_granted_capability_floods_opaquely_through_a_relay_to_the_grantee() {
let p = build_node_fast_announce(EntityKeypair::from_bytes([0x83u8; 32])).await;
let r = build_node_fast_announce(EntityKeypair::from_bytes([0x84u8; 32])).await;
let a = build_node_fast_announce(EntityKeypair::from_bytes([0x85u8; 32])).await;
let p_entity = p.entity_id().clone();
let org_b = OrgKeypair::from_bytes([0x8Bu8; 32]); let org_a = OrgKeypair::from_bytes([0x8Au8; 32]); let base = std::env::temp_dir().join(format!("net-oa34b2-relay-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&base);
let (grant, secret) = cross_org_grant(&org_b, &org_a, &p_entity, "cross");
let grant_id = grant.grant_id;
let p_cert = OrgMembershipCert::try_issue(&org_b, p_entity.clone(), 1, 3600).expect("P cert");
let p_authority = Arc::new(
NodeAuthority::adopt(&base.join("p"), p_cert, &p_entity, 0, None).expect("adopt P"),
);
p.install_node_authority(p_authority)
.expect("install P authority");
p.set_owner_cert_emission(true).expect("enable P emission");
let _svc = p
.serve_rpc_granted("cross", Arc::new(TrivialHandler), Arc::new(|_| true))
.expect("P granted serve");
p.install_provider_grant_audience(grant.clone(), copy_secret(&secret))
.expect("install P provider grant");
let a_entity = a.entity_id().clone();
let a_cert = OrgMembershipCert::try_issue(&org_a, a_entity.clone(), 1, 3600).expect("A cert");
let a_authority = Arc::new(
NodeAuthority::adopt(&base.join("a"), a_cert, &a_entity, 0, None).expect("adopt A"),
);
a.install_node_authority(a_authority)
.expect("install A authority");
a.install_consumer_grant_audience(grant.clone(), copy_secret(&secret))
.expect("install A consumer grant");
connect_no_start(&p, &r).await;
connect_no_start(&r, &a).await;
p.start();
r.start();
a.start();
let now = unix_now();
assert!(
wait_until(Duration::from_secs(5), || {
p.announcement_scoped_for_send_for_test().len() == 1
})
.await,
"P emits exactly one granted envelope",
);
let mut a_resolved_p = false;
for _ in 0..40 {
p.announce_capabilities(CapabilitySet::new()).await.ok();
if a.scoped_granted_providers_for_test(&grant_id, now)
.iter()
.any(|prov| prov == &p_entity)
{
a_resolved_p = true;
break;
}
tokio::time::sleep(Duration::from_millis(150)).await;
}
assert!(
a_resolved_p,
"the grantee A resolves P through the relay despite different orgs and no direct session",
);
assert!(
r.scoped_relay_gate_len_for_test() >= 1,
"the relay admitted and forwarded the opaque envelope",
);
assert!(
r.scoped_owner_providers_for_test(now).is_empty()
&& r.scoped_granted_providers_for_test(&grant_id, now)
.is_empty(),
"the authority-less non-grantee relay forwards but never decrypts or stores",
);
assert!(
p.local_announcement_for_test()
.map(|ann| !ann.capabilities.has_tag("nrpc:cross"))
.unwrap_or(false),
"the granted tag never appears in P's plaintext announcement",
);
}
#[tokio::test]
async fn a_send_after_a_visibility_change_refuses_the_stale_emission() {
let server = build_node_with(EntityKeypair::from_bytes([0x52u8; 32])).await;
let (_org_b, _dir) = install_authority(&server, "vis-race");
let _svc = server
.serve_rpc("svc", Arc::new(TrivialHandler))
.expect("public serve");
server
.announce_capabilities(CapabilitySet::new())
.await
.expect("announce");
assert!(
wait_until(Duration::from_secs(5), || {
server.announcement_bytes_for_send_for_test().is_some()
})
.await,
"an emission is published",
);
server.test_advance_visibility_generation();
assert!(
server.announcement_bytes_for_send_for_test().is_none(),
"a send after a visibility change must not ship the stale emission",
);
}
#[tokio::test]
async fn a_visibility_change_during_serialization_refuses_the_stale_bytes() {
let server = build_node_with(EntityKeypair::from_bytes([0x53u8; 32])).await;
let (_org_b, _dir) = install_authority(&server, "vis-serialize-race");
let _svc = server
.serve_rpc("svc", Arc::new(TrivialHandler))
.expect("public serve");
server
.announce_capabilities(CapabilitySet::new())
.await
.expect("announce");
assert!(
wait_until(Duration::from_secs(5), || {
server.announcement_bytes_for_send_for_test().is_some()
})
.await,
"an emission is published",
);
let server_probe = server.clone();
let probe = move || server_probe.test_advance_visibility_generation();
assert!(
server
.announcement_bytes_for_send_probed_for_test(&probe)
.is_none(),
"a visibility change during serialization must refuse the stale bytes",
);
}
#[tokio::test]
async fn a_real_registry_transition_during_serialization_refuses_the_stale_bytes() {
let server = build_node_with(EntityKeypair::from_bytes([0x54u8; 32])).await;
let (_org_b, _dir) = install_authority(&server, "vis-real-transition");
let public_handle = server
.serve_rpc("svc", Arc::new(TrivialHandler))
.expect("public serve");
server
.announce_capabilities(CapabilitySet::new())
.await
.expect("announce");
assert!(
wait_until(Duration::from_secs(5), || {
server.announcement_bytes_for_send_for_test().is_some()
})
.await,
"an emission is published",
);
let published = server
.announcement_bytes_for_send_for_test()
.expect("published emission");
let decoded = CapabilityAnnouncement::from_bytes(&published).expect("decode");
assert!(
decoded
.capabilities
.tags
.iter()
.any(|t| t.to_string() == "nrpc:svc"),
"the published plaintext must carry nrpc:svc before the transition",
);
let slot = Arc::new(parking_lot::Mutex::new(Some(public_handle)));
let installed: Arc<parking_lot::Mutex<Option<net::adapter::net::mesh_rpc::ServeHandle>>> =
Arc::new(parking_lot::Mutex::new(None));
let server_probe = server.clone();
let slot_probe = Arc::clone(&slot);
let installed_probe = Arc::clone(&installed);
let probe = move || {
if let Some(handle) = slot_probe.lock().take() {
drop(handle);
let replacement = server_probe
.serve_rpc_owner_scoped("svc", Arc::new(TrivialHandler), Arc::new(|_| true))
.expect("owner-scoped re-serve");
*installed_probe.lock() = Some(replacement);
}
};
assert!(
server
.announcement_bytes_for_send_probed_for_test(&probe)
.is_none(),
"a real Public -> OwnerScoped transition during serialization must refuse the stale plaintext bytes",
);
assert!(
installed.lock().is_some(),
"the probe must actually have performed the transition",
);
}
#[tokio::test]
async fn grant_audience_registries_install_and_remove_on_a_live_node() {
let bare = build_node_with(EntityKeypair::from_bytes([0x60u8; 32])).await;
let org_b = OrgKeypair::from_bytes([0x42u8; 32]); let org_a = OrgKeypair::from_bytes([0x6Au8; 32]); let (provider_grant, provider_secret) = OrgCapabilityGrant::try_issue(
&org_b,
org_a.org_id(),
CapabilityAuthorityId::for_tag("nrpc:reconcile"),
GrantRights::DISCOVER,
GrantTargetScope::ExactNode(bare.entity_id().clone()),
3600,
)
.expect("issue provider grant");
let provider_secret = provider_secret.expect("DISCOVER mints a secret");
assert_eq!(
bare.install_provider_grant_audience(provider_grant.clone(), provider_secret)
.unwrap_err(),
GrantAudienceInstallError::NoAuthority,
"a node without authority refuses a grant-audience install",
);
let server = build_node_with(EntityKeypair::from_bytes([0x61u8; 32])).await;
let node_entity = server.entity_id().clone();
let node_cert =
OrgMembershipCert::try_issue(&org_b, node_entity.clone(), 1, 3600).expect("node cert");
let dir = std::env::temp_dir().join(format!("net-oa34b2-store-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let authority =
NodeAuthority::adopt(&dir, node_cert, &node_entity, 0, None).expect("adopt authority");
server
.install_node_authority(Arc::new(authority))
.expect("install authority");
let (p_grant, p_secret) = OrgCapabilityGrant::try_issue(
&org_b,
org_a.org_id(),
CapabilityAuthorityId::for_tag("nrpc:reconcile"),
GrantRights::DISCOVER.union(GrantRights::INVOKE),
GrantTargetScope::ExactNode(node_entity.clone()),
3600,
)
.expect("issue provider grant");
let p_secret = p_secret.expect("secret");
let p_grant_id = p_grant.grant_id;
let p_secret_copy = copy_secret(&p_secret);
assert_eq!(
server
.install_provider_grant_audience(p_grant.clone(), p_secret)
.expect("install provider grant"),
GrantAudienceInstalled::Installed,
);
assert_eq!(server.provider_grant_audiences_len_for_test(), 1);
assert_eq!(
server
.install_provider_grant_audience(p_grant.clone(), p_secret_copy)
.expect("idempotent re-install"),
GrantAudienceInstalled::AlreadyPresent,
);
assert_eq!(server.provider_grant_audiences_len_for_test(), 1);
let (foreign_grant, foreign_secret) = OrgCapabilityGrant::try_issue(
&org_a,
org_a.org_id(),
CapabilityAuthorityId::for_tag("nrpc:reconcile"),
GrantRights::DISCOVER,
GrantTargetScope::AnyNodeOwnedBy(org_a.org_id()),
3600,
)
.expect("issue foreign grant");
assert_eq!(
server
.install_provider_grant_audience(foreign_grant, foreign_secret.expect("secret"))
.unwrap_err(),
GrantAudienceInstallError::WrongProviderIssuer,
);
let (c_grant, c_secret) = OrgCapabilityGrant::try_issue(
&org_a,
org_b.org_id(),
CapabilityAuthorityId::for_tag("nrpc:remote-svc"),
GrantRights::DISCOVER,
GrantTargetScope::AnyNodeOwnedBy(org_a.org_id()),
3600,
)
.expect("issue consumer grant");
let c_grant_id = c_grant.grant_id;
assert_eq!(
server
.install_consumer_grant_audience(c_grant, c_secret.expect("secret"))
.expect("install consumer grant"),
GrantAudienceInstalled::Installed,
);
assert_eq!(server.consumer_grant_audiences_len_for_test(), 1);
assert_eq!(server.provider_grant_audiences_len_for_test(), 1);
assert!(server.remove_provider_grant_audience(&p_grant_id));
assert_eq!(server.provider_grant_audiences_len_for_test(), 0);
assert!(!server.remove_provider_grant_audience(&p_grant_id));
assert_eq!(server.consumer_grant_audiences_len_for_test(), 1);
assert!(server.remove_consumer_grant_audience(&c_grant_id));
assert_eq!(server.consumer_grant_audiences_len_for_test(), 0);
}
fn cross_org_invoke_intent(
caller_kp: EntityKeypair,
org_a: &OrgKeypair,
org_b: &OrgKeypair,
provider: EntityId,
service: &str,
target_scope: GrantTargetScope,
) -> OrgProofIntent {
let cap = CapabilityAuthorityId::for_tag(&format!("nrpc:{service}"));
let (grant, secret) = OrgCapabilityGrant::try_issue(
org_b,
org_a.org_id(),
cap,
GrantRights::INVOKE,
target_scope,
3600,
)
.expect("issue cross-org INVOKE grant");
assert!(
secret.is_none(),
"an INVOKE-only grant carries no audience material by construction",
);
cross_org_intent_with_grant(caller_kp, org_a, provider, org_b.org_id(), service, grant)
}
fn cross_org_intent_with_grant(
caller_kp: EntityKeypair,
org_a: &OrgKeypair,
provider: EntityId,
provider_owner_org: OrgId,
service: &str,
grant: OrgCapabilityGrant,
) -> OrgProofIntent {
let caller_entity = caller_kp.entity_id().clone();
let cap = CapabilityAuthorityId::for_tag(&format!("nrpc:{service}"));
let membership =
OrgMembershipCert::try_issue(org_a, caller_entity.clone(), 1, 3600).expect("membership");
let dispatcher =
OrgDispatcherGrant::try_issue(org_a, caller_entity, DispatcherScope::Exact(cap), 3600)
.expect("dispatcher");
OrgProofIntent {
caller: Arc::new(caller_kp),
membership,
dispatcher,
capability_grant: Some(grant),
acting_org: org_a.org_id(),
provider_owner_org,
provider,
capability: cap,
proof_ttl_secs: 30,
}
}
#[tokio::test]
async fn live_two_node_cross_org_granted_admit() {
const CALLER_SEED: [u8; 32] = [0x1au8; 32];
let server = build_node_with(EntityKeypair::generate()).await;
let caller = build_node_with(EntityKeypair::from_bytes(CALLER_SEED)).await;
bring_up(&caller, &server).await;
let (org_b, _dir) = install_authority(&server, "xorg-admit");
let org_a = OrgKeypair::from_bytes([0x7au8; 32]);
assert_ne!(org_a.org_id(), org_b.org_id(), "A and B are distinct orgs");
let provider = server.entity_id().clone();
let caller_entity = caller.entity_id().clone();
let calls = Arc::new(AtomicUsize::new(0));
let saw = Arc::new(AtomicBool::new(false));
let attribution_ok = Arc::new(AtomicBool::new(false));
let stripped = Arc::new(AtomicBool::new(false));
let _serve = server
.serve_rpc_protected(
"svc",
Arc::new(AdmitHandler {
calls: calls.clone(),
saw_admission: saw.clone(),
attribution_ok: attribution_ok.clone(),
proof_stripped: stripped.clone(),
expected_caller: caller_entity,
expected_acting_org: org_a.org_id(),
expected_provider_org: org_b.org_id(),
expected_provider: provider.clone(),
expected_capability: CapabilityAuthorityId::for_tag("nrpc:svc"),
}),
OrgAdmission::CrossOrgGranted,
Arc::new(|_| true),
)
.expect("serve protected cross-org");
let intent = cross_org_invoke_intent(
EntityKeypair::from_bytes(CALLER_SEED),
&org_a,
&org_b,
provider.clone(),
"svc",
GrantTargetScope::ExactNode(provider),
);
let opts = CallOptions {
org_proof_intent: Some(intent),
deadline: Some(Instant::now() + Duration::from_secs(5)),
..Default::default()
};
let reply = caller
.call(server.node_id(), "svc", Bytes::from_static(b"ping"), opts)
.await
.expect("admitted cross-org call returns Ok");
assert_eq!(reply.body.as_ref(), b"pong", "the reply returns to S");
assert_eq!(calls.load(Ordering::SeqCst), 1, "handler ran exactly once");
assert!(
saw.load(Ordering::SeqCst),
"handler observed org_admission attribution (Some)",
);
assert!(
attribution_ok.load(Ordering::SeqCst),
"all five attribution fields (caller S, acting org A, provider org B, provider P₂, \
capability nrpc:svc) match — no caller-claimed field is used as attribution",
);
assert!(
stripped.load(Ordering::SeqCst),
"the raw net-org-admission proof header was stripped from the handler view",
);
}
#[tokio::test]
async fn live_two_node_protected_missing_local_capability_denies() {
const CALLER_SEED: [u8; 32] = [0x1du8; 32];
let server = build_node_with(EntityKeypair::generate()).await;
let caller = build_node_with(EntityKeypair::from_bytes(CALLER_SEED)).await;
bring_up(&caller, &server).await;
let (org_b, _dir) = install_authority(&server, "notag");
let org_a = OrgKeypair::from_bytes([0x7au8; 32]);
let provider = server.entity_id().clone();
let calls = Arc::new(AtomicUsize::new(0));
let _serve = server
.serve_rpc_protected(
"svc",
Arc::new(DarkHandler {
calls: calls.clone(),
}),
OrgAdmission::CrossOrgGranted,
Arc::new(|_| true),
)
.expect("serve protected");
server.test_inject_capability_announcement(CapabilityAnnouncement::new(
server.node_id(),
server.entity_id().clone(),
100,
CapabilitySet::new(),
));
let intent = cross_org_invoke_intent(
EntityKeypair::from_bytes(CALLER_SEED),
&org_a,
&org_b,
provider.clone(),
"svc",
GrantTargetScope::ExactNode(provider),
);
let err = caller
.call(
server.node_id(),
"svc",
Bytes::from_static(b"ping"),
CallOptions {
org_proof_intent: Some(intent),
deadline: Some(Instant::now() + Duration::from_secs(5)),
..Default::default()
},
)
.await
.expect_err("a valid proof to a provider missing its local tag must be denied");
match err {
RpcError::ServerError { status, .. } => {
assert_eq!(
status, 0x0009,
"AdmissionDenied (0x0009), never a timeout masquerade",
);
}
other => panic!("expected AdmissionDenied, got {other:?}"),
}
assert_handler_stays_dark(
&calls,
"the handler stayed dark — the possession precheck denied before admission",
)
.await;
}
#[tokio::test]
async fn live_cross_org_any_node_owned_by_reuse_and_deny() {
const CALLER_SEED: [u8; 32] = [0x1eu8; 32];
let org_b = OrgKeypair::from_bytes([0x42u8; 32]);
let org_a = OrgKeypair::from_bytes([0x7au8; 32]);
let org_c = OrgKeypair::from_bytes([0x33u8; 32]);
let caller = build_node_with(EntityKeypair::from_bytes(CALLER_SEED)).await;
let (p2, dir2) = adopted_node(0x51, &org_b, "anyb-p2").await;
let (p2b, dir2b) = adopted_node(0x52, &org_b, "anyb-p2b").await;
let (pc, dirc) = adopted_node(0x53, &org_c, "anyb-pc").await;
bring_up(&caller, &p2).await;
bring_up(&caller, &p2b).await;
bring_up(&caller, &pc).await;
let cap = CapabilityAuthorityId::for_tag("nrpc:svc");
let (grant, secret) = OrgCapabilityGrant::try_issue(
&org_b,
org_a.org_id(),
cap,
GrantRights::INVOKE,
GrantTargetScope::AnyNodeOwnedBy(org_b.org_id()),
3600,
)
.expect("issue AnyNodeOwnedBy grant");
assert!(secret.is_none(), "INVOKE-only mints no audience material");
let p2_calls = Arc::new(AtomicUsize::new(0));
let _s2 = p2
.serve_rpc_protected(
"svc",
Arc::new(DarkHandler {
calls: p2_calls.clone(),
}),
OrgAdmission::CrossOrgGranted,
Arc::new(|_| true),
)
.expect("serve p2");
let p2b_calls = Arc::new(AtomicUsize::new(0));
let _s2b = p2b
.serve_rpc_protected(
"svc",
Arc::new(DarkHandler {
calls: p2b_calls.clone(),
}),
OrgAdmission::CrossOrgGranted,
Arc::new(|_| true),
)
.expect("serve p2b");
let pc_calls = Arc::new(AtomicUsize::new(0));
let _sc = pc
.serve_rpc_protected(
"svc",
Arc::new(DarkHandler {
calls: pc_calls.clone(),
}),
OrgAdmission::CrossOrgGranted,
Arc::new(|_| true),
)
.expect("serve pc");
let opts = || CallOptions {
deadline: Some(Instant::now() + Duration::from_secs(5)),
..Default::default()
};
caller
.call(
p2.node_id(),
"svc",
Bytes::from_static(b"ping"),
CallOptions {
org_proof_intent: Some(cross_org_intent_with_grant(
EntityKeypair::from_bytes(CALLER_SEED),
&org_a,
p2.entity_id().clone(),
org_b.org_id(),
"svc",
grant.clone(),
)),
..opts()
},
)
.await
.expect("admit at the first B-owned node");
caller
.call(
p2b.node_id(),
"svc",
Bytes::from_static(b"ping"),
CallOptions {
org_proof_intent: Some(cross_org_intent_with_grant(
EntityKeypair::from_bytes(CALLER_SEED),
&org_a,
p2b.entity_id().clone(),
org_b.org_id(),
"svc",
grant.clone(),
)),
..opts()
},
)
.await
.expect("admit at the second B-owned node (grant reuse)");
assert_eq!(p2_calls.load(Ordering::SeqCst), 1, "P₂ handler ran once");
assert_eq!(
p2b_calls.load(Ordering::SeqCst),
1,
"the second B-owned node's handler ran once — one grant, reused",
);
let err = caller
.call(
pc.node_id(),
"svc",
Bytes::from_static(b"ping"),
CallOptions {
org_proof_intent: Some(cross_org_intent_with_grant(
EntityKeypair::from_bytes(CALLER_SEED),
&org_a,
pc.entity_id().clone(),
org_c.org_id(),
"svc",
grant.clone(),
)),
..opts()
},
)
.await
.expect_err("a B-issued grant must be denied at a non-B provider");
match err {
RpcError::ServerError { status, .. } => {
assert_eq!(status, 0x0009, "AdmissionDenied at the non-B provider");
}
other => panic!("expected AdmissionDenied, got {other:?}"),
}
assert_handler_stays_dark(&pc_calls, "the non-B provider's handler stayed dark").await;
for _d in [dir2, dir2b, dirc] {}
}
#[tokio::test]
async fn live_two_node_owner_delegated_membership_only_denied() {
const CALLER_SEED: [u8; 32] = [0x2du8; 32];
let server = build_node_with(EntityKeypair::generate()).await;
let caller = build_node_with(EntityKeypair::from_bytes(CALLER_SEED)).await;
bring_up(&caller, &server).await;
let (org_b, _dir) = install_authority(&server, "memonly");
let provider = server.entity_id().clone();
let caller_entity = caller.entity_id().clone();
let calls = Arc::new(AtomicUsize::new(0));
let _serve = server
.serve_rpc_protected(
"svc",
Arc::new(DarkHandler {
calls: calls.clone(),
}),
OrgAdmission::OwnerDelegated,
Arc::new(|_| true),
)
.expect("serve protected");
let membership =
OrgMembershipCert::try_issue(&org_b, caller_entity.clone(), 1, 3600).expect("membership");
let dispatcher = OrgDispatcherGrant::try_issue(
&org_b,
caller_entity,
DispatcherScope::Exact(CapabilityAuthorityId::for_tag("nrpc:elsewhere")),
3600,
)
.expect("dispatcher");
let intent = OrgProofIntent {
caller: Arc::new(EntityKeypair::from_bytes(CALLER_SEED)),
membership,
dispatcher,
capability_grant: None,
acting_org: org_b.org_id(),
provider_owner_org: org_b.org_id(),
provider,
capability: CapabilityAuthorityId::for_tag("nrpc:svc"),
proof_ttl_secs: 30,
};
let err = caller
.call(
server.node_id(),
"svc",
Bytes::from_static(b"ping"),
CallOptions {
org_proof_intent: Some(intent),
deadline: Some(Instant::now() + Duration::from_secs(5)),
..Default::default()
},
)
.await
.expect_err("a membership-only proof must be denied");
match err {
RpcError::ServerError { status, .. } => assert_eq!(status, 0x0009, "AdmissionDenied"),
other => panic!("expected AdmissionDenied, got {other:?}"),
}
assert_handler_stays_dark(&calls, "the handler stayed dark").await;
}
#[tokio::test]
async fn live_two_node_public_capability_unchanged_beside_protected() {
const CALLER_SEED: [u8; 32] = [0x2eu8; 32];
let server = build_node_with(EntityKeypair::generate()).await;
let caller = build_node_with(EntityKeypair::from_bytes(CALLER_SEED)).await;
bring_up(&caller, &server).await;
let (_org_b, _dir) = install_authority(&server, "pubunchanged");
let pub_calls = Arc::new(AtomicUsize::new(0));
let saw_proof = Arc::new(AtomicBool::new(false));
let _pub = server
.serve_rpc(
"pub",
Arc::new(HeaderSpyHandler {
calls: pub_calls.clone(),
saw_proof: saw_proof.clone(),
}),
)
.expect("serve public");
let prot_calls = Arc::new(AtomicUsize::new(0));
let _prot = server
.serve_rpc_protected(
"svc",
Arc::new(DarkHandler {
calls: prot_calls.clone(),
}),
OrgAdmission::OwnerDelegated,
Arc::new(|_| true),
)
.expect("serve protected");
let reply = caller
.call(
server.node_id(),
"pub",
Bytes::from_static(b"ping"),
CallOptions {
deadline: Some(Instant::now() + Duration::from_secs(5)),
..Default::default()
},
)
.await
.expect("the public call succeeds without a proof");
assert_eq!(reply.body.as_ref(), b"pong");
assert_eq!(
pub_calls.load(Ordering::SeqCst),
1,
"public handler ran once"
);
assert!(
!saw_proof.load(Ordering::SeqCst),
"the public handler saw no org-admission proof header",
);
let err = caller
.call(
server.node_id(),
"svc",
Bytes::from_static(b"ping"),
CallOptions {
deadline: Some(Instant::now() + Duration::from_secs(5)),
..Default::default()
},
)
.await
.expect_err("the protected service still requires a proof");
match err {
RpcError::ServerError { status, .. } => assert_eq!(status, 0x0009, "AdmissionDenied"),
other => panic!("expected AdmissionDenied, got {other:?}"),
}
assert_handler_stays_dark(&prot_calls, "the protected handler stayed dark").await;
}
#[tokio::test]
async fn live_two_node_owner_delegated_floor_survives_restart_denies() {
const CALLER_SEED: [u8; 32] = [0x2fu8; 32];
let server = build_node_with(EntityKeypair::generate()).await;
let caller = build_node_with(EntityKeypair::from_bytes(CALLER_SEED)).await;
bring_up(&caller, &server).await;
let node_entity = server.entity_id().clone();
let caller_entity = caller.entity_id().clone();
let org_b = OrgKeypair::from_bytes([0x42u8; 32]);
let node_cert =
OrgMembershipCert::try_issue(&org_b, node_entity.clone(), 1, 3600).expect("node cert");
let dir = std::env::temp_dir().join(format!("net-oa4-restart-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let authority =
NodeAuthority::adopt(&dir, node_cert, &node_entity, 0, None).expect("adopt authority");
let mut floors = std::collections::BTreeMap::new();
floors.insert(caller_entity.clone(), 5u32);
authority
.revocation
.apply_bundle(&OrgRevocationBundle::try_issue(&org_b, &floors).expect("bundle 5"))
.expect("apply floor 5");
server
.install_node_authority(Arc::new(authority))
.expect("install authority");
let calls = Arc::new(AtomicUsize::new(0));
let _serve = server
.serve_rpc_protected(
"svc",
Arc::new(DarkHandler {
calls: calls.clone(),
}),
OrgAdmission::OwnerDelegated,
Arc::new(|_| true),
)
.expect("serve protected");
let reopened = NodeAuthority::open(&dir, &node_entity).expect("reopen authority");
assert_eq!(
reopened
.revocation
.floor_for(&org_b.org_id(), &caller_entity),
5,
"floor 5 survives the restart",
);
let mut lower = std::collections::BTreeMap::new();
lower.insert(caller_entity.clone(), 3u32);
reopened
.revocation
.apply_bundle(&OrgRevocationBundle::try_issue(&org_b, &lower).expect("bundle 3"))
.expect("lower bundle is a no-op");
assert_eq!(
reopened
.revocation
.floor_for(&org_b.org_id(), &caller_entity),
5,
"floor stays 5 after a lower bundle",
);
server
.install_node_authority(Arc::new(reopened))
.expect("install reopened authority");
let intent4 = owner_delegated_intent_gen(
EntityKeypair::from_bytes(CALLER_SEED),
&org_b,
node_entity.clone(),
"svc",
4,
);
let err = caller
.call(
server.node_id(),
"svc",
Bytes::from_static(b"ping"),
CallOptions {
org_proof_intent: Some(intent4),
deadline: Some(Instant::now() + Duration::from_secs(5)),
..Default::default()
},
)
.await
.expect_err("a below-floor cert must be denied after the restart");
match err {
RpcError::ServerError { status, .. } => {
assert_eq!(status, 0x0009, "MembershipRevoked at the live gate");
}
other => panic!("expected AdmissionDenied, got {other:?}"),
}
assert_handler_stays_dark(&calls, "the handler stayed dark for the below-floor call").await;
let intent5 = owner_delegated_intent_gen(
EntityKeypair::from_bytes(CALLER_SEED),
&org_b,
node_entity.clone(),
"svc",
5,
);
caller
.call(
server.node_id(),
"svc",
Bytes::from_static(b"ping"),
CallOptions {
org_proof_intent: Some(intent5),
deadline: Some(Instant::now() + Duration::from_secs(5)),
..Default::default()
},
)
.await
.expect("an at-floor cert admits");
assert_eq!(
calls.load(Ordering::SeqCst),
1,
"the at-floor call reached the handler",
);
}
async fn granted_service_provider<H: RpcHandler>(
seed: u8,
svc: &str,
handler: Arc<H>,
) -> (Arc<MeshNode>, ServeHandle, EntityId, std::path::PathBuf) {
let p = build_node_fast_announce(EntityKeypair::from_bytes([seed; 32])).await;
let entity = p.entity_id().clone();
let org_b = OrgKeypair::from_bytes([0x42u8; 32]);
let cert = OrgMembershipCert::try_issue(&org_b, entity.clone(), 1, 3600).expect("cert");
let dir = std::env::temp_dir().join(format!(
"net-oa4-granted-{svc}-{seed}-{}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&dir);
let authority = NodeAuthority::adopt(&dir, cert, &entity, 0, None).expect("adopt");
p.install_node_authority(Arc::new(authority))
.expect("install authority");
p.set_owner_cert_emission(true).expect("enable emission");
let handle = p
.serve_rpc_granted(svc, handler, Arc::new(|_| true))
.expect("granted serve");
(p, handle, entity, dir)
}
async fn converge_granted_resolution(
provider: &Arc<MeshNode>,
consumer: &Arc<MeshNode>,
grant_id: &[u8; 32],
provider_entity: &EntityId,
now: u64,
) -> bool {
for _ in 0..100 {
provider
.announce_capabilities(CapabilitySet::new())
.await
.ok();
if consumer.scoped_granted_providers_for_test(grant_id, now)
== vec![provider_entity.clone()]
{
return true;
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
false
}
#[tokio::test]
async fn live_granted_audience_discovers_then_invokes() {
let org_b = OrgKeypair::from_bytes([0x42u8; 32]);
let org_a = OrgKeypair::from_bytes([0x7au8; 32]);
let now = unix_now();
let cap = CapabilityAuthorityId::for_tag("nrpc:svc");
let calls = Arc::new(AtomicUsize::new(0));
let (p2, _serve, p_entity, _p_dir) = granted_service_provider(
0x61,
"svc",
Arc::new(DarkHandler {
calls: calls.clone(),
}),
)
.await;
let (a, _a_dir) = adopted_node(0x62, &org_a, "gdisc").await;
bring_up(&a, &p2).await;
let (grant, secret) = OrgCapabilityGrant::try_issue(
&org_b,
org_a.org_id(),
cap,
GrantRights::DISCOVER.union(GrantRights::INVOKE),
GrantTargetScope::ExactNode(p_entity.clone()),
3600,
)
.expect("grant");
let secret = secret.expect("a DISCOVER grant mints an audience secret");
let grant_id = grant.grant_id;
a.install_consumer_grant_audience(grant.clone(), copy_secret(&secret))
.expect("install consumer grant");
p2.install_provider_grant_audience(grant.clone(), copy_secret(&secret))
.expect("install provider grant");
assert!(
converge_granted_resolution(&p2, &a, &grant_id, &p_entity, now).await,
"A privately resolves exactly P₂ over the live 0x0C04 scoped send",
);
let intent = cross_org_intent_with_grant(
EntityKeypair::from_bytes([0x62u8; 32]),
&org_a,
p_entity.clone(),
org_b.org_id(),
"svc",
grant.clone(),
);
a.call(
p2.node_id(),
"svc",
Bytes::from_static(b"ping"),
CallOptions {
org_proof_intent: Some(intent),
deadline: Some(Instant::now() + Duration::from_secs(5)),
..Default::default()
},
)
.await
.expect("the granted invocation admits");
assert_eq!(
calls.load(Ordering::SeqCst),
1,
"the exact P₂ handler ran once under the same grant",
);
}
#[tokio::test]
async fn live_granted_audience_discover_only_resolves_but_cannot_invoke() {
let org_b = OrgKeypair::from_bytes([0x42u8; 32]);
let org_a = OrgKeypair::from_bytes([0x7au8; 32]);
let now = unix_now();
let calls = Arc::new(AtomicUsize::new(0));
let (p2, _serve, p_entity, _p_dir) = granted_service_provider(
0x63,
"svc",
Arc::new(DarkHandler {
calls: calls.clone(),
}),
)
.await;
let (a, _a_dir) = adopted_node(0x64, &org_a, "gdonly").await;
bring_up(&a, &p2).await;
let (grant, secret) = cross_org_grant(&org_b, &org_a, &p_entity, "svc");
let grant_id = grant.grant_id;
a.install_consumer_grant_audience(grant.clone(), copy_secret(&secret))
.expect("install consumer grant");
p2.install_provider_grant_audience(grant.clone(), copy_secret(&secret))
.expect("install provider grant");
assert!(
converge_granted_resolution(&p2, &a, &grant_id, &p_entity, now).await,
"A resolves P₂ under the DISCOVER right over the live 0x0C04 scoped send",
);
let intent = cross_org_intent_with_grant(
EntityKeypair::from_bytes([0x64u8; 32]),
&org_a,
p_entity.clone(),
org_b.org_id(),
"svc",
grant.clone(),
);
let err = a
.call(
p2.node_id(),
"svc",
Bytes::from_static(b"ping"),
CallOptions {
org_proof_intent: Some(intent),
deadline: Some(Instant::now() + Duration::from_secs(5)),
..Default::default()
},
)
.await
.expect_err("a DISCOVER-only grant cannot invoke");
match err {
RpcError::ServerError { status, .. } => assert_eq!(status, 0x0009, "AdmissionDenied"),
other => panic!("expected AdmissionDenied, got {other:?}"),
}
assert_handler_stays_dark(
&calls,
"the handler stayed dark — discovery did not confer invocation",
)
.await;
}
#[tokio::test]
async fn live_granted_audience_wrong_dispatcher_resolves_but_invocation_denied() {
let org_b = OrgKeypair::from_bytes([0x42u8; 32]);
let org_a = OrgKeypair::from_bytes([0x7au8; 32]);
let now = unix_now();
let cap = CapabilityAuthorityId::for_tag("nrpc:svc");
let calls = Arc::new(AtomicUsize::new(0));
let (p2, _serve, p_entity, _p_dir) = granted_service_provider(
0x65,
"svc",
Arc::new(DarkHandler {
calls: calls.clone(),
}),
)
.await;
let (a, _a_dir) = adopted_node(0x66, &org_a, "gwrongdisp").await;
bring_up(&a, &p2).await;
let (grant, secret) = OrgCapabilityGrant::try_issue(
&org_b,
org_a.org_id(),
cap,
GrantRights::DISCOVER.union(GrantRights::INVOKE),
GrantTargetScope::ExactNode(p_entity.clone()),
3600,
)
.expect("grant");
let secret = secret.expect("secret");
let grant_id = grant.grant_id;
a.install_consumer_grant_audience(grant.clone(), copy_secret(&secret))
.expect("install consumer grant");
p2.install_provider_grant_audience(grant.clone(), copy_secret(&secret))
.expect("install provider grant");
assert!(
converge_granted_resolution(&p2, &a, &grant_id, &p_entity, now).await,
"A resolves P₂ over the live 0x0C04 scoped send",
);
let a_entity = a.entity_id().clone();
let intent = OrgProofIntent {
caller: Arc::new(EntityKeypair::from_bytes([0x66u8; 32])),
membership: OrgMembershipCert::try_issue(&org_a, a_entity.clone(), 1, 3600)
.expect("membership"),
dispatcher: OrgDispatcherGrant::try_issue(
&org_a,
a_entity,
DispatcherScope::Exact(CapabilityAuthorityId::for_tag("nrpc:elsewhere")),
3600,
)
.expect("dispatcher"),
capability_grant: Some(grant.clone()),
acting_org: org_a.org_id(),
provider_owner_org: org_b.org_id(),
provider: p_entity.clone(),
capability: cap,
proof_ttl_secs: 30,
};
let err = a
.call(
p2.node_id(),
"svc",
Bytes::from_static(b"ping"),
CallOptions {
org_proof_intent: Some(intent),
deadline: Some(Instant::now() + Duration::from_secs(5)),
..Default::default()
},
)
.await
.expect_err("a dispatcher scoped elsewhere cannot invoke");
match err {
RpcError::ServerError { status, .. } => assert_eq!(status, 0x0009, "AdmissionDenied"),
other => panic!("expected AdmissionDenied, got {other:?}"),
}
assert_handler_stays_dark(&calls, "the handler stayed dark").await;
}
#[tokio::test]
async fn live_granted_audience_provider_policy_final() {
let org_b = OrgKeypair::from_bytes([0x42u8; 32]);
let org_a = OrgKeypair::from_bytes([0x7au8; 32]);
let now = unix_now();
let cap = CapabilityAuthorityId::for_tag("nrpc:svc");
let p2 = build_node_fast_announce(EntityKeypair::from_bytes([0x67u8; 32])).await;
let p_entity = p2.entity_id().clone();
let node_cert =
OrgMembershipCert::try_issue(&org_b, p_entity.clone(), 1, 3600).expect("node cert");
let p_dir = std::env::temp_dir().join(format!("net-oa4-granted-veto-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&p_dir);
let authority =
NodeAuthority::adopt(&p_dir, node_cert, &p_entity, 0, None).expect("adopt authority");
p2.install_node_authority(Arc::new(authority))
.expect("install authority");
p2.set_owner_cert_emission(true).expect("enable emission");
let calls = Arc::new(AtomicUsize::new(0));
let _serve = p2
.serve_rpc_granted(
"svc",
Arc::new(DarkHandler {
calls: calls.clone(),
}),
Arc::new(|_| false),
)
.expect("granted serve with veto policy");
let (a, _a_dir) = adopted_node(0x68, &org_a, "gveto").await;
bring_up(&a, &p2).await;
let (grant, secret) = OrgCapabilityGrant::try_issue(
&org_b,
org_a.org_id(),
cap,
GrantRights::DISCOVER.union(GrantRights::INVOKE),
GrantTargetScope::ExactNode(p_entity.clone()),
3600,
)
.expect("grant");
let secret = secret.expect("secret");
let grant_id = grant.grant_id;
a.install_consumer_grant_audience(grant.clone(), copy_secret(&secret))
.expect("install consumer grant");
p2.install_provider_grant_audience(grant.clone(), copy_secret(&secret))
.expect("install provider grant");
assert!(
converge_granted_resolution(&p2, &a, &grant_id, &p_entity, now).await,
"A resolves P₂ over the live 0x0C04 scoped send",
);
let intent = cross_org_intent_with_grant(
EntityKeypair::from_bytes([0x68u8; 32]),
&org_a,
p_entity.clone(),
org_b.org_id(),
"svc",
grant.clone(),
);
let err = a
.call(
p2.node_id(),
"svc",
Bytes::from_static(b"ping"),
CallOptions {
org_proof_intent: Some(intent),
deadline: Some(Instant::now() + Duration::from_secs(5)),
..Default::default()
},
)
.await
.expect_err("the provider policy vetoes the call");
match err {
RpcError::ServerError { status, .. } => assert_eq!(status, 0x0009, "AdmissionDenied"),
other => panic!("expected AdmissionDenied, got {other:?}"),
}
assert_handler_stays_dark(&calls, "the handler stayed dark under the policy veto").await;
}
#[test]
fn invoke_only_grant_carries_no_discovery_material() {
let org_b = OrgKeypair::from_bytes([0x42u8; 32]);
let org_a = OrgKeypair::from_bytes([0x7au8; 32]);
let provider = EntityKeypair::from_bytes([0x69u8; 32]).entity_id().clone();
let (grant, secret) = OrgCapabilityGrant::try_issue(
&org_b,
org_a.org_id(),
CapabilityAuthorityId::for_tag("nrpc:svc"),
GrantRights::INVOKE,
GrantTargetScope::ExactNode(provider),
3600,
)
.expect("issue INVOKE-only grant");
assert!(
secret.is_none(),
"an INVOKE-only grant mints no audience secret",
);
assert!(
!grant.permits_discover(),
"an INVOKE-only grant confers no discovery right",
);
assert!(grant.permits_invoke(), "it does carry INVOKE");
}
#[tokio::test]
async fn live_two_node_proof_for_another_identity_is_denied() {
const SESSION_SEED: [u8; 32] = [0x71u8; 32];
const OTHER_SEED: [u8; 32] = [0x72u8; 32];
let server = build_node_with(EntityKeypair::generate()).await;
let caller = build_node_with(EntityKeypair::from_bytes(SESSION_SEED)).await;
bring_up(&caller, &server).await;
let (org_b, _dir) = install_authority(&server, "bindmismatch");
let provider = server.entity_id().clone();
let calls = Arc::new(AtomicUsize::new(0));
let _serve = server
.serve_rpc_protected(
"svc",
Arc::new(DarkHandler {
calls: calls.clone(),
}),
OrgAdmission::OwnerDelegated,
Arc::new(|_| true),
)
.expect("serve protected");
let other_kp = EntityKeypair::from_bytes(OTHER_SEED);
assert_ne!(
other_kp.entity_id(),
caller.entity_id(),
"the proof subject must differ from the session peer",
);
let intent = owner_delegated_intent(other_kp, &org_b, provider.clone(), "svc");
let opts = CallOptions {
org_proof_intent: Some(intent),
deadline: Some(Instant::now() + Duration::from_secs(5)),
..Default::default()
};
let err = caller
.call(server.node_id(), "svc", Bytes::from_static(b"ping"), opts)
.await
.expect_err("a proof minted for another identity must be denied");
match err {
RpcError::ServerError { status, .. } => assert_eq!(
status, 0x0009,
"a proof/session identity mismatch denies AdmissionDenied",
),
other => panic!("expected AdmissionDenied, got {other:?}"),
}
assert_handler_stays_dark(
&calls,
"the handler must stay dark for a proof bound to another identity",
)
.await;
let self_intent = owner_delegated_intent(
EntityKeypair::from_bytes(SESSION_SEED),
&org_b,
provider,
"svc",
);
let opts = CallOptions {
org_proof_intent: Some(self_intent),
deadline: Some(Instant::now() + Duration::from_secs(5)),
..Default::default()
};
caller
.call(server.node_id(), "svc", Bytes::from_static(b"ping"), opts)
.await
.expect("a correctly-bound proof on the same session admits");
assert_eq!(
calls.load(Ordering::SeqCst),
1,
"only the correctly-bound call reached the handler",
);
}
#[tokio::test]
async fn a_proof_ttl_outside_the_ceiling_fails_locally() {
const CALLER_SEED: [u8; 32] = [0x73u8; 32];
let server = build_node_with(EntityKeypair::generate()).await;
let caller = build_node_with(EntityKeypair::from_bytes(CALLER_SEED)).await;
bring_up(&caller, &server).await;
let (org_b, _dir) = install_authority(&server, "ttlceiling");
let provider = server.entity_id().clone();
let calls = Arc::new(AtomicUsize::new(0));
let _serve = server
.serve_rpc_protected(
"svc",
Arc::new(DarkHandler {
calls: calls.clone(),
}),
OrgAdmission::OwnerDelegated,
Arc::new(|_| true),
)
.expect("serve protected");
for bad_ttl in [0u64, MAX_ORG_PROOF_TTL_SECS + 1] {
let mut intent = owner_delegated_intent(
EntityKeypair::from_bytes(CALLER_SEED),
&org_b,
provider.clone(),
"svc",
);
intent.proof_ttl_secs = bad_ttl;
let opts = CallOptions {
org_proof_intent: Some(intent),
deadline: Some(Instant::now() + Duration::from_secs(5)),
..Default::default()
};
let err = caller
.call(server.node_id(), "svc", Bytes::from_static(b"ping"), opts)
.await
.expect_err("a TTL outside 1..=MAX must be refused");
assert!(
matches!(
err,
RpcError::Codec {
direction: CodecDirection::Encode,
..
}
),
"ttl {bad_ttl} must be refused by the CALLER at encode time, \
before any frame is emitted; got: {err:?}",
);
}
assert_handler_stays_dark(
&calls,
"no frame reached the provider for an out-of-range TTL",
)
.await;
let mut intent = owner_delegated_intent(
EntityKeypair::from_bytes(CALLER_SEED),
&org_b,
provider,
"svc",
);
intent.proof_ttl_secs = MAX_ORG_PROOF_TTL_SECS;
let opts = CallOptions {
org_proof_intent: Some(intent),
deadline: Some(Instant::now() + Duration::from_secs(5)),
..Default::default()
};
caller
.call(server.node_id(), "svc", Bytes::from_static(b"ping"), opts)
.await
.expect("a TTL at the ceiling admits");
assert_eq!(calls.load(Ordering::SeqCst), 1, "only the valid TTL called");
}