Some checks failed
CI / rustfmt (push) Failing after 17s
CI / wasm build (push) Failing after 1m13s
CI / clippy (push) Failing after 1m17s
CI / example (push) Failing after 1m30s
CI / test (push) Successful in 1m50s
CI / duplicate code (push) Failing after 31s
CI / web client (push) Failing after 31s
CI / cargo-machete (push) Successful in 1m15s
CI / cargo-deny (push) Failing after 2m26s
201 lines
6.3 KiB
Rust
201 lines
6.3 KiB
Rust
use std::net::{IpAddr, Ipv4Addr};
|
|
|
|
use mtp_codec::{CommunicationType, CommunicationValue, DataType, DataValue, TypeMap};
|
|
use mtp_transport::{Host, Policy, Receiver, Sender, connect, host};
|
|
|
|
fn generate_self_signed_cert() -> (Vec<u8>, Vec<u8>) {
|
|
let key_pair = rcgen::KeyPair::generate().unwrap();
|
|
let params =
|
|
rcgen::CertificateParams::new(vec!["localhost".into(), "127.0.0.1".into()]).unwrap();
|
|
let cert = params.self_signed(&key_pair).unwrap();
|
|
let cert_pem = cert.pem();
|
|
let key_pem = key_pair.serialize_pem();
|
|
(cert_pem.into_bytes(), key_pem.into_bytes())
|
|
}
|
|
|
|
async fn start_test_host(cert_pem: Vec<u8>, key_pem: Vec<u8>) -> Host {
|
|
host(
|
|
IpAddr::V4(Ipv4Addr::LOCALHOST),
|
|
0,
|
|
cert_pem,
|
|
key_pem,
|
|
Policy::default(),
|
|
)
|
|
.await
|
|
.unwrap()
|
|
}
|
|
|
|
async fn connect_to_host(h: &Host, cert_pem: Vec<u8>) -> (Sender, Receiver) {
|
|
let url = format!("https://127.0.0.1:{}", h.local_addr().port());
|
|
connect(&url, Some(cert_pem), Policy::default())
|
|
.await
|
|
.unwrap()
|
|
}
|
|
|
|
async fn connected_pair() -> (Host, Sender, Receiver, Sender, Receiver) {
|
|
let (cert_pem, key_pem) = generate_self_signed_cert();
|
|
let mut h = start_test_host(cert_pem.clone(), key_pem).await;
|
|
let (client_tx, client_rx) = connect_to_host(&h, cert_pem).await;
|
|
let (host_tx, host_rx) = h.next().await.unwrap();
|
|
(h, client_tx, client_rx, host_tx, host_rx)
|
|
}
|
|
|
|
fn numbered_message(comm_type: CommunicationType, value: u128, tm: &TypeMap) -> CommunicationValue {
|
|
CommunicationValue::new(comm_type).add_data(
|
|
DataType::PqSignature.to_id(tm),
|
|
DataValue::UnsignedNumber(value),
|
|
)
|
|
}
|
|
|
|
fn assert_numbered_message(
|
|
message: &CommunicationValue,
|
|
comm_type: CommunicationType,
|
|
value: u128,
|
|
tm: &TypeMap,
|
|
) {
|
|
assert_eq!(message.get_type(), comm_type.to_id(tm));
|
|
assert_eq!(
|
|
message.get_data(DataType::PqSignature).clone(),
|
|
DataValue::UnsignedNumber(value)
|
|
);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_host_start_and_stop() {
|
|
let (cert_pem, key_pem) = generate_self_signed_cert();
|
|
let h = start_test_host(cert_pem, key_pem).await;
|
|
let addr = h.local_addr();
|
|
// Port should be non-zero (OS-assigned)
|
|
assert!(addr.port() > 0);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_send_receive_roundtrip() {
|
|
let (_h, client_tx, client_rx, host_tx, host_rx) = connected_pair().await;
|
|
|
|
let tm = TypeMap::latest();
|
|
|
|
// Client sends a simple message
|
|
let msg = numbered_message(CommunicationType::Ping, 42, &tm);
|
|
client_tx.send(&msg).await.unwrap();
|
|
|
|
// Host receives it
|
|
let received = host_rx.receive().await.unwrap();
|
|
assert_numbered_message(&received, CommunicationType::Ping, 42, &tm);
|
|
|
|
// Host sends a response
|
|
let resp = numbered_message(CommunicationType::Pong, 99, &tm);
|
|
host_tx.send(&resp).await.unwrap();
|
|
|
|
// Client receives it
|
|
let client_received = client_rx.receive().await.unwrap();
|
|
assert_numbered_message(&client_received, CommunicationType::Pong, 99, &tm);
|
|
|
|
// Close both sides
|
|
client_tx.close();
|
|
host_tx.close();
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_concurrent_messages() {
|
|
let (_h, client_tx, _client_rx, _host_tx, host_rx) = connected_pair().await;
|
|
|
|
let tm = TypeMap::latest();
|
|
|
|
// Send 5 messages in sequence
|
|
for i in 0..5u128 {
|
|
let msg = numbered_message(CommunicationType::Ping, i, &tm);
|
|
client_tx.send(&msg).await.unwrap();
|
|
}
|
|
|
|
// Receive all 5 in order
|
|
for i in 0..5u128 {
|
|
let received = host_rx.receive().await.unwrap();
|
|
assert_numbered_message(&received, CommunicationType::Ping, i, &tm);
|
|
}
|
|
|
|
// Send 3 responses back
|
|
for i in 0..3u128 {
|
|
let msg = numbered_message(CommunicationType::Pong, i * 10, &tm);
|
|
client_tx.send(&msg).await.unwrap();
|
|
}
|
|
|
|
for i in 0..3u128 {
|
|
let received = host_rx.receive().await.unwrap();
|
|
assert_numbered_message(&received, CommunicationType::Pong, i * 10, &tm);
|
|
}
|
|
|
|
client_tx.close();
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_close_detection() {
|
|
let (_h, client_tx, _client_rx, _host_tx, host_rx) = connected_pair().await;
|
|
|
|
// Send a message then close
|
|
let msg = CommunicationValue::new(CommunicationType::Ping);
|
|
client_tx.send(&msg).await.unwrap();
|
|
client_tx.close();
|
|
|
|
// Host should still receive the message
|
|
let tm = TypeMap::latest();
|
|
let received = host_rx.receive().await.unwrap();
|
|
assert_eq!(received.get_type(), CommunicationType::Ping.to_id(&tm));
|
|
|
|
// Host should get an error or closed signal on next receive
|
|
let result = host_rx.receive().await;
|
|
assert!(result.is_err());
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_host_shutdown_stops_accepting() {
|
|
let (cert_pem, key_pem) = generate_self_signed_cert();
|
|
let mut h = start_test_host(cert_pem.clone(), key_pem).await;
|
|
let url = format!("https://127.0.0.1:{}", h.local_addr().port());
|
|
|
|
// A connection succeeds while the host is accepting.
|
|
let (_c_tx, _c_rx) = connect(&url, Some(cert_pem.clone()), Policy::default())
|
|
.await
|
|
.unwrap();
|
|
let _accepted = h.next().await.unwrap();
|
|
|
|
// After shutdown the accept task is aborted and its endpoint is dropped, so
|
|
// new connections no longer succeed. Guard with a timeout so a hung connect
|
|
// still fails the assertion rather than blocking the test.
|
|
h.shutdown();
|
|
|
|
let result = tokio::time::timeout(
|
|
std::time::Duration::from_secs(5),
|
|
connect(&url, Some(cert_pem), Policy::default()),
|
|
)
|
|
.await;
|
|
assert!(
|
|
matches!(result, Err(_) | Ok(Err(_))),
|
|
"connect should not succeed after host shutdown"
|
|
);
|
|
}
|
|
|
|
#[tokio::test]
|
|
async fn test_drop_receiver_keeps_sender_alive() {
|
|
let (_h, client_tx, client_rx, host_tx, host_rx) = connected_pair().await;
|
|
|
|
// Client sends a message the host receives.
|
|
let msg = CommunicationValue::new(CommunicationType::Ping);
|
|
client_tx.send(&msg).await.unwrap();
|
|
let _ = host_rx.receive().await.unwrap();
|
|
|
|
// Dropping the host Receiver aborts only its accept task; the Sender shares
|
|
// the same connection and must keep working.
|
|
drop(host_rx);
|
|
|
|
let tm = TypeMap::latest();
|
|
|
|
let resp = numbered_message(CommunicationType::Pong, 7, &tm);
|
|
host_tx.send(&resp).await.unwrap();
|
|
|
|
let got = client_rx.receive().await.unwrap();
|
|
assert_numbered_message(&got, CommunicationType::Pong, 7, &tm);
|
|
|
|
client_tx.close();
|
|
host_tx.close();
|
|
}
|