use std::{net::SocketAddr, sync::Arc}; use disasmer_core::{ ArtifactId, DataPlaneObject, DataPlaneScope, Digest, NativeQuicTransport, NodeEndpoint, NodeId, ProcessId, ProjectId, RendezvousRequest, TenantId, Transport, }; use quinn::rustls::{ pki_types::{CertificateDer, PrivateKeyDer, PrivatePkcs8KeyDer}, RootCertStore, }; use quinn::{ClientConfig, Endpoint, ServerConfig}; use serde::{Deserialize, Serialize}; use serde_json::json; #[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] struct QuicTransferRequest { scope: DataPlaneScope, authorization_digest: Digest, requested_bytes: u64, } #[tokio::main] async fn main() -> Result<(), Box> { let payload = b"artifact-bytes-over-rust-native-quic".to_vec(); let (cert, key) = self_signed_localhost_cert()?; let server_endpoint = Endpoint::server( ServerConfig::with_single_cert(vec![cert.clone()], key)?, "127.0.0.1:0".parse()?, )?; let server_addr = server_endpoint.local_addr()?; let scope = DataPlaneScope { tenant: TenantId::from("tenant"), project: ProjectId::from("project"), process: ProcessId::from("vp-quic"), object: DataPlaneObject::Artifact(ArtifactId::from("quic-artifact")), authorization_subject: "node-a-to-node-b".to_owned(), }; let transport = NativeQuicTransport; let plan = transport.plan_authenticated_direct_bulk_transfer( RendezvousRequest { scope: scope.clone(), source: endpoint("node-a", server_addr), destination: endpoint("node-b", "127.0.0.1:0".parse()?), }, true, "", )?; let expected_scope = plan.scope.clone(); let expected_digest = plan.authorization_digest.clone(); let expected_payload = payload.clone(); let server = tokio::spawn(async move { let incoming = server_endpoint .accept() .await .ok_or("server endpoint closed before accepting a QUIC connection")?; let connection = incoming.await?; let (mut send, mut recv) = connection.accept_bi().await?; let request_bytes = recv.read_to_end(64 * 1024).await?; let request: QuicTransferRequest = serde_json::from_slice(&request_bytes)?; if request.scope != expected_scope { return Err("QUIC request scope did not match the authorized data-plane scope".into()); } if request.authorization_digest != expected_digest { return Err( "QUIC request authorization digest did not match the rendezvous plan".into(), ); } send.write_all(&expected_payload).await?; send.finish()?; server_endpoint.wait_idle().await; Ok::>(request_bytes.len()) }); let mut roots = RootCertStore::empty(); roots.add(cert)?; let client_config = ClientConfig::with_root_certificates(Arc::new(roots))?; let mut client_endpoint = Endpoint::client("127.0.0.1:0".parse()?)?; client_endpoint.set_default_client_config(client_config); let connection = client_endpoint.connect(server_addr, "localhost")?.await?; let (mut send, mut recv) = connection.open_bi().await?; let request = QuicTransferRequest { scope: plan.scope.clone(), authorization_digest: plan.authorization_digest.clone(), requested_bytes: payload.len() as u64, }; let request_bytes = serde_json::to_vec(&request)?; send.write_all(&request_bytes).await?; send.finish()?; let received = recv.read_to_end(64 * 1024).await?; connection.close(0u32.into(), b"done"); client_endpoint.wait_idle().await; let server_received_request_bytes = server.await??; if received != payload { return Err("QUIC artifact payload did not round trip".into()); } println!( "{}", json!({ "kind": "disasmer_quic_smoke", "transport": format!("{:?}", transport.kind()), "rust_native_quic": true, "authenticated_direct_connection": transport.authenticated_direct_connections(), "coordinator_assisted_rendezvous": plan.coordinator_assisted_rendezvous, "coordinator_bulk_relay_allowed": plan.coordinator_bulk_relay_allowed, "source_node": plan.source.node, "destination_node": plan.destination.node, "scope": plan.scope, "request_bytes": request_bytes.len(), "server_received_request_bytes": server_received_request_bytes, "payload_bytes": received.len(), "authorization_digest": plan.authorization_digest, }) ); Ok(()) } fn self_signed_localhost_cert() -> Result< (CertificateDer<'static>, PrivateKeyDer<'static>), Box, > { let cert = rcgen::generate_simple_self_signed(vec!["localhost".to_owned()])?; let key = PrivateKeyDer::Pkcs8(PrivatePkcs8KeyDer::from(cert.signing_key.serialize_der())); Ok((cert.cert.into(), key)) } fn endpoint(name: &str, addr: SocketAddr) -> NodeEndpoint { NodeEndpoint { node: NodeId::from(name), advertised_addr: addr.to_string(), public_key_fingerprint: Digest::sha256(format!("{name}-public-key")), } }