Public dry run dryrun-224e0a3c717e
Source commit: 224e0a3c717e58cdb7bb61dcebbf5dadc4225d3a Public tree identity: sha256:2ea0b60d516db02a3ecdeb49982be1782cc3e7364a0e0eadc541c2e095d8b126
This commit is contained in:
commit
815fd6392d
111 changed files with 35316 additions and 0 deletions
138
crates/disasmer-node/src/bin/disasmer-quic-smoke.rs
Normal file
138
crates/disasmer-node/src/bin/disasmer-quic-smoke.rs
Normal file
|
|
@ -0,0 +1,138 @@
|
|||
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<dyn std::error::Error + Send + Sync>> {
|
||||
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::<usize, Box<dyn std::error::Error + Send + Sync>>(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<dyn std::error::Error + Send + Sync>,
|
||||
> {
|
||||
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")),
|
||||
}
|
||||
}
|
||||
Loading…
Add table
Add a link
Reference in a new issue