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
This commit is contained in:
deepseek-v4-pro committed 2026-07-17 14:01:51 +00:00
1 parent a0dbe4fd3c
commit 4fc6854846
7 files changed
+32 -1172

No files matched your search

Generated
-6
View File
@@ -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",
+2 -11
View File
@@ -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"
+9 -734
View File
@@ -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<TlsIdentity>,
/// Opaque call-protocol-level auth token, decrypted from the vault.
pub auth_token: Option<alknet_core::auth::AuthToken>,
/// 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<RemoteIdentity>,
}
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<CallConnection, ClientError> {
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<quinn::ClientConfig, String> {
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<rustls::crypto::CryptoProvider>,
tls_identity: &Option<TlsIdentity>,
) -> Result<Arc<dyn rustls::client::ResolvesClientCert>, 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<rustls::crypto::CryptoProvider>,
remote_identity: &Option<RemoteIdentity>,
) -> Result<Arc<dyn rustls::client::danger::ServerCertVerifier>, 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<rustls::RootCertStore, String> {
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<Vec<rustls::pki_types::CertificateDer<'static>>, 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::<Result<Vec<_>, _>>()
.map_err(|e| e.to_string())
}
#[cfg(feature = "quinn")]
fn load_private_key(
path: &std::path::Path,
) -> Result<rustls::pki_types::PrivateKeyDer<'static>, 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<rustls::sign::CertifiedKey>,
raw_public_keys: bool,
}
#[cfg(feature = "quinn")]
impl RawKeyClientCertResolver {
fn new(key: Arc<rustls::sign::CertifiedKey>) -> 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<Arc<rustls::sign::CertifiedKey>> {
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<Arc<rustls::sign::CertifiedKey>> {
None
}
fn has_certs(&self) -> bool {
false
}
}
/// `ServerCertVerifier` that pins a specific fingerprint (ADR-034 §3, the
/// known-peer path). For `ed25519:<hex>` remotes the raw Ed25519 pub key is
/// extracted from the presented cert and matched against the pinned
/// fingerprint; for `SHA256:<hex>` 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<rustls::client::danger::ServerCertVerified, rustls::Error> {
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<rustls::client::danger::HandshakeSignatureValid, rustls::Error> {
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<rustls::client::danger::HandshakeSignatureValid, rustls::Error> {
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<rustls::SignatureScheme> {
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<Box<dyn rustls::sign::Signer>> {
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<rustls::pki_types::SubjectPublicKeyInfoDer<'_>> {
Some(self.spki_public_key())
}
}
#[cfg(feature = "quinn")]
impl rustls::sign::Signer for Ed25519SigningKey {
fn sign(&self, message: &[u8]) -> Result<Vec<u8>, 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<T: Send + Sync>() {}
assert_send_sync::<CallClient>();
assert_send_sync::<CallCredentials>();
assert_send_sync::<RemoteIdentity>();
}
#[cfg(feature = "quinn")]
fn build_ed25519_spki_der(raw_key: &[u8; 32]) -> Vec<u8> {
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<rustls::crypto::CryptoProvider> {
Arc::new(rustls::crypto::aws_lc_rs::default_provider())
}
#[cfg(feature = "quinn")]
fn verify_pin(
verifier: &FingerprintPinVerifier,
cert_der: rustls::pki_types::CertificateDer<'_>,
) -> Result<rustls::client::danger::ServerCertVerified, rustls::Error> {
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<RemoteIdentity> = 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<alknet_core::config::TlsIdentity> = 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"
);
}
}
+19 -88
View File
@@ -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<String>,
}
async fn discover_operations(connection: &CallConnection) -> Result<Vec<OpSummary>, AdapterError> {
@@ -182,7 +179,6 @@ async fn discover_operations(connection: &CallConnection) -> Result<Vec<OpSummar
name: name.to_string(),
schema,
connection: connection.clone(),
credentials_auth_token: None,
});
}
Ok(summaries)
@@ -329,10 +325,8 @@ fn parse_access_control(v: &Value) -> 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<CallConnection>,
remote_name: String,
credentials_auth_token: Option<String>,
) -> 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<CallConnection>,
remote_name: String,
credentials_auth_token: Option<String>,
) -> 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::<Value>));
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(
+1 -1
View File
@@ -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;
+1 -1
View File
@@ -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`.
//!
-331
View File
@@ -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<Identity> {
None
}
fn resolve_from_token(&self, _: &alknet_core::auth::AuthToken) -> Option<Identity> {
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:<hex>` 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<OperationRegistry>,
) -> (std::net::SocketAddr, String, tokio::task::JoinHandle<()>) {
let provider: Arc<dyn IdentityProvider> = Arc::new(NoopIdentityProvider);
let adapter = Arc::new(CallAdapter::new(
Arc::clone(&registry),
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<OperationRegistry> {
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"
);
}