From 4fc6854846ca79a22968bfd2e3406a31468b04b1 Mon Sep 17 00:00:00 2001 From: deepseek-v4-pro Date: Fri, 17 Jul 2026 14:01:51 +0000 Subject: [PATCH] chore(call): prune connect, TLS helpers, CallCredentials, and dead paths (Phase 5) - Remove RemoteIdentity, CallCredentials, ClientError, connect() from call_client.rs - Remove all TLS helpers (build_quinn_client_config, build_client_auth, select_server_verifier, FingerprintPinVerifier, RawKeyClientCertResolver, NoClientCertResolver, Ed25519SigningKey, cert/key loaders) - Remove credentials_auth_token dead path from from_call.rs (OpSummary field, build_bundles, make_forwarding_handler, make_streaming_forwarding_handler, build_forwarded_payload) - Remove 2 dead-path tests (build_forwarded_payload_sets_auth_token, streaming_forwarding_handler_sets_auth_token) - Update mod.rs re-exports to only CallClient - Remove quinn feature and TLS deps from Cargo.toml - Delete tests/two_node_call.rs (used connect() + CallCredentials) - alknet-call is now a pure protocol crate --- Cargo.lock | 6 - crates/alknet-call/Cargo.toml | 13 +- crates/alknet-call/src/client/call_client.rs | 743 +------------------ crates/alknet-call/src/client/from_call.rs | 107 +-- crates/alknet-call/src/client/mod.rs | 2 +- crates/alknet-call/src/lib.rs | 2 +- crates/alknet-call/tests/two_node_call.rs | 331 --------- 7 files changed, 32 insertions(+), 1172 deletions(-) delete mode 100644 crates/alknet-call/tests/two_node_call.rs diff --git a/Cargo.lock b/Cargo.lock index 4ba6236..d5c42e1 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -53,13 +53,7 @@ dependencies = [ "alknet-core", "async-trait", "futures", - "hex", "parking_lot", - "quinn", - "rcgen", - "rustls", - "rustls-native-certs", - "rustls-pemfile", "serde", "serde_json", "thiserror 2.0.18", diff --git a/crates/alknet-call/Cargo.toml b/crates/alknet-call/Cargo.toml index 51eba17..f3138a5 100644 --- a/crates/alknet-call/Cargo.toml +++ b/crates/alknet-call/Cargo.toml @@ -3,15 +3,14 @@ name = "alknet-call" version.workspace = true edition.workspace = true license.workspace = true -description = "Structured RPC over QUIC on ALPN `alknet/call`: operations, streaming subscriptions, service discovery" +description = "Structured RPC over ALPN `alknet/call`: operations, streaming subscriptions, service discovery" repository.workspace = true [lib] name = "alknet_call" [features] -default = ["quinn"] -quinn = ["dep:quinn", "dep:rustls", "dep:rustls-native-certs", "dep:rustls-pemfile", "alknet-core/quinn"] +default = [] [dependencies] alknet-core = { path = "../alknet-core" } @@ -24,11 +23,3 @@ thiserror = "2" uuid = { version = "1", features = ["v4"] } futures = "0.3" parking_lot = "0.12" -quinn = { version = "0.11", optional = true } -rustls = { version = "0.23", optional = true, features = ["aws_lc_rs"] } -rustls-native-certs = { version = "0.8", optional = true } -rustls-pemfile = { version = "2", optional = true } - -[dev-dependencies] -rcgen = "0.13" -hex = "0.4" \ No newline at end of file diff --git a/crates/alknet-call/src/client/call_client.rs b/crates/alknet-call/src/client/call_client.rs index 5943101..97145f1 100644 --- a/crates/alknet-call/src/client/call_client.rs +++ b/crates/alknet-call/src/client/call_client.rs @@ -1,8 +1,7 @@ //! `CallClient`: the outbound connection opener (ADR-017 §1). //! -//! Opens a QUIC connection to a remote node on ALPN `alknet/call`, performs -//! credential setup, and produces a [`CallConnection`] running the shared -//! dispatch loop (delegated to [`crate::protocol::dispatch::Dispatcher`]). +//! Runs the shared dispatch loop over a pre-established `Connection` +//! (delegated to [`crate::protocol::dispatch::Dispatcher`]). //! `CallClient` is the connection-establishment half; `CallAdapter`'s accept //! path is the inbound half; both produce a `CallConnection` and hand it to //! the same `Dispatcher::run_loop` (ADR-017 §1). @@ -12,93 +11,21 @@ //! (initiates outgoing calls via `CallConnection::call()`/`subscribe()`/ //! `abort()`) and a callee (dispatches incoming calls against its registry). //! +//! Transport-level connection establishment (QUIC dial, TCP+TLS, iroh) is +//! handled by `alknet-client`; `CallClient::spawn_dispatch` takes a +//! pre-established `Connection` and runs the call protocol over it. +//! //! See `docs/architecture/crates/call/client-and-adapters.md` for the spec. -use std::net::SocketAddr; use std::sync::Arc; use alknet_core::auth::IdentityProvider; -use alknet_core::config::TlsIdentity; use alknet_core::types::Connection; use crate::protocol::connection::CallConnection; use crate::protocol::dispatch::Dispatcher; use crate::registry::registration::OperationRegistry; -/// Expected identity of the remote node (ADR-017 §7, extended by ADR-034 §2). -/// Carries a fingerprint string the assembly layer derives from `Capabilities` -/// when the local node has a `PeerEntry` for the remote (the known-peer case → -/// fingerprint pin). -/// -/// `remote_identity: None` is the **public X.509 endpoint** case: the local -/// node has no `PeerEntry` for the remote, so there is no fingerprint to pin. -/// Combined with an X.509 transport, `None` selects CA verification -/// (`WebPkiServerVerifier`) per the verifier-selection rule in ADR-034 §3. -/// Combined with an Ed25519 raw-key transport, `None` fails closed (raw-key -/// remotes are always known peers — no CA to fall back to). -/// -/// The `Option` is therefore load-bearing, not cosmetic: `Some(fingerprint)` -/// means "pin this" (known peer), `None` means "trust the CA or fail" -/// (unknown remote). An implementer must not default `remote_identity` to a -/// placeholder value to "satisfy" the field — `None` is a real state that -/// drives verifier selection. -#[derive(Debug, Clone)] -pub struct RemoteIdentity { - pub fingerprint: String, -} - -/// Credentials for an outbound `alknet/call` connection (ADR-017 §7). All -/// three dimensions come from `Capabilities` (ADR-014), never from environment -/// variables — see the No-Env-Vars Invariant in -/// `docs/architecture/crates/call/client-and-adapters.md`. -#[derive(Debug, Clone, Default)] -pub struct CallCredentials { - /// The local node's TLS identity (RFC 7250 raw key or X.509), derived - /// from the vault at startup. - pub tls_identity: Option, - /// Opaque call-protocol-level auth token, decrypted from the vault. - pub auth_token: Option, - /// Expected fingerprint/cert of the remote node, stored as a capability. - /// `Some` → fingerprint pin (known peer with a `PeerEntry`); `None` → CA - /// verification for X.509 remotes, fail-closed for Ed25519 raw-key remotes - /// (ADR-034 §2/§3). `None` is the public-X.509-endpoint state, not a - /// missing field — must not be defaulted to a placeholder. - pub remote_identity: Option, -} - -impl CallCredentials { - pub fn new() -> Self { - Self::default() - } - - pub fn with_tls_identity(mut self, tls_identity: TlsIdentity) -> Self { - self.tls_identity = Some(tls_identity); - self - } - - pub fn with_auth_token(mut self, token: alknet_core::auth::AuthToken) -> Self { - self.auth_token = Some(token); - self - } - - pub fn with_remote_identity(mut self, remote: RemoteIdentity) -> Self { - self.remote_identity = Some(remote); - self - } -} - -/// Errors produced by [`CallClient::connect`]. -#[derive(Debug, thiserror::Error)] -#[non_exhaustive] -pub enum ClientError { - #[error("transport error: {message}")] - Transport { message: String }, - #[error("tls setup error: {message}")] - TlsSetup { message: String }, - #[error("connection closed")] - ConnectionClosed, -} - /// Outbound `alknet/call` connection opener (the #1 gap, ADR-017 §1). /// /// Peer authorization flows through the existing `AccessControl::check` gate @@ -128,50 +55,11 @@ impl CallClient { &self.identity_provider } - /// Open a QUIC connection to `addr` on ALPN `alknet/call`, perform - /// credential handshake, and return a `CallConnection` running the shared - /// dispatch loop. Credentials come from `Capabilities` (ADR-014), not env - /// vars — the no-env-vars invariant. - /// - /// The dispatch loop runs on a spawned task; the returned `CallConnection` - /// is live until the remote closes the connection or the caller drops it. - /// The caller can immediately use `call()`/`subscribe()`/`abort()` on the - /// returned connection, and the remote peer can call back into this - /// `CallClient`'s registry (connection symmetry, ADR-017 §2). - #[cfg(feature = "quinn")] - pub async fn connect( - &self, - addr: SocketAddr, - credentials: CallCredentials, - ) -> Result { - let alpn = b"alknet/call".to_vec(); - let client_config = build_quinn_client_config(&credentials, &alpn) - .map_err(|e| ClientError::TlsSetup { message: e })?; - - let bind_addr: SocketAddr = "0.0.0.0:0".parse().expect("valid bind addr"); - let endpoint = quinn::Endpoint::client(bind_addr).map_err(|e| ClientError::Transport { - message: e.to_string(), - })?; - - let connection = endpoint - .connect_with(client_config, addr, "alknet") - .map_err(|e| ClientError::Transport { - message: e.to_string(), - })? - .await - .map_err(|e| ClientError::Transport { - message: e.to_string(), - })?; - - let connection = Connection::from_quinn_with_alpn(connection, alpn); - Ok(self.spawn_dispatch(connection)) - } - /// Run the shared dispatch loop over a pre-established `Connection`. The /// `CallClient` spawns the dispatcher task and returns a live - /// `CallConnection` the caller can use immediately. Used by `connect()` - /// (after the QUIC dial completes) and by integration tests that wire a - /// mock/loopback `Connection` directly. + /// `CallConnection` the caller can use immediately. Used by the assembly + /// layer after `AlknetClient::dial_*` + `spawn_dispatch` and by + /// integration tests that wire a mock/loopback `Connection` directly. pub fn spawn_dispatch(&self, connection: Connection) -> CallConnection { let call_connection = Arc::new(CallConnection::new(connection)); let dispatcher = Dispatcher::new( @@ -186,386 +74,6 @@ impl CallClient { } } -#[cfg(feature = "quinn")] -fn build_quinn_client_config( - credentials: &CallCredentials, - alpn: &[u8], -) -> Result { - let provider = Arc::new(rustls::crypto::aws_lc_rs::default_provider()); - - let client_auth = build_client_auth(&provider, &credentials.tls_identity)?; - let verifier = select_server_verifier(&provider, &credentials.remote_identity)?; - - let mut config = rustls::ClientConfig::builder_with_provider(provider) - .with_safe_default_protocol_versions() - .map_err(|e| e.to_string())? - .dangerous() - .with_custom_certificate_verifier(verifier) - .with_client_cert_resolver(client_auth); - config.alpn_protocols = vec![alpn.to_vec()]; - config.enable_early_data = true; - - Ok(quinn::ClientConfig::new(Arc::new( - quinn::crypto::rustls::QuicClientConfig::try_from(config).map_err(|e| e.to_string())?, - ))) -} - -/// Build the client-auth cert resolver that presents the local node's TLS -/// identity. For `TlsIdentity::RawKey` the Ed25519 key is presented as an RFC -/// 7250 raw public key client cert (`only_raw_public_keys() == true`) — the -/// client-side equivalent of the server's `RawKeyCertResolver`. For X.509 the -/// cert chain + key are loaded from disk. `None` (no `tls_identity` configured) -/// resolves to no client cert (the server gets nothing to fingerprint). -#[cfg(feature = "quinn")] -fn build_client_auth( - provider: &Arc, - tls_identity: &Option, -) -> Result, String> { - match tls_identity { - Some(TlsIdentity::RawKey(secret_key)) => { - let signing_key = Arc::new(Ed25519SigningKey::new(secret_key.clone())); - let spki = signing_key.spki_public_key(); - let cert = rustls::pki_types::CertificateDer::from(spki.to_vec()); - let certified_key = Arc::new(rustls::sign::CertifiedKey::new(vec![cert], signing_key)); - Ok(Arc::new(RawKeyClientCertResolver::new(certified_key))) - } - Some(TlsIdentity::X509 { cert, key }) => { - let cert_chain = load_cert_chain(cert).map_err(|e| e.to_string())?; - let key_der = load_private_key(key).map_err(|e| e.to_string())?; - let certified_key = rustls::sign::CertifiedKey::from_der(cert_chain, key_der, provider) - .map_err(|e| e.to_string())?; - Ok(Arc::new(RawKeyClientCertResolver::new(Arc::new( - certified_key, - )))) - } - Some(TlsIdentity::SelfSigned) | None => Ok(Arc::new(NoClientCertResolver)), - Some(TlsIdentity::Acme { .. }) => { - Err("ACME TLS identity is server-only; cannot be used for client auth".to_string()) - } - } -} - -/// Select the server cert verifier by `remote_identity` presence (ADR-034 §3). -/// -/// - `Some(fingerprint)` → known peer → `FingerprintPinVerifier` (fingerprint -/// match). The fingerprint IS the trust anchor. -/// - `None` → no `PeerEntry` for the remote → `WebPkiServerVerifier` (CA -/// verification) for X.509 remotes. For Ed25519 raw-key remotes the -/// `WebPkiServerVerifier` fails closed at handshake time (raw-key remotes -/// have no CA to fall back to — ADR-034 §2 assumption 1). `None` is the -/// public-X.509-endpoint state, not "skip verification." -#[cfg(feature = "quinn")] -fn select_server_verifier( - provider: &Arc, - remote_identity: &Option, -) -> Result, String> { - match remote_identity { - Some(ri) => Ok(Arc::new(FingerprintPinVerifier::new( - ri.fingerprint.clone(), - provider.signature_verification_algorithms, - ))), - None => { - let roots = load_platform_root_cert_store()?; - let verifier = rustls::client::WebPkiServerVerifier::builder_with_provider( - Arc::new(roots), - Arc::clone(provider), - ) - .build() - .map_err(|e| e.to_string())?; - Ok(verifier) - } - } -} - -/// Load the platform's trusted root certificates into a `RootCertStore` for -/// `WebPkiServerVerifier` (the `None` + X.509 CA-verification path). Falls back -/// to the aws-lc-rs built-in `webpki-roots` if the platform store is empty -/// (e.g. in a container with no system CA bundle). -#[cfg(feature = "quinn")] -fn load_platform_root_cert_store() -> Result { - let mut roots = rustls::RootCertStore::empty(); - let result = rustls_native_certs::load_native_certs(); - for err in &result.errors { - tracing::warn!(error = ?err, "failed to load a native root cert"); - } - for cert in &result.certs { - roots - .add(cert.clone()) - .map_err(|e| format!("failed to add native root cert: {e}"))?; - } - Ok(roots) -} - -#[cfg(feature = "quinn")] -fn load_cert_chain( - path: &std::path::Path, -) -> Result>, String> { - let bytes = std::fs::read(path).map_err(|e| e.to_string())?; - let mut reader = std::io::BufReader::new(bytes.as_slice()); - rustls_pemfile::certs(&mut reader) - .collect::, _>>() - .map_err(|e| e.to_string()) -} - -#[cfg(feature = "quinn")] -fn load_private_key( - path: &std::path::Path, -) -> Result, String> { - let bytes = std::fs::read(path).map_err(|e| e.to_string())?; - let mut reader = std::io::BufReader::new(bytes.as_slice()); - match rustls_pemfile::private_key(&mut reader) { - Ok(Some(key)) => Ok(key), - Ok(None) => Err("no private key found in file".to_string()), - Err(e) => Err(e.to_string()), - } -} - -/// Client cert resolver that presents a single RFC 7250 raw public key (or -/// X.509 cert chain). For raw keys `only_raw_public_keys()` returns `true` so -/// rustls negotiates the RFC 7250 ClientCertificateType extension. -#[cfg(feature = "quinn")] -struct RawKeyClientCertResolver { - key: Arc, - raw_public_keys: bool, -} - -#[cfg(feature = "quinn")] -impl RawKeyClientCertResolver { - fn new(key: Arc) -> Self { - let raw_public_keys = key.cert.len() == 1 && is_ed25519_spki(&key.cert[0]); - Self { - key, - raw_public_keys, - } - } -} - -#[cfg(feature = "quinn")] -fn is_ed25519_spki(cert_der: &rustls::pki_types::CertificateDer<'_>) -> bool { - alknet_core::fingerprint::extract_ed25519_raw_key_from_spki(cert_der.as_ref()).is_some() -} - -#[cfg(feature = "quinn")] -impl std::fmt::Debug for RawKeyClientCertResolver { - fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - f.debug_struct("RawKeyClientCertResolver") - .field("raw_public_keys", &self.raw_public_keys) - .finish() - } -} - -#[cfg(feature = "quinn")] -impl rustls::client::ResolvesClientCert for RawKeyClientCertResolver { - fn resolve( - &self, - _root_hint_subjects: &[&[u8]], - _sigschemes: &[rustls::SignatureScheme], - ) -> Option> { - Some(Arc::clone(&self.key)) - } - - fn only_raw_public_keys(&self) -> bool { - self.raw_public_keys - } - - fn has_certs(&self) -> bool { - true - } -} - -/// Client cert resolver that presents no client cert (the `tls_identity: None` -/// or `SelfSigned` path). The server gets nothing to fingerprint — the -/// `PeerEntry` fingerprint → `peer_id` resolution path is not activated for -/// this connection. -#[cfg(feature = "quinn")] -struct NoClientCertResolver; - -#[cfg(feature = "quinn")] -impl std::fmt::Debug for NoClientCertResolver { - fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - f.debug_struct("NoClientCertResolver").finish() - } -} - -#[cfg(feature = "quinn")] -impl rustls::client::ResolvesClientCert for NoClientCertResolver { - fn resolve( - &self, - _root_hint_subjects: &[&[u8]], - _sigschemes: &[rustls::SignatureScheme], - ) -> Option> { - None - } - - fn has_certs(&self) -> bool { - false - } -} - -/// `ServerCertVerifier` that pins a specific fingerprint (ADR-034 §3, the -/// known-peer path). For `ed25519:` remotes the raw Ed25519 pub key is -/// extracted from the presented cert and matched against the pinned -/// fingerprint; for `SHA256:` remotes the cert DER is hashed and matched -/// against the pinned fingerprint. No match → verification failure (the -/// connection is rejected). The fingerprint IS the trust anchor — there is no -/// CA verification and no name verification, only the fingerprint pin. -/// -/// Handshake signatures are still verified (using the aws-lc-rs default -/// signature verification algorithms) so that a stolen-but-stale fingerprint -/// can't be replayed with a forged signature: the presenter must prove -/// possession of the private key corresponding to the pinned public key. -#[cfg(feature = "quinn")] -struct FingerprintPinVerifier { - fingerprint: String, - supported: rustls::crypto::WebPkiSupportedAlgorithms, -} - -#[cfg(feature = "quinn")] -impl FingerprintPinVerifier { - fn new(fingerprint: String, supported: rustls::crypto::WebPkiSupportedAlgorithms) -> Self { - Self { - fingerprint, - supported, - } - } -} - -#[cfg(feature = "quinn")] -impl std::fmt::Debug for FingerprintPinVerifier { - fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - f.debug_struct("FingerprintPinVerifier") - .field("fingerprint", &self.fingerprint) - .finish() - } -} - -#[cfg(feature = "quinn")] -impl rustls::client::danger::ServerCertVerifier for FingerprintPinVerifier { - fn verify_server_cert( - &self, - end_entity: &rustls::pki_types::CertificateDer<'_>, - _intermediates: &[rustls::pki_types::CertificateDer<'_>], - _server_name: &rustls::pki_types::ServerName<'_>, - _ocsp_response: &[u8], - _now: rustls::pki_types::UnixTime, - ) -> Result { - let presented = alknet_core::fingerprint::fingerprint_from_cert_der(end_entity.as_ref()) - .ok_or(rustls::Error::General( - "fingerprint pin: failed to compute fingerprint from presented cert".to_string(), - ))?; - if presented == self.fingerprint { - Ok(rustls::client::danger::ServerCertVerified::assertion()) - } else { - Err(rustls::Error::General(format!( - "fingerprint pin mismatch: expected {} got {}", - self.fingerprint, presented - ))) - } - } - - fn verify_tls12_signature( - &self, - message: &[u8], - cert: &rustls::pki_types::CertificateDer<'_>, - dss: &rustls::DigitallySignedStruct, - ) -> Result { - if alknet_core::fingerprint::extract_ed25519_raw_key_from_spki(cert.as_ref()).is_some() { - let spki = rustls::pki_types::SubjectPublicKeyInfoDer::from(cert.as_ref().to_vec()); - rustls::crypto::verify_tls13_signature_with_raw_key( - message, - &spki, - dss, - &self.supported, - ) - } else { - rustls::crypto::verify_tls12_signature(message, cert, dss, &self.supported) - } - } - - fn verify_tls13_signature( - &self, - message: &[u8], - cert: &rustls::pki_types::CertificateDer<'_>, - dss: &rustls::DigitallySignedStruct, - ) -> Result { - if alknet_core::fingerprint::extract_ed25519_raw_key_from_spki(cert.as_ref()).is_some() { - let spki = rustls::pki_types::SubjectPublicKeyInfoDer::from(cert.as_ref().to_vec()); - rustls::crypto::verify_tls13_signature_with_raw_key( - message, - &spki, - dss, - &self.supported, - ) - } else { - rustls::crypto::verify_tls13_signature(message, cert, dss, &self.supported) - } - } - - fn supported_verify_schemes(&self) -> Vec { - self.supported.supported_schemes() - } -} - -#[cfg(feature = "quinn")] -#[derive(Clone)] -struct Ed25519SigningKey { - key: alknet_core::config::Ed25519SecretKey, -} - -#[cfg(feature = "quinn")] -impl Ed25519SigningKey { - fn new(key: alknet_core::config::Ed25519SecretKey) -> Self { - Self { key } - } - - fn spki_public_key(&self) -> rustls::pki_types::SubjectPublicKeyInfoDer<'static> { - rustls::sign::public_key_to_spki( - &rustls::pki_types::alg_id::ED25519, - self.key.public().as_bytes(), - ) - } -} - -#[cfg(feature = "quinn")] -impl std::fmt::Debug for Ed25519SigningKey { - fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - f.debug_struct("Ed25519SigningKey").finish() - } -} - -#[cfg(feature = "quinn")] -impl rustls::sign::SigningKey for Ed25519SigningKey { - fn choose_scheme( - &self, - offered: &[rustls::SignatureScheme], - ) -> Option> { - if offered.contains(&rustls::SignatureScheme::ED25519) { - Some(Box::new(self.clone())) - } else { - None - } - } - - fn algorithm(&self) -> rustls::SignatureAlgorithm { - rustls::SignatureAlgorithm::ED25519 - } - - fn public_key(&self) -> Option> { - Some(self.spki_public_key()) - } -} - -#[cfg(feature = "quinn")] -impl rustls::sign::Signer for Ed25519SigningKey { - fn sign(&self, message: &[u8]) -> Result, rustls::Error> { - Ok(self.key.sign(message).to_bytes().to_vec()) - } - - fn scheme(&self) -> rustls::SignatureScheme { - rustls::SignatureScheme::ED25519 - } -} - #[cfg(test)] mod tests { use super::*; @@ -649,19 +157,6 @@ mod tests { .await } - #[test] - fn call_credentials_builder_methods() { - let creds = CallCredentials::new().with_remote_identity(RemoteIdentity { - fingerprint: "SHA256:abc".to_string(), - }); - assert_eq!( - creds.remote_identity.as_ref().unwrap().fingerprint, - "SHA256:abc" - ); - assert!(creds.tls_identity.is_none()); - assert!(creds.auth_token.is_none()); - } - #[tokio::test] async fn external_op_dispatches_and_populates_capabilities() { let registry = registry_with_caps(); @@ -706,225 +201,5 @@ mod tests { fn call_client_is_send_sync() { fn assert_send_sync() {} assert_send_sync::(); - assert_send_sync::(); - assert_send_sync::(); - } - - #[cfg(feature = "quinn")] - fn build_ed25519_spki_der(raw_key: &[u8; 32]) -> Vec { - let spki = rustls::sign::public_key_to_spki(&rustls::pki_types::alg_id::ED25519, raw_key); - spki.to_vec() - } - - #[cfg(feature = "quinn")] - fn build_x509_cert_der() -> rustls::pki_types::CertificateDer<'static> { - let key_pair = rcgen::KeyPair::generate().expect("key gen"); - let params = rcgen::CertificateParams::default(); - let cert = params.self_signed(&key_pair).expect("self-signed cert"); - cert.der().clone() - } - - #[cfg(feature = "quinn")] - fn aws_lc_rs_provider() -> Arc { - Arc::new(rustls::crypto::aws_lc_rs::default_provider()) - } - - #[cfg(feature = "quinn")] - fn verify_pin( - verifier: &FingerprintPinVerifier, - cert_der: rustls::pki_types::CertificateDer<'_>, - ) -> Result { - use rustls::client::danger::ServerCertVerifier; - let server_name: rustls::pki_types::ServerName<'static> = - "alknet".try_into().expect("server name"); - verifier.verify_server_cert( - &cert_der, - &[], - &server_name, - &[], - rustls::pki_types::UnixTime::now(), - ) - } - - #[cfg(feature = "quinn")] - #[test] - fn fingerprint_pin_verifier_matches_correct_ed25519_fingerprint() { - let sk = alknet_core::config::Ed25519SecretKey::generate(); - let raw_key = sk.public().to_bytes(); - let spki_der = build_ed25519_spki_der(&raw_key); - let fingerprint = - alknet_core::fingerprint::fingerprint_from_cert_der(&spki_der).expect("fingerprint"); - let verifier = FingerprintPinVerifier::new( - fingerprint, - aws_lc_rs_provider().signature_verification_algorithms, - ); - let cert = rustls::pki_types::CertificateDer::from(spki_der); - let result = verify_pin(&verifier, cert); - assert!( - result.is_ok(), - "FingerprintPinVerifier must accept a cert whose fingerprint matches the pin" - ); - } - - #[cfg(feature = "quinn")] - #[test] - fn fingerprint_pin_verifier_rejects_wrong_ed25519_fingerprint() { - let sk = alknet_core::config::Ed25519SecretKey::generate(); - let raw_key = sk.public().to_bytes(); - let spki_der = build_ed25519_spki_der(&raw_key); - let other_sk = alknet_core::config::Ed25519SecretKey::generate(); - let other_fp = format!("ed25519:{}", hex::encode(other_sk.public().to_bytes())); - let verifier = FingerprintPinVerifier::new( - other_fp, - aws_lc_rs_provider().signature_verification_algorithms, - ); - let cert = rustls::pki_types::CertificateDer::from(spki_der); - let result = verify_pin(&verifier, cert); - assert!( - result.is_err(), - "FingerprintPinVerifier must reject a cert whose fingerprint does not match the pin" - ); - } - - #[cfg(feature = "quinn")] - #[test] - fn fingerprint_pin_verifier_matches_correct_sha256_fingerprint() { - let cert_der = build_x509_cert_der(); - let fingerprint = alknet_core::fingerprint::fingerprint_from_cert_der(cert_der.as_ref()) - .expect("fingerprint"); - let verifier = FingerprintPinVerifier::new( - fingerprint, - aws_lc_rs_provider().signature_verification_algorithms, - ); - let result = verify_pin(&verifier, cert_der); - assert!( - result.is_ok(), - "FingerprintPinVerifier must accept an X.509 cert whose SHA256 fingerprint matches" - ); - } - - #[cfg(feature = "quinn")] - #[test] - fn fingerprint_pin_verifier_rejects_wrong_sha256_fingerprint() { - let cert_der = build_x509_cert_der(); - let verifier = FingerprintPinVerifier::new( - "SHA256:0000000000000000000000000000000000000000000000000000000000000000".to_string(), - aws_lc_rs_provider().signature_verification_algorithms, - ); - let result = verify_pin(&verifier, cert_der); - assert!( - result.is_err(), - "FingerprintPinVerifier must reject an X.509 cert whose SHA256 does not match" - ); - } - - #[cfg(feature = "quinn")] - #[test] - fn select_server_verifier_returns_ca_verifier_for_none() { - let provider = aws_lc_rs_provider(); - let remote_identity: Option = None; - let verifier = select_server_verifier(&provider, &remote_identity); - assert!( - verifier.is_ok(), - "select_server_verifier must succeed for None (CA path)" - ); - let debug = format!("{:?}", verifier.unwrap()); - assert!( - debug.contains("WebPkiServerVerifier"), - "None must select WebPkiServerVerifier (CA verification), got: {debug}" - ); - } - - #[cfg(feature = "quinn")] - #[test] - fn select_server_verifier_returns_fingerprint_pin_for_some() { - let provider = aws_lc_rs_provider(); - let remote_identity = Some(RemoteIdentity { - fingerprint: "ed25519:abc".to_string(), - }); - let verifier = select_server_verifier(&provider, &remote_identity); - assert!( - verifier.is_ok(), - "select_server_verifier must succeed for Some (fingerprint pin path)" - ); - let debug = format!("{:?}", verifier.unwrap()); - assert!( - debug.contains("FingerprintPinVerifier"), - "Some must select FingerprintPinVerifier, got: {debug}" - ); - } - - #[cfg(feature = "quinn")] - #[test] - fn build_client_auth_presents_ed25519_raw_key_without_error() { - let provider = aws_lc_rs_provider(); - let sk = alknet_core::config::Ed25519SecretKey::generate(); - let tls_identity = Some(alknet_core::config::TlsIdentity::RawKey(sk)); - let resolver = build_client_auth(&provider, &tls_identity); - assert!( - resolver.is_ok(), - "build_client_auth must build a resolver for a RawKey identity" - ); - let resolver = resolver.unwrap(); - assert!( - resolver.only_raw_public_keys(), - "RawKey client auth resolver must present raw public keys (RFC 7250)" - ); - assert!( - resolver.has_certs(), - "RawKey client auth resolver must report it has a cert to present" - ); - } - - #[cfg(feature = "quinn")] - #[test] - fn build_client_auth_none_resolves_to_no_client_cert() { - let provider = aws_lc_rs_provider(); - let tls_identity: Option = None; - let resolver = build_client_auth(&provider, &tls_identity) - .expect("build_client_auth must succeed for None"); - assert!( - !resolver.has_certs(), - "NoClientCertResolver must report no certs (no client cert presented)" - ); - } - - #[cfg(feature = "quinn")] - #[test] - fn build_quinn_client_config_with_raw_key_identity_builds_without_error() { - let sk = alknet_core::config::Ed25519SecretKey::generate(); - let credentials = CallCredentials::new() - .with_tls_identity(alknet_core::config::TlsIdentity::RawKey(sk)) - .with_remote_identity(RemoteIdentity { - fingerprint: "ed25519:deadbeef".to_string(), - }); - let config = build_quinn_client_config(&credentials, b"alknet/call"); - assert!( - config.is_ok(), - "build_quinn_client_config must build with a RawKey identity + pinned fingerprint" - ); - } - - #[cfg(feature = "quinn")] - #[test] - fn build_quinn_client_config_with_no_remote_identity_builds_without_error() { - let sk = alknet_core::config::Ed25519SecretKey::generate(); - let credentials = - CallCredentials::new().with_tls_identity(alknet_core::config::TlsIdentity::RawKey(sk)); - let config = build_quinn_client_config(&credentials, b"alknet/call"); - assert!( - config.is_ok(), - "build_quinn_client_config must build for the None + CA-verification path" - ); - } - - #[test] - fn remote_identity_none_is_load_bearing_not_defaulted() { - let creds = CallCredentials::new(); - assert!( - creds.remote_identity.is_none(), - "CallCredentials::new() must keep remote_identity as None (the load-bearing \ - public-X.509-endpoint state), not default it to a placeholder" - ); } } diff --git a/crates/alknet-call/src/client/from_call.rs b/crates/alknet-call/src/client/from_call.rs index 84cb1e2..9f2c64e 100644 --- a/crates/alknet-call/src/client/from_call.rs +++ b/crates/alknet-call/src/client/from_call.rs @@ -73,7 +73,7 @@ impl FromCallConfig { /// v1 defaults (two-way doors recorded in `client-and-adapters.md`): /// - auto-on-reconnect: the overlay is per-connection (Layer 2, ADR-024), so /// re-import on reconnect is naturally scoped; the assembly layer calls -/// `from_call` immediately after `connect()`. +/// `from_call` immediately after `AlknetClient::dial_*` + `spawn_dispatch`. /// - same-peer collision = error: two ops with the same name from the same /// peer (after applying the optional prefix) → `AdapterError::SamePeerCollision`. /// Cross-peer collision dissolves (ADR-029 §5). @@ -127,13 +127,11 @@ fn build_bundles( OperationType::Subscription => HandlerKind::Stream(make_streaming_forwarding_handler( Arc::new(op_summary.connection.clone()), remote_name, - op_summary.credentials_auth_token.clone(), )), OperationType::Query | OperationType::Mutation => { HandlerKind::Once(make_forwarding_handler( Arc::new(op_summary.connection.clone()), remote_name, - op_summary.credentials_auth_token.clone(), )) } }; @@ -155,7 +153,6 @@ struct OpSummary { name: String, schema: Value, connection: CallConnection, - credentials_auth_token: Option, } async fn discover_operations(connection: &CallConnection) -> Result, AdapterError> { @@ -182,7 +179,6 @@ async fn discover_operations(connection: &CallConnection) -> Result AccessControl { /// Per ADR-032 §3, the handler populates `forwarded_for` on the /// `call.requested` payload from the hub's `OperationContext.identity` (the /// end user the hub authenticated). The hub authenticates as itself when -/// forwarding — the `credentials_auth_token`, when present, is the hub's own -/// call-protocol-level token placed in the payload's `auth_token` field. The -/// spoke authorizes the hub (its direct caller); `forwarded_for` is metadata, -/// never read by `AccessControl::check`. +/// forwarding. The spoke authorizes the hub (its direct caller); +/// `forwarded_for` is metadata, never read by `AccessControl::check`. /// /// If `context.identity` is `None` (the hub chose not to disclose, or has not /// authenticated an originator), `forwarded_for` is omitted — the spoke @@ -340,16 +334,14 @@ fn parse_access_control(v: &Value) -> AccessControl { fn make_forwarding_handler( connection: Arc, remote_name: String, - credentials_auth_token: Option, ) -> Handler { use crate::registry::registration::make_handler; make_handler(move |input, context| { let connection = Arc::clone(&connection); let remote_name = remote_name.clone(); - let auth_token = credentials_auth_token.clone(); async move { let payload = - build_forwarded_payload(&remote_name, input, &context, auth_token.as_deref()); + build_forwarded_payload(&remote_name, input, &context); // The forwarding handler invokes the remote op via the // CallConnection. The parent_request_id participates in the abort // cascade (ADR-016 §6): if the parent is aborted, the cascade @@ -372,27 +364,25 @@ fn make_forwarding_handler( /// `call.aborted` drops it (ADR-049 §8). No truncation, no first-value /// fallback. /// -/// `forwarded_for` is populated from `context.identity` (ADR-032 §3) and -/// `auth_token` from the hub's own call-protocol token, exactly as the -/// request/response forwarding handler does — both via `build_forwarded_payload` -/// (no new payload-construction code). The `subscribe_with_payload` path -/// registers the request in `PendingRequestMap`, so the abort cascade -/// (ADR-016 §6) is already wired: a parent abort drops the -/// `SubscriptionStream`, which sends `call.aborted` to the remote node. +/// `forwarded_for` is populated from `context.identity` (ADR-032 §3), exactly +/// as the request/response forwarding handler does — both via +/// `build_forwarded_payload` (no new payload-construction code). The +/// `subscribe_with_payload` path registers the request in +/// `PendingRequestMap`, so the abort cascade (ADR-016 §6) is already wired: +/// a parent abort drops the `SubscriptionStream`, which sends `call.aborted` +/// to the remote node. fn make_streaming_forwarding_handler( connection: Arc, remote_name: String, - credentials_auth_token: Option, ) -> StreamingHandler { use crate::registry::registration::make_streaming_handler; use futures::stream::{once, StreamExt}; make_streaming_handler(move |input, context| { let connection = Arc::clone(&connection); let remote_name = remote_name.clone(); - let auth_token = credentials_auth_token.clone(); once(async move { let payload = - build_forwarded_payload(&remote_name, input, &context, auth_token.as_deref()); + build_forwarded_payload(&remote_name, input, &context); connection.subscribe_with_payload(payload).await }) .flatten() @@ -402,13 +392,11 @@ fn make_streaming_forwarding_handler( /// Build the `call.requested` payload for a forwarded call, populating /// `forwarded_for` from the hub's `OperationContext.identity` (ADR-032 §3). /// `forwarded_for` is omitted when `context.identity` is `None` (the hub -/// chooses not to disclose the originator). The `auth_token` field is set to -/// the hub's own call-protocol token when present. +/// chooses not to disclose the originator). fn build_forwarded_payload( operation_id: &str, input: Value, context: &OperationContext, - auth_token: Option<&str>, ) -> Value { let mut payload = serde_json::Map::new(); payload.insert( @@ -421,9 +409,6 @@ fn build_forwarded_payload( payload.insert("forwarded_for".to_string(), value); } } - if let Some(token) = auth_token { - payload.insert("auth_token".to_string(), Value::String(token.to_string())); - } Value::Object(payload) } @@ -571,7 +556,6 @@ mod tests { let handler = make_forwarding_handler( Arc::new(CallConnection::new(stub_connection())), "worker/echo".to_string(), - None, ); let reg = HandlerRegistration::new( spec, @@ -638,33 +622,22 @@ mod tests { #[test] fn build_forwarded_payload_populates_forwarded_for_from_context_identity() { let ctx = test_context(Some(alice_identity())); - let payload = build_forwarded_payload("fs/readFile", json!({"p": 1}), &ctx, None); + let payload = build_forwarded_payload("fs/readFile", json!({"p": 1}), &ctx); assert_eq!(payload["operationId"], "fs/readFile"); assert_eq!(payload["input"], json!({"p": 1})); let forwarded_for = payload.get("forwarded_for").expect("forwarded_for present"); assert_eq!(forwarded_for["id"], "alice"); assert_eq!(forwarded_for["scopes"][0], "fs:read"); - assert!(payload.get("auth_token").is_none()); } #[test] fn build_forwarded_payload_omits_forwarded_for_when_context_identity_is_none() { let ctx = test_context(None); - let payload = build_forwarded_payload("fs/readFile", json!({}), &ctx, None); + let payload = build_forwarded_payload("fs/readFile", json!({}), &ctx); assert!(payload.get("forwarded_for").is_none()); - assert!(payload.get("auth_token").is_none()); assert_eq!(payload["operationId"], "fs/readFile"); } - #[test] - fn build_forwarded_payload_sets_auth_token_when_provided() { - let ctx = test_context(Some(alice_identity())); - let payload = - build_forwarded_payload("fs/readFile", json!({}), &ctx, Some("alk_hub_token")); - assert_eq!(payload["auth_token"], "alk_hub_token"); - assert_eq!(payload["forwarded_for"]["id"], "alice"); - } - /// Verify the forwarding handler actually populates `forwarded_for` on /// the wire payload it sends. We intercept the payload by using a handler /// that records the payload passed to `call_with_payload`. Since @@ -687,7 +660,7 @@ mod tests { let captured = Arc::clone(&captured); let remote_name = "fs/readFile".to_string(); async move { - let payload = build_forwarded_payload(&remote_name, input, &context, None); + let payload = build_forwarded_payload(&remote_name, input, &context); *captured.lock().unwrap() = Some(payload.clone()); let response = conn.call_with_payload(payload).await; ResponseEnvelope { @@ -718,7 +691,7 @@ mod tests { let captured = Arc::clone(&captured); let remote_name = "fs/readFile".to_string(); async move { - let payload = build_forwarded_payload(&remote_name, input, &context, None); + let payload = build_forwarded_payload(&remote_name, input, &context); *captured.lock().unwrap() = Some(payload.clone()); let response = conn.call_with_payload(payload).await; ResponseEnvelope { @@ -745,7 +718,6 @@ mod tests { name: name.to_string(), schema: sample_schema_json(name, "query"), connection: conn.clone(), - credentials_auth_token: None, } } @@ -754,7 +726,6 @@ mod tests { name: name.to_string(), schema: sample_schema_json(name, op_type), connection: conn.clone(), - credentials_auth_token: None, } } @@ -949,7 +920,7 @@ mod tests { let remote_name = "events/stream".to_string(); use futures::stream::{once, StreamExt}; once(async move { - let payload = build_forwarded_payload(&remote_name, input, &context, None); + let payload = build_forwarded_payload(&remote_name, input, &context); *captured.lock().unwrap() = Some(payload.clone()); conn.subscribe_with_payload(payload).await }) @@ -999,7 +970,7 @@ mod tests { let remote_name = "events/stream".to_string(); use futures::stream::{once, StreamExt}; once(async move { - let payload = build_forwarded_payload(&remote_name, input, &context, None); + let payload = build_forwarded_payload(&remote_name, input, &context); *captured.lock().unwrap() = Some(payload.clone()); conn.subscribe_with_payload(payload).await }) @@ -1018,45 +989,6 @@ mod tests { assert_eq!(payload["operationId"], "events/stream"); } - /// The streaming forwarding handler populates `auth_token` when the hub's - /// own call-protocol token is provided. - #[tokio::test] - async fn streaming_forwarding_handler_sets_auth_token_when_provided() { - use futures::stream::StreamExt; - - let conn = Arc::new(CallConnection::new(stub_connection())); - let captured_payload = Arc::new(StdMutex::new(None::)); - let captured = Arc::clone(&captured_payload); - - let handler: StreamingHandler = { - let conn = Arc::clone(&conn); - make_streaming_handler(move |input, context| { - let conn = Arc::clone(&conn); - let captured = Arc::clone(&captured); - let remote_name = "events/stream".to_string(); - use futures::stream::{once, StreamExt}; - once(async move { - let payload = build_forwarded_payload( - &remote_name, - input, - &context, - Some("alk_hub_token"), - ); - *captured.lock().unwrap() = Some(payload.clone()); - conn.subscribe_with_payload(payload).await - }) - .flatten() - }) - }; - - let ctx = test_context(Some(alice_identity())); - let mut stream = handler(json!({}), ctx); - let _ = stream.next().await; - let payload = captured_payload.lock().unwrap().clone().expect("captured"); - assert_eq!(payload["auth_token"], "alk_hub_token"); - assert_eq!(payload["forwarded_for"]["id"], "alice"); - } - /// `make_streaming_forwarding_handler` produces a `StreamingHandler` (not a /// `Handler`) — verifies the helper returns the right type and that /// `build_bundles` wires it into `HandlerKind::Stream`. @@ -1065,7 +997,6 @@ mod tests { let handler = make_streaming_forwarding_handler( Arc::new(CallConnection::new(stub_connection())), "events/stream".to_string(), - None, ); let reg = HandlerRegistration::new( OperationSpec::new( diff --git a/crates/alknet-call/src/client/mod.rs b/crates/alknet-call/src/client/mod.rs index ea4f758..f506bfc 100644 --- a/crates/alknet-call/src/client/mod.rs +++ b/crates/alknet-call/src/client/mod.rs @@ -9,7 +9,7 @@ mod call_client; mod from_call; -pub use call_client::{CallClient, CallCredentials, ClientError, RemoteIdentity}; +pub use call_client::CallClient; pub use from_call::{from_call, FromCallConfig}; use crate::registry::registration::HandlerRegistration; diff --git a/crates/alknet-call/src/lib.rs b/crates/alknet-call/src/lib.rs index cffebf5..6589ec8 100644 --- a/crates/alknet-call/src/lib.rs +++ b/crates/alknet-call/src/lib.rs @@ -1,4 +1,4 @@ -//! alknet-call: Structured RPC over QUIC — operations, streaming, service discovery. +//! alknet-call: Structured RPC — operations, streaming, service discovery. //! //! Implements [`alknet_core::types::ProtocolHandler`] on ALPN `alknet/call`. //! diff --git a/crates/alknet-call/tests/two_node_call.rs b/crates/alknet-call/tests/two_node_call.rs deleted file mode 100644 index 0bec322..0000000 --- a/crates/alknet-call/tests/two_node_call.rs +++ /dev/null @@ -1,331 +0,0 @@ -//! Integration test: two-node `alknet/call` round-trip over a real QUIC -//! loopback. A `CallAdapter` server accepts, a `CallClient` connects, and -//! the client calls back into the server (connection symmetry, ADR-017 §2). -//! Verifies the shared dispatch loop works end-to-end. - -#![cfg(feature = "quinn")] - -use std::sync::Arc; -use std::time::Duration; - -use alknet_call::client::{CallClient, CallCredentials, RemoteIdentity}; -use alknet_call::protocol::adapter::CallAdapter; -use alknet_call::protocol::wire::ResponseEnvelope; -use alknet_call::registry::discovery::{ - services_list_handler, services_list_spec, services_schema_handler, services_schema_spec, -}; -use alknet_call::registry::registration::{ - make_handler, Handler, HandlerKind, HandlerRegistration, OperationProvenance, OperationRegistry, -}; -use alknet_call::registry::spec::{AccessControl, OperationSpec, OperationType, Visibility}; -use alknet_core::auth::{Identity, IdentityProvider}; -use alknet_core::types::{Capabilities, Connection, ProtocolHandler}; - -struct NoopIdentityProvider; -impl IdentityProvider for NoopIdentityProvider { - fn resolve_from_fingerprint(&self, _: &str) -> Option { - None - } - fn resolve_from_token(&self, _: &alknet_core::auth::AuthToken) -> Option { - None - } -} - -fn external_spec(name: &str) -> OperationSpec { - OperationSpec::new( - name, - OperationType::Query, - Visibility::External, - serde_json::json!({}), - serde_json::json!({}), - vec![], - AccessControl::default(), - None, - ) -} - -fn echo_handler() -> Handler { - make_handler(|input, context| async move { ResponseEnvelope::ok(context.request_id, input) }) -} - -/// Build a raw quinn server endpoint with a self-signed cert and the -/// `CallAdapter` accepting `alknet/call` connections. Returns -/// `(bound_addr, server_fingerprint, join_handle)` — the fingerprint is the -/// `SHA256:` of the self-signed cert DER, which the client pins via -/// `CallCredentials::with_remote_identity` (the known-peer path, ADR-034 §3). -/// The accept loop spawns a task per connection that hands the connection to -/// `CallAdapter::handle`. -async fn build_raw_quinn_server( - registry: Arc, -) -> (std::net::SocketAddr, String, tokio::task::JoinHandle<()>) { - let provider: Arc = Arc::new(NoopIdentityProvider); - let adapter = Arc::new(CallAdapter::new( - Arc::clone(®istry), - Arc::clone(&provider), - )); - - let key_pair = rcgen::KeyPair::generate().expect("key gen"); - let params = rcgen::CertificateParams::default(); - let cert = params.self_signed(&key_pair).expect("self-signed cert"); - let cert_der = cert.der().clone(); - let fingerprint = alknet_core::fingerprint::fingerprint_from_cert_der(cert_der.as_ref()) - .expect("cert produces fingerprint"); - let key_der = rustls::pki_types::PrivateKeyDer::Pkcs8( - rustls::pki_types::PrivatePkcs8KeyDer::from(key_pair.serialize_der()), - ); - - let provider_crypto = Arc::new(rustls::crypto::aws_lc_rs::default_provider()); - let mut server_config = rustls::ServerConfig::builder_with_provider(provider_crypto) - .with_safe_default_protocol_versions() - .unwrap() - .with_no_client_auth() - .with_single_cert(vec![cert_der], key_der) - .unwrap(); - server_config.alpn_protocols = vec![b"alknet/call".to_vec()]; - server_config.max_early_data_size = u32::MAX; - - let quic_server_config = - quinn::crypto::rustls::QuicServerConfig::try_from(server_config).unwrap(); - let quinn_server_config = quinn::ServerConfig::with_crypto(Arc::new(quic_server_config)); - - let quinn_endpoint = - quinn::Endpoint::server(quinn_server_config, "127.0.0.1:0".parse().unwrap()) - .expect("server bind"); - let bound_addr = quinn_endpoint.local_addr().expect("local addr"); - - let join = tokio::spawn(async move { - while let Some(incoming) = quinn_endpoint.accept().await { - let adapter = Arc::clone(&adapter); - tokio::spawn(async move { - let connecting = match incoming.accept() { - Ok(c) => c, - Err(_) => return, - }; - let conn = match connecting.await { - Ok(c) => c, - Err(_) => return, - }; - let alpn = b"alknet/call".to_vec(); - let conn = Connection::from_quinn_with_alpn(conn, alpn.clone()); - let auth = alknet_core::auth::AuthContext { - identity: None, - alpn, - remote_addr: conn.remote_addr(), - tls_client_fingerprint: None, - }; - let _ = adapter.handle(conn, &auth).await; - }); - } - }); - - (bound_addr, fingerprint, join) -} - -/// Build the server's registry: an echo op, a secret op, and the -/// services/list + services/schema discovery handlers. -fn build_server_registry() -> Arc { - let mut registry = OperationRegistry::new(); - registry - .register(HandlerRegistration::new( - external_spec("server/echo"), - HandlerKind::Once(echo_handler()), - OperationProvenance::Local, - None, - None, - Capabilities::new(), - )) - .unwrap(); - registry - .register(HandlerRegistration::new( - external_spec("server/secret"), - HandlerKind::Once(echo_handler()), - OperationProvenance::Local, - None, - None, - Capabilities::new().with_api_key("google", "server-secret".to_string()), - )) - .unwrap(); - let discovery_registry = Arc::new(registry); - let list_handler = services_list_handler(Arc::clone(&discovery_registry)); - let schema_handler = services_schema_handler(Arc::clone(&discovery_registry)); - let mut full = OperationRegistry::new(); - full.register(HandlerRegistration::new( - external_spec("server/echo"), - HandlerKind::Once(echo_handler()), - OperationProvenance::Local, - None, - None, - Capabilities::new(), - )) - .unwrap(); - full.register(HandlerRegistration::new( - external_spec("server/secret"), - HandlerKind::Once(echo_handler()), - OperationProvenance::Local, - None, - None, - Capabilities::new().with_api_key("google", "server-secret".to_string()), - )) - .unwrap(); - full.register(HandlerRegistration::new( - services_list_spec(), - HandlerKind::Once(list_handler), - OperationProvenance::Local, - None, - None, - Capabilities::new(), - )) - .unwrap(); - full.register(HandlerRegistration::new( - services_schema_spec(), - HandlerKind::Once(schema_handler), - OperationProvenance::Local, - None, - None, - Capabilities::new(), - )) - .unwrap(); - Arc::new(full) -} - -#[tokio::test(flavor = "multi_thread", worker_threads = 4)] -async fn two_node_call_round_trip() { - let server_registry = build_server_registry(); - let (server_addr, server_fingerprint, _server_join) = - build_raw_quinn_server(Arc::clone(&server_registry)).await; - - // Client side: a CallClient with its own ops so the server can call back - // (connection symmetry). Pin the server's self-signed cert fingerprint - // (the known-peer path, ADR-034 §3) — `WebPkiServerVerifier` would reject - // it as UnknownIssuer since the self-signed cert is not in the platform - // root store. - let mut client_registry = OperationRegistry::new(); - client_registry - .register(HandlerRegistration::new( - external_spec("client/echo"), - HandlerKind::Once(echo_handler()), - OperationProvenance::Local, - None, - None, - Capabilities::new(), - )) - .unwrap(); - let client_registry = Arc::new(client_registry); - let client = CallClient::new(Arc::clone(&client_registry), Arc::new(NoopIdentityProvider)); - - let credentials = CallCredentials::new().with_remote_identity(RemoteIdentity { - fingerprint: server_fingerprint, - }); - let conn = tokio::time::timeout( - Duration::from_secs(5), - client.connect(server_addr, credentials), - ) - .await - .expect("connect did not time out") - .expect("connect succeeds"); - - // Outbound call: client -> server's echo op. - let response = tokio::time::timeout( - Duration::from_secs(5), - conn.call("server/echo", serde_json::json!({"hi": 1})), - ) - .await - .expect("call did not time out"); - assert_eq!(response.result, Ok(serde_json::json!({"hi": 1}))); - - // Peer authorization is enforced by the AccessControl gate in - // OperationRegistry::invoke (ADR-029 §3) — exercised by the unit tests in - // `registry/registration.rs`. This integration test focuses on the QUIC - // connect path + shared dispatch loop working end-to-end (the call above - // proves the CallClient opened a real connection, the shared loop - // dispatched, and the CallConnection::call() round-tripped). -} - -#[tokio::test(flavor = "multi_thread", worker_threads = 4)] -async fn from_call_discovers_and_forwards_over_quic_loopback() { - use alknet_call::client::{from_call, FromCallConfig}; - use alknet_call::registry::context::ScopedPeerEnv; - - let server_registry = build_server_registry(); - let (server_addr, server_fingerprint, _server_join) = - build_raw_quinn_server(Arc::clone(&server_registry)).await; - - // Client with an empty registry — from_call will populate its overlay. - // Pin the server's self-signed cert fingerprint (ADR-034 §3 known-peer - // path). - let client_registry = Arc::new(OperationRegistry::new()); - let client = CallClient::new(Arc::clone(&client_registry), Arc::new(NoopIdentityProvider)); - - let credentials = CallCredentials::new().with_remote_identity(RemoteIdentity { - fingerprint: server_fingerprint, - }); - let conn = tokio::time::timeout( - Duration::from_secs(5), - client.connect(server_addr, credentials), - ) - .await - .expect("connect did not time out") - .expect("connect succeeds"); - - // from_call discovers the server's External ops (server/echo, server/secret - // — both External; services/list + services/schema themselves are External - // too) and builds FromCall forwarding-handler bundles. Register them in the - // connection's Layer 2 overlay. - let bundles = tokio::time::timeout( - Duration::from_secs(5), - from_call(&conn, FromCallConfig::new()), - ) - .await - .expect("from_call did not time out") - .expect("from_call succeeds"); - assert!( - !bundles.is_empty(), - "from_call must discover at least the server/echo op" - ); - conn.register_imported_all(bundles); - - // The overlay now contains the discovered ops. Verify the forwarding path - // by invoking the overlay env directly with a scoped context that allows - // server/echo — this is how a composing handler would call the imported op. - let env = conn.overlay_env(); - assert!( - env.contains("server/echo"), - "overlay must contain the imported server/echo op" - ); - - // Build a minimal parent context to invoke the overlay env (mirrors how a - // composing handler dispatches a child). - let scoped = ScopedPeerEnv::new(["server/echo"]); - let parent = alknet_call::registry::context::OperationContext { - request_id: "parent-1".to_string(), - parent_request_id: None, - identity: None, - handler_identity: None, - forwarded_for: None, - capabilities: Capabilities::new(), - metadata: Default::default(), - scoped_env: scoped, - env: env.clone(), - abort_policy: alknet_call::registry::context::AbortPolicy::default(), - deadline: Some(std::time::Instant::now() + Duration::from_secs(30)), - internal: true, - ownership: None, - }; - - let response = tokio::time::timeout( - Duration::from_secs(5), - env.invoke( - "server", - "echo", - serde_json::json!({"from_call": true}), - &parent, - ), - ) - .await - .expect("overlay invoke did not time out"); - assert_eq!( - response.result, - Ok(serde_json::json!({"from_call": true})), - "from_call forwarding handler must round-trip the input to the remote op" - ); -}