use macula_rust::cbor::Value;
use macula_rust::cert::ed25519_pubkey_from_cert;
use macula_rust::connection;
use macula_rust::identity::KeyPair;
use macula_rust::transport::{connect, Trust};
const STATION_HOST: &str = "station-de-frankfurt.macula.io";
const STATION_PORT: u16 = 4433;
const TORONTO_HOST: &str = "2600:3c04::2000:f0ff:feb9:e155";
const TORONTO_PORT: u16 = 4433;
const TORONTO_NODE_ID_HEX: &str =
"5748e81d89a6ea4b619fecda394ffac9f8f58a05d7a7234034783b6e1fd043d5";
const MILAN_HOST: &str = "station-it-milan.macula.io";
const MILAN_PORT: u16 = 4433;
#[tokio::test]
#[ignore = "requires network access to a live macula-station"]
async fn probe_what_frankfurt_presents() {
let connection = connect(STATION_HOST, STATION_PORT, Trust::Insecure)
.await
.expect("QUIC/TLS handshake with ALPN=macula should succeed against a live station");
println!(
"connected: alpn={:?} remote={}",
connection
.handshake_data()
.and_then(|d| d.downcast::<quinn::crypto::rustls::HandshakeData>().ok())
.and_then(|d| d.protocol)
.map(|p| String::from_utf8_lossy(&p).into_owned()),
connection.remote_address(),
);
let identity = connection
.peer_identity()
.expect("server cert chain should be present after a completed handshake");
let certs = identity
.downcast::<Vec<rustls::pki_types::CertificateDer<'static>>>()
.expect("peer_identity for a rustls-backed QUIC connection is a cert chain");
println!("station presented {} certificate(s)", certs.len());
let leaf = certs.first().expect("at least one cert in the chain");
match ed25519_pubkey_from_cert(leaf.as_ref()) {
Ok(pubkey) => println!("leaf is Ed25519, pubkey = {}", hex::encode(pubkey)),
Err(e) => println!("leaf is NOT a bare Ed25519 SPKI cert: {e}"),
}
connection.close(0u32.into(), b"probe complete");
}
#[tokio::test]
#[ignore = "requires network access to a live macula-station"]
async fn webpki_trust_succeeds_against_the_real_fleet() {
let connection = connect(STATION_HOST, STATION_PORT, Trust::WebPki)
.await
.expect("CA-chain validation should succeed against a real Let's Encrypt cert");
connection.close(0u32.into(), b"done");
}
#[tokio::test]
#[ignore = "requires network access to a live macula-station"]
async fn full_handshake_succeeds_against_the_real_fleet() {
let identity = KeyPair::generate_with_default_puzzle();
let session = connection::connect(STATION_HOST, STATION_PORT, Trust::WebPki, &identity)
.await
.expect("CONNECT/HELLO handshake should succeed against a live station");
println!(
"handshake accepted: remote={} station_node_id={} negotiated_capabilities={}",
session.remote_address(),
hex::encode(session.station.node_id),
session.station.negotiated_capabilities,
);
assert!(session.station.accepted);
session
.close("normal", Some("integration test done"), &identity)
.await;
}
#[tokio::test]
#[ignore = "requires network access to a live macula-station"]
async fn unhardened_identity_against_the_real_fleet_is_observed_not_assumed() {
let identity = KeyPair::generate();
let result = connection::connect(STATION_HOST, STATION_PORT, Trust::WebPki, &identity).await;
match result {
Ok(session) => {
println!(
"OBSERVED: unhardened identity was ACCEPTED (accepted={}) -- see this \
test's doc comment for why that's a fleet-configuration fact, not \
evidence this crate's puzzle handling is wrong",
session.station.accepted
);
session.close("normal", None, &identity).await;
}
Err(e) => {
println!("OBSERVED: unhardened identity was rejected, as: {e}");
}
}
}
#[tokio::test]
#[ignore = "requires network access to a live macula-station"]
async fn call_round_trip_against_the_real_fleet() {
let identity = KeyPair::generate_with_default_puzzle();
let mut session = connection::connect(STATION_HOST, STATION_PORT, Trust::WebPki, &identity)
.await
.expect("handshake should succeed");
let response = session
.call(
"macula_rust.test_probe",
[0u8; 32], macula_rust::cbor::Value::Null,
(now_ms() + 10_000) as i128,
&identity,
std::time::Duration::from_secs(10),
)
.await
.expect("should get SOME response (result or a well-formed error), not a timeout");
match response {
macula_rust::frame::CallResponse::Result {
payload,
responded_by,
} => {
println!("OBSERVED: got a RESULT (unexpected for a made-up procedure, but valid): payload={payload:?} responded_by={}", hex::encode(responded_by));
}
macula_rust::frame::CallResponse::Error {
code,
name,
reported_by,
detail,
} => {
println!(
"OBSERVED: got an ERROR (expected for a nonexistent procedure): code={code} name={name} reported_by={} detail={detail:?}",
hex::encode(reported_by)
);
}
}
session
.close("normal", Some("call test done"), &identity)
.await;
}
fn now_ms() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.expect("system clock after epoch")
.as_millis() as u64
}
#[tokio::test]
#[ignore = "requires network access to a live macula-station"]
async fn pubsub_round_trip_against_the_real_fleet() {
let identity = KeyPair::generate_with_default_puzzle();
let mut session = connection::connect(STATION_HOST, STATION_PORT, Trust::WebPki, &identity)
.await
.expect("handshake should succeed");
let realm: [u8; 32] = rand::random();
let topic = format!(
"macula-rust.test.{}",
hex::encode(rand::random::<[u8; 8]>())
);
session
.subscribe(
&macula_rust::frame::SubscribeSpec::new(topic.clone(), realm, identity.node_id()),
&identity,
)
.await
.expect("SUBSCRIBE should send without error");
session
.publish(
&macula_rust::frame::PublishSpec::new(
topic.clone(),
realm,
identity.node_id(),
1,
macula_rust::cbor::Value::text("hello from macula-rust"),
now_ms(),
),
&identity,
)
.await
.expect("PUBLISH should send without error");
match session.recv_event(std::time::Duration::from_secs(5)).await {
Ok(event) => {
println!(
"OBSERVED: received our own EVENT back — topic={} seq={} delivered_via={} payload={:?}",
event.topic, event.seq, event.delivered_via, event.payload
);
assert_eq!(event.topic, topic);
}
Err(e) => {
println!(
"OBSERVED: no EVENT arrived within 5s ({e}) — a subscriber may not receive its \
own publish, or delivery may simply be slower than this test waits. Not \
asserted as a failure either way; see this test's doc comment."
);
}
}
session
.close("normal", Some("pubsub test done"), &identity)
.await;
}
#[tokio::test]
#[ignore = "requires network access to a live macula-station"]
async fn run_subscriber_and_run_publisher_against_the_real_fleet() {
let realm: [u8; 32] = rand::random();
let topic = format!(
"macula-rust.test.{}",
hex::encode(rand::random::<[u8; 8]>())
);
let watcher_id = KeyPair::generate_with_default_puzzle();
let mut watcher = connection::connect(STATION_HOST, STATION_PORT, Trust::WebPki, &watcher_id)
.await
.expect("watcher handshake should succeed");
watcher
.subscribe(
&macula_rust::frame::SubscribeSpec::new(
"pubsub.publish_completed_v1",
realm,
watcher_id.node_id(),
),
&watcher_id,
)
.await
.expect("watcher SUBSCRIBE should send without error");
let sub_id = KeyPair::generate_with_default_puzzle();
let mut sub_session = connection::connect(STATION_HOST, STATION_PORT, Trust::WebPki, &sub_id)
.await
.expect("subscriber handshake should succeed");
let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
let sub_topic = topic.clone();
let subscribe_task = tokio::spawn(async move {
let spec = macula_rust::frame::SubscribeSpec::new(sub_topic, realm, sub_id.node_id());
let stop = tokio::time::sleep(std::time::Duration::from_secs(8));
sub_session
.run_subscriber(&spec, &sub_id, stop, |evt| {
let _ = tx.send(evt);
})
.await
});
tokio::time::sleep(std::time::Duration::from_millis(500)).await;
let pub_id = KeyPair::generate_with_default_puzzle();
let mut pub_session = connection::connect(STATION_HOST, STATION_PORT, Trust::WebPki, &pub_id)
.await
.expect("publisher handshake should succeed");
let spec = macula_rust::frame::PublishSpec::new(
topic.clone(),
realm,
pub_id.node_id(),
1,
Value::text("hello from run_publisher"),
now_ms(),
);
let publish_result = pub_session.run_publisher(&spec, &pub_id, true).await;
assert!(
publish_result.is_ok(),
"run_publisher should succeed: {publish_result:?}"
);
println!("OBSERVED: run_publisher completed cleanly");
pub_session
.close("normal", Some("publisher test done"), &pub_id)
.await;
match tokio::time::timeout(std::time::Duration::from_secs(5), rx.recv()).await {
Ok(Some(evt)) => {
println!(
"OBSERVED: run_subscriber's handler received the real EVENT -- topic={} payload={:?}",
evt.topic, evt.payload
);
assert_eq!(evt.topic, topic);
}
_ => println!(
"OBSERVED: no EVENT arrived via the subscriber's callback within 5s -- a subscriber \
may not receive its own publish, same caveat as pubsub_round_trip_against_the_real_fleet"
),
}
let sub_result = subscribe_task
.await
.expect("subscriber task should not panic");
assert!(
sub_result.is_ok(),
"run_subscriber should return Ok after its stop future resolves: {sub_result:?}"
);
println!("OBSERVED: run_subscriber returned cleanly after its stop future resolved");
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(10);
let mut confirmed = false;
while std::time::Instant::now() < deadline {
match watcher.recv_event(std::time::Duration::from_secs(2)).await {
Ok(evt) if evt.topic == "pubsub.publish_completed_v1" => {
let outcome = evt.payload.get("outcome");
println!(
"OBSERVED: independent watcher confirmed a real pubsub.publish_completed_v1 fact landed -- outcome={outcome:?}"
);
assert_eq!(outcome, Some(&Value::text("completed")));
confirmed = true;
break;
}
Ok(other) => {
println!(
"(watcher skipping unrelated event on topic {})",
other.topic
);
}
Err(e) => {
println!("(watcher skipping a non-event frame or timeout: {e})");
}
}
}
assert!(
confirmed,
"independent watcher never observed a pubsub.publish_completed_v1 fact within the deadline"
);
watcher
.close("normal", Some("watcher test done"), &watcher_id)
.await;
}
#[tokio::test]
#[ignore = "requires network access to a live macula-station"]
async fn publish_survives_immediate_close_against_the_real_fleet() {
let sub_identity = KeyPair::generate_with_default_puzzle();
let mut sub_session =
connection::connect(STATION_HOST, STATION_PORT, Trust::WebPki, &sub_identity)
.await
.expect("handshake should succeed (subscriber)");
let realm: [u8; 32] = rand::random();
let topic = format!(
"macula-rust.test.immediate-close.{}",
hex::encode(rand::random::<[u8; 8]>())
);
sub_session
.subscribe(
&macula_rust::frame::SubscribeSpec::new(topic.clone(), realm, sub_identity.node_id()),
&sub_identity,
)
.await
.expect("SUBSCRIBE should send without error");
tokio::time::sleep(std::time::Duration::from_millis(500)).await;
{
let pub_identity = KeyPair::generate_with_default_puzzle();
let mut pub_session =
connection::connect(STATION_HOST, STATION_PORT, Trust::WebPki, &pub_identity)
.await
.expect("handshake should succeed (publisher)");
pub_session
.publish(
&macula_rust::frame::PublishSpec::new(
topic.clone(),
realm,
pub_identity.node_id(),
1,
macula_rust::cbor::Value::text(
"hello from the immediate-close regression test",
),
now_ms(),
),
&pub_identity,
)
.await
.expect("PUBLISH should send without error");
pub_session.close("normal", None, &pub_identity).await;
}
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
loop {
let remaining = deadline.saturating_duration_since(std::time::Instant::now());
assert!(
!remaining.is_zero(),
"EVENT for our topic never arrived after a publish immediately followed by close \
(this is the exact race this test exists to catch)"
);
match sub_session.recv_event(remaining).await {
Ok(event) if event.topic == topic => break,
Ok(event) => println!("skipping unrelated live EVENT: topic={}", event.topic),
Err(connection::RecvEventError::Recv(connection::RecvFrameError::Timeout)) => {
panic!(
"EVENT for our topic never arrived after a publish immediately followed by \
close (this is the exact race this test exists to catch)"
);
}
Err(e) => println!("skipping a non-EVENT frame on this shared, busy station: {e}"),
}
}
sub_session
.close("normal", Some("immediate-close test done"), &sub_identity)
.await;
}
#[tokio::test]
#[ignore = "requires network access to a live macula-station"]
async fn single_block_put_get_round_trip_against_the_real_fleet() {
let identity = KeyPair::generate_with_default_puzzle();
let mut session = connection::connect(STATION_HOST, STATION_PORT, Trust::WebPki, &identity)
.await
.expect("handshake should succeed");
let data: Vec<u8> = (0..4096).map(|_| rand::random::<u8>()).collect();
let mcid = macula_rust::content::put(&mut session, &data, "test-block", &identity)
.await
.expect("put should succeed");
assert!(
!macula_rust::manifest::mcid_is_chunked(&mcid),
"4096 bytes is well under the chunking threshold"
);
println!(
"OBSERVED: stored single block under mcid={}",
hex::encode(mcid)
);
let fetched = macula_rust::content::get(&mut session, mcid, &identity)
.await
.expect("get should succeed for content this session just put");
assert_eq!(
fetched, data,
"fetched bytes must match what was put, exactly"
);
session
.close("normal", Some("content single-block test done"), &identity)
.await;
}
#[tokio::test]
#[ignore = "requires network access to a live macula-station"]
async fn chunked_put_get_round_trip_against_the_real_fleet() {
let identity = KeyPair::generate_with_default_puzzle();
let mut session = connection::connect(STATION_HOST, STATION_PORT, Trust::WebPki, &identity)
.await
.expect("handshake should succeed");
let size = macula_rust::manifest::DEFAULT_CHUNK_SIZE * 2 + 12_345;
let data: Vec<u8> = (0..size).map(|_| rand::random::<u8>()).collect();
let mcid = macula_rust::content::put(&mut session, &data, "test-chunked", &identity)
.await
.expect("chunked put should succeed");
assert!(
macula_rust::manifest::mcid_is_chunked(&mcid),
"{size} bytes is well over the chunking threshold"
);
println!(
"OBSERVED: stored {size} bytes as a manifest under mcid={}",
hex::encode(mcid)
);
let fetched = macula_rust::content::get(&mut session, mcid, &identity)
.await
.expect("chunked get should succeed for content this session just put");
assert_eq!(
fetched, data,
"reassembled bytes must match what was put, exactly"
);
session
.close("normal", Some("content chunked test done"), &identity)
.await;
}
#[tokio::test]
#[ignore = "requires network access to a live macula-station"]
async fn get_of_an_unknown_block_reports_not_found_against_the_real_fleet() {
let identity = KeyPair::generate_with_default_puzzle();
let mut session = connection::connect(STATION_HOST, STATION_PORT, Trust::WebPki, &identity)
.await
.expect("handshake should succeed");
let random_hash: [u8; 32] = rand::random();
let mcid = macula_rust::manifest::block_mcid(&random_hash);
match macula_rust::content::get(&mut session, mcid, &identity).await {
Err(macula_rust::content::GetError::NotFound) => {
println!("OBSERVED: not_found reported correctly for an unknown mcid");
}
other => panic!("expected GetError::NotFound, got {other:?}"),
}
session
.close("normal", Some("content not-found test done"), &identity)
.await;
}
#[tokio::test]
#[ignore = "requires network access to a live macula-station"]
async fn stream_open_round_trip_against_the_real_fleet() {
let identity = KeyPair::generate_with_default_puzzle();
let mut session = connection::connect(STATION_HOST, STATION_PORT, Trust::WebPki, &identity)
.await
.expect("handshake should succeed");
let mut handle = macula_rust::stream::StreamHandle::open(
&mut session,
"macula_rust.test_stream",
[0u8; 32],
macula_rust::frame::StreamMode::ClientStream,
macula_rust::cbor::Value::Null,
(now_ms() + 10_000) as i128,
&identity,
)
.await
.expect("opening a dedicated stream and sending STREAM_OPEN should succeed");
handle
.send_data(
macula_rust::frame::StreamEncoding::Raw,
macula_rust::cbor::Value::Bytes(b"hello from macula-rust".to_vec()),
&identity,
)
.await
.expect("sending a chunk should succeed");
handle
.close_send(&identity)
.await
.expect("half-closing should succeed");
match handle.await_reply(std::time::Duration::from_secs(5)).await {
Ok((payload, responded_by)) => {
println!(
"OBSERVED: got a STREAM_REPLY (unexpected for a made-up procedure, but valid): payload={payload:?} responded_by={}",
hex::encode(responded_by)
);
}
Err(e) => {
println!("OBSERVED: no reply within 5s, as: {e}");
}
}
session
.close("normal", Some("stream test done"), &identity)
.await;
}
#[tokio::test]
#[ignore = "requires network access to a live macula-station"]
async fn streaming_provider_round_trip_against_the_real_fleet() {
let provider_identity = KeyPair::generate_with_default_puzzle();
let caller_identity = KeyPair::generate_with_default_puzzle();
let mut provider_session = connection::connect(
STATION_HOST,
STATION_PORT,
Trust::WebPki,
&provider_identity,
)
.await
.expect("provider handshake should succeed");
let mut caller_session =
connection::connect(STATION_HOST, STATION_PORT, Trust::WebPki, &caller_identity)
.await
.expect("caller handshake should succeed");
let realm: [u8; 32] = rand::random();
let procedure = format!(
"macula_rust.test_provider.{}",
hex::encode(rand::random::<[u8; 8]>())
);
let advertise_spec = macula_rust::frame::AdvertiseSpec::new(
realm,
procedure.clone(),
provider_identity.node_id(),
);
provider_session
.advertise(&advertise_spec, &provider_identity)
.await
.expect("advertise should send");
tokio::time::sleep(std::time::Duration::from_millis(500)).await;
let accept_task = tokio::spawn(async move {
let result = macula_rust::stream::StreamHandle::accept(
&mut provider_session,
std::time::Duration::from_secs(10),
)
.await;
(result, provider_session)
});
let mut caller_handle = macula_rust::stream::StreamHandle::open(
&mut caller_session,
&procedure,
realm,
macula_rust::frame::StreamMode::ServerStream,
macula_rust::cbor::Value::Null,
(now_ms() + 10_000) as i128,
&caller_identity,
)
.await
.expect("caller should open a stream");
let (accept_result, provider_session) =
accept_task.await.expect("accept task should not panic");
let (mut provider_handle, open_info) =
accept_result.expect("provider should accept the inbound STREAM_OPEN");
println!(
"OBSERVED: provider accepted stream_open for procedure={} mode={:?}",
open_info.procedure, open_info.mode
);
assert_eq!(open_info.procedure, procedure);
assert_eq!(open_info.mode, macula_rust::frame::StreamMode::ServerStream);
provider_handle
.send_data(
macula_rust::frame::StreamEncoding::Raw,
macula_rust::cbor::Value::Bytes(b"hello from the provider".to_vec()),
&provider_identity,
)
.await
.expect("provider should push a chunk");
provider_handle
.close_send(&provider_identity)
.await
.expect("provider should close its send side");
match caller_handle
.recv(std::time::Duration::from_secs(5))
.await
.expect("caller should receive the pushed chunk")
{
macula_rust::stream::StreamItem::Data { body, .. } => {
assert_eq!(
body,
macula_rust::cbor::Value::Bytes(b"hello from the provider".to_vec())
);
}
other => panic!("expected Data, got {other:?}"),
}
match caller_handle
.recv(std::time::Duration::from_secs(5))
.await
.expect("caller should see end-of-stream")
{
macula_rust::stream::StreamItem::Eof => {}
other => panic!("expected Eof, got {other:?}"),
}
provider_session
.close("normal", Some("provider test done"), &provider_identity)
.await;
caller_session
.close("normal", Some("caller test done"), &caller_identity)
.await;
}
#[tokio::test]
#[ignore = "requires network access to a live macula-station"]
async fn unary_call_provider_round_trip_against_the_real_fleet() {
let provider_identity = KeyPair::generate_with_default_puzzle();
let caller_identity = KeyPair::generate_with_default_puzzle();
let mut provider_session = connection::connect(
STATION_HOST,
STATION_PORT,
Trust::WebPki,
&provider_identity,
)
.await
.expect("provider handshake should succeed");
let mut caller_session =
connection::connect(STATION_HOST, STATION_PORT, Trust::WebPki, &caller_identity)
.await
.expect("caller handshake should succeed");
let realm: [u8; 32] = rand::random();
let procedure = format!(
"macula_rust.test_add.{}",
hex::encode(rand::random::<[u8; 8]>())
);
let advertise_spec = macula_rust::frame::AdvertiseSpec::new(
realm,
procedure.clone(),
provider_identity.node_id(),
);
provider_session
.advertise(&advertise_spec, &provider_identity)
.await
.expect("advertise should send");
tokio::time::sleep(std::time::Duration::from_millis(500)).await;
let target_procedure = procedure.clone();
let lookup = move |_realm: &[u8; 32], proc: &str| -> Option<connection::CallHandler> {
if proc != target_procedure {
return None;
}
let handler: connection::CallHandler = std::sync::Arc::new(|payload: Value| {
Box::pin(async move {
let a = match payload.get("a") {
Some(Value::Int(n)) => *n,
_ => return Err("missing or non-integer field \"a\"".to_string()),
};
let b = match payload.get("b") {
Some(Value::Int(n)) => *n,
_ => return Err("missing or non-integer field \"b\"".to_string()),
};
Ok(Value::Int(a + b))
}) as connection::BoxFuture<'static, Result<Value, String>>
});
Some(handler)
};
let serve_task = tokio::spawn(async move {
let result = provider_session
.serve_one_call(
lookup,
&provider_identity,
std::time::Duration::from_secs(15),
)
.await;
(result, provider_session, provider_identity)
});
let payload = Value::Map(vec![
(Value::text("a"), Value::Int(3)),
(Value::text("b"), Value::Int(4)),
]);
let response = caller_session
.call(
&procedure,
realm,
payload,
(now_ms() + 10_000) as i128,
&caller_identity,
std::time::Duration::from_secs(10),
)
.await
.expect("call should succeed");
let (serve_result, provider_session, provider_identity) =
serve_task.await.expect("serve task should not panic");
serve_result.expect("provider should serve the inbound CALL");
match response {
macula_rust::frame::CallResponse::Result { payload, .. } => {
assert_eq!(payload, Value::Int(7), "3 + 4 should reply with RESULT 7");
}
other => panic!("expected a RESULT, got {other:?}"),
}
println!(
"OBSERVED: provider served the inbound CALL for procedure={procedure}, caller got RESULT 7"
);
provider_session
.close(
"normal",
Some("unary provider test done"),
&provider_identity,
)
.await;
caller_session
.close("normal", Some("unary caller test done"), &caller_identity)
.await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
#[ignore = "requires network access to a live macula-station"]
async fn unary_call_provider_round_trip_multi_thread_runtime() {
let provider_identity = KeyPair::generate_with_default_puzzle();
let caller_identity = KeyPair::generate_with_default_puzzle();
let mut provider_session = connection::connect(
STATION_HOST,
STATION_PORT,
Trust::WebPki,
&provider_identity,
)
.await
.expect("provider handshake should succeed");
let mut caller_session =
connection::connect(STATION_HOST, STATION_PORT, Trust::WebPki, &caller_identity)
.await
.expect("caller handshake should succeed");
let realm: [u8; 32] = rand::random();
let procedure = format!(
"macula_rust.test_multithread.{}",
hex::encode(rand::random::<[u8; 8]>())
);
let advertise_spec = macula_rust::frame::AdvertiseSpec::new(
realm,
procedure.clone(),
provider_identity.node_id(),
);
provider_session
.advertise(&advertise_spec, &provider_identity)
.await
.expect("advertise should send");
tokio::time::sleep(std::time::Duration::from_millis(500)).await;
let target_procedure = procedure.clone();
let lookup = move |_realm: &[u8; 32], proc: &str| -> Option<connection::CallHandler> {
if proc != target_procedure {
return None;
}
let handler: connection::CallHandler = std::sync::Arc::new(|payload: Value| {
Box::pin(async move { Ok(payload) })
as connection::BoxFuture<'static, Result<Value, String>>
});
Some(handler)
};
let serve_task = tokio::spawn(async move {
let result = provider_session
.serve_one_call(
lookup,
&provider_identity,
std::time::Duration::from_secs(15),
)
.await;
provider_session
.close(
"normal",
Some("multi-thread provider test done"),
&provider_identity,
)
.await;
result
});
let response = caller_session
.call(
&procedure,
realm,
Value::text("hello"),
(now_ms() + 10_000) as i128,
&caller_identity,
std::time::Duration::from_secs(10),
)
.await
.expect("call should succeed");
let serve_result = serve_task.await.expect("serve task should not panic");
serve_result.expect("provider should serve the inbound CALL");
match response {
macula_rust::frame::CallResponse::Result { payload, .. } => {
assert_eq!(payload, Value::text("hello"));
}
other => panic!("expected a RESULT, got {other:?}"),
}
caller_session
.close(
"normal",
Some("multi-thread caller test done"),
&caller_identity,
)
.await;
}
#[tokio::test]
#[ignore = "requires network access to a live macula-station"]
async fn rpc_telemetry_facts_against_the_real_fleet() {
let provider_identity = KeyPair::generate_with_default_puzzle();
let caller_identity = KeyPair::generate_with_default_puzzle();
let watcher_identity = KeyPair::generate_with_default_puzzle();
let realm: [u8; 32] = rand::random();
let procedure = format!(
"macula_rust.test_rpc_facts.{}",
hex::encode(rand::random::<[u8; 8]>())
);
let mut watcher =
connection::connect(STATION_HOST, STATION_PORT, Trust::WebPki, &watcher_identity)
.await
.expect("watcher handshake should succeed");
for topic in [
"rpc.sent_v1",
"rpc.completed_v1",
"rpc.received_v1",
"rpc.replied_v1",
] {
watcher
.subscribe(
&macula_rust::frame::SubscribeSpec::new(topic, realm, watcher_identity.node_id()),
&watcher_identity,
)
.await
.expect("watcher SUBSCRIBE should send without error");
}
let mut provider_session = connection::connect(
STATION_HOST,
STATION_PORT,
Trust::WebPki,
&provider_identity,
)
.await
.expect("provider handshake should succeed");
let mut caller_session =
connection::connect(STATION_HOST, STATION_PORT, Trust::WebPki, &caller_identity)
.await
.expect("caller handshake should succeed");
let advertise_spec = macula_rust::frame::AdvertiseSpec::new(
realm,
procedure.clone(),
provider_identity.node_id(),
);
provider_session
.advertise(&advertise_spec, &provider_identity)
.await
.expect("advertise should send");
tokio::time::sleep(std::time::Duration::from_millis(500)).await;
let target_procedure = procedure.clone();
let lookup = move |_realm: &[u8; 32], proc: &str| -> Option<connection::CallHandler> {
if proc != target_procedure {
return None;
}
Some(std::sync::Arc::new(|payload: Value| {
Box::pin(async move { Ok(payload) })
as connection::BoxFuture<'static, Result<Value, String>>
}))
};
let serve_task = tokio::spawn(async move {
let result = provider_session
.serve_one_call(
lookup,
&provider_identity,
std::time::Duration::from_secs(15),
)
.await;
(result, provider_session, provider_identity)
});
let response = caller_session
.call(
&procedure,
realm,
Value::text("hello"),
(now_ms() + 10_000) as i128,
&caller_identity,
std::time::Duration::from_secs(10),
)
.await
.expect("call should succeed");
assert!(
matches!(response, macula_rust::frame::CallResponse::Result { .. }),
"expected a RESULT, got {response:?}"
);
let (serve_result, provider_session, provider_identity) =
serve_task.await.expect("serve task should not panic");
serve_result.expect("provider should serve the inbound CALL");
provider_session
.close("normal", Some("rpc facts test done"), &provider_identity)
.await;
let mut seen = std::collections::HashSet::new();
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
while seen.len() < 4 && std::time::Instant::now() < deadline {
let Ok(value) =
tokio::time::timeout(std::time::Duration::from_secs(1), watcher.recv_frame()).await
else {
continue;
};
let Ok(evt) = macula_rust::frame::parse_event(&value.expect("recv_frame should not error"))
else {
continue;
};
if seen.insert(evt.topic.clone()) {
println!(
"OBSERVED: {} fact landed with payload={:?}",
evt.topic, evt.payload
);
}
}
assert_eq!(
seen,
std::collections::HashSet::from([
"rpc.sent_v1".to_string(),
"rpc.completed_v1".to_string(),
"rpc.received_v1".to_string(),
"rpc.replied_v1".to_string(),
]),
"expected all 4 RPC telemetry facts to land, only saw: {seen:?}"
);
caller_session
.close("normal", Some("rpc facts test done"), &caller_identity)
.await;
watcher
.close("normal", Some("rpc facts test done"), &watcher_identity)
.await;
}
#[tokio::test]
#[ignore = "requires network access to a live macula-station"]
async fn unary_call_provider_reports_unknown_next_peer_on_lookup_miss_against_the_real_fleet() {
let provider_identity = KeyPair::generate_with_default_puzzle();
let caller_identity = KeyPair::generate_with_default_puzzle();
let mut provider_session = connection::connect(
STATION_HOST,
STATION_PORT,
Trust::WebPki,
&provider_identity,
)
.await
.expect("provider handshake should succeed");
let mut caller_session =
connection::connect(STATION_HOST, STATION_PORT, Trust::WebPki, &caller_identity)
.await
.expect("caller handshake should succeed");
let realm: [u8; 32] = rand::random();
let procedure = format!(
"macula_rust.test_miss.{}",
hex::encode(rand::random::<[u8; 8]>())
);
let advertise_spec = macula_rust::frame::AdvertiseSpec::new(
realm,
procedure.clone(),
provider_identity.node_id(),
);
provider_session
.advertise(&advertise_spec, &provider_identity)
.await
.expect("advertise should send");
tokio::time::sleep(std::time::Duration::from_millis(500)).await;
let no_handlers = |_realm: &[u8; 32], _proc: &str| -> Option<connection::CallHandler> { None };
let serve_task = tokio::spawn(async move {
let result = provider_session
.serve_one_call(
no_handlers,
&provider_identity,
std::time::Duration::from_secs(15),
)
.await;
(result, provider_session, provider_identity)
});
let response = caller_session
.call(
&procedure,
realm,
Value::Null,
(now_ms() + 10_000) as i128,
&caller_identity,
std::time::Duration::from_secs(10),
)
.await
.expect("call should succeed");
let (serve_result, provider_session, provider_identity) =
serve_task.await.expect("serve task should not panic");
serve_result.expect("provider should serve the inbound CALL (with an error reply)");
match response {
macula_rust::frame::CallResponse::Error { code, name, .. } => {
assert_eq!(code, macula_rust::bolt4::Code::UnknownNextPeer.as_u8());
println!("OBSERVED: lookup miss correctly reported as ERROR code={code} name={name}");
}
other => panic!("expected an ERROR, got {other:?}"),
}
provider_session
.close(
"normal",
Some("unary provider miss test done"),
&provider_identity,
)
.await;
caller_session
.close(
"normal",
Some("unary caller miss test done"),
&caller_identity,
)
.await;
}
#[tokio::test]
#[ignore = "requires network access to a live macula-station"]
async fn pinned_trust_full_handshake_succeeds_against_toronto() {
let node_id: [u8; 32] = hex::decode(TORONTO_NODE_ID_HEX)
.expect("valid hex")
.try_into()
.expect("32 bytes");
let identity = KeyPair::generate_with_default_puzzle();
let session = connection::connect(
TORONTO_HOST,
TORONTO_PORT,
Trust::Pinned(node_id),
&identity,
)
.await
.expect("Pinned-trust CONNECT/HELLO handshake should succeed against a live no-DNS station");
println!(
"handshake accepted: remote={} station_node_id={} negotiated_capabilities={}",
session.remote_address(),
hex::encode(session.station.node_id),
session.station.negotiated_capabilities,
);
assert!(session.station.accepted);
assert_eq!(
session.station.node_id, node_id,
"the station's own reported node_id should match the one we pinned"
);
session
.close("normal", Some("pinned trust test done"), &identity)
.await;
}
#[tokio::test]
#[ignore = "requires network access to a live macula-station"]
async fn cross_station_streaming_round_trip_frankfurt_provider_milan_caller() {
let provider_identity = KeyPair::generate_with_default_puzzle();
let caller_identity = KeyPair::generate_with_default_puzzle();
let mut provider_session = connection::connect(
STATION_HOST,
STATION_PORT,
Trust::WebPki,
&provider_identity,
)
.await
.expect("provider handshake against Frankfurt should succeed");
let mut caller_session =
connection::connect(MILAN_HOST, MILAN_PORT, Trust::WebPki, &caller_identity)
.await
.expect("caller handshake against Milan should succeed");
let realm: [u8; 32] = rand::random();
let procedure = format!(
"macula_rust.test_call.{}",
hex::encode(rand::random::<[u8; 8]>())
);
let advertise_spec = macula_rust::frame::AdvertiseSpec::new(
realm,
procedure.clone(),
provider_identity.node_id(),
);
provider_session
.advertise(&advertise_spec, &provider_identity)
.await
.expect("advertise on Frankfurt should send");
tokio::time::sleep(std::time::Duration::from_secs(5)).await;
let accept_task = tokio::spawn(async move {
let result = macula_rust::stream::StreamHandle::accept(
&mut provider_session,
std::time::Duration::from_secs(15),
)
.await;
(result, provider_session)
});
let open_result = macula_rust::stream::StreamHandle::open(
&mut caller_session,
&procedure,
realm,
macula_rust::frame::StreamMode::Bidi,
macula_rust::cbor::Value::Null,
(now_ms() + 10_000) as i128,
&caller_identity,
)
.await;
let mut caller_handle = match open_result {
Ok(h) => h,
Err(e) => {
println!(
"OBSERVED: cross-station STREAM_OPEN failed as: {e} -- Milan could not \
route to a procedure only advertised on Frankfurt within 5s. This is the \
real, useful answer to whether a call feature can rely on cross-station \
routing working promptly; see this test's doc comment."
);
let _ = accept_task.await;
return;
}
};
println!("OBSERVED: cross-station STREAM_OPEN succeeded -- Milan routed it to Frankfurt");
let (accept_result, provider_session) =
accept_task.await.expect("accept task should not panic");
let (mut provider_handle, open_info) =
accept_result.expect("provider should accept the inbound STREAM_OPEN");
assert_eq!(open_info.procedure, procedure);
assert_eq!(open_info.mode, macula_rust::frame::StreamMode::Bidi);
caller_handle
.send_data(
macula_rust::frame::StreamEncoding::Raw,
macula_rust::cbor::Value::Bytes(b"audio frame from phone2 (milan)".to_vec()),
&caller_identity,
)
.await
.expect("caller should push a frame");
provider_handle
.send_data(
macula_rust::frame::StreamEncoding::Raw,
macula_rust::cbor::Value::Bytes(b"audio frame from phone1 (frankfurt)".to_vec()),
&provider_identity,
)
.await
.expect("provider should push a frame");
match provider_handle
.recv(std::time::Duration::from_secs(5))
.await
.expect(
"provider should receive the caller's frame -- see this test's doc comment, \
fixed 2026-08-29 by stamping `signer` on stream data frames",
) {
macula_rust::stream::StreamItem::Data { body, .. } => {
assert_eq!(
body,
macula_rust::cbor::Value::Bytes(b"audio frame from phone2 (milan)".to_vec())
);
println!("OBSERVED: provider (Frankfurt) received phone2's frame from Milan");
}
other => panic!("expected Data, got {other:?}"),
}
match caller_handle
.recv(std::time::Duration::from_secs(5))
.await
.expect("caller should receive the provider's frame")
{
macula_rust::stream::StreamItem::Data { body, .. } => {
assert_eq!(
body,
macula_rust::cbor::Value::Bytes(b"audio frame from phone1 (frankfurt)".to_vec())
);
println!("OBSERVED: caller (Milan) received phone1's frame from Frankfurt");
}
other => panic!("expected Data, got {other:?}"),
}
caller_handle
.close_send(&caller_identity)
.await
.expect("caller should half-close");
provider_handle
.close_send(&provider_identity)
.await
.expect("provider should half-close");
provider_session
.close(
"normal",
Some("cross-station call test done"),
&provider_identity,
)
.await;
caller_session
.close(
"normal",
Some("cross-station call test done"),
&caller_identity,
)
.await;
}
#[tokio::test]
#[ignore = "requires network access to a live macula-station"]
async fn cross_station_unary_call_round_trip_frankfurt_provider_milan_caller() {
let provider_identity = KeyPair::generate_with_default_puzzle();
let caller_identity = KeyPair::generate_with_default_puzzle();
let mut provider_session = connection::connect(
STATION_HOST,
STATION_PORT,
Trust::WebPki,
&provider_identity,
)
.await
.expect("provider handshake against Frankfurt should succeed");
let mut caller_session =
connection::connect(MILAN_HOST, MILAN_PORT, Trust::WebPki, &caller_identity)
.await
.expect("caller handshake against Milan should succeed");
let realm: [u8; 32] = rand::random();
let procedure = format!(
"macula_rust.test_signal.{}",
hex::encode(rand::random::<[u8; 8]>())
);
let advertise_spec = macula_rust::frame::AdvertiseSpec::new(
realm,
procedure.clone(),
provider_identity.node_id(),
);
provider_session
.advertise(&advertise_spec, &provider_identity)
.await
.expect("advertise on Frankfurt should send");
tokio::time::sleep(std::time::Duration::from_secs(5)).await;
let target_procedure = procedure.clone();
let lookup = move |_realm: &[u8; 32], proc: &str| -> Option<connection::CallHandler> {
if proc != target_procedure {
return None;
}
let handler: connection::CallHandler = std::sync::Arc::new(|payload: Value| {
Box::pin(async move {
match payload {
Value::Text(s) if s == "offer from phone2 (milan)" => {
Ok(Value::text("answer from phone1 (frankfurt)"))
}
other => Err(format!("unexpected payload: {other:?}")),
}
}) as connection::BoxFuture<'static, Result<Value, String>>
});
Some(handler)
};
let serve_task = tokio::spawn(async move {
let result = provider_session
.serve_one_call(
lookup,
&provider_identity,
std::time::Duration::from_secs(20),
)
.await;
(result, provider_session, provider_identity)
});
let response = caller_session
.call(
&procedure,
realm,
Value::text("offer from phone2 (milan)"),
(now_ms() + 15_000) as i128,
&caller_identity,
std::time::Duration::from_secs(15),
)
.await;
let (serve_result, provider_session, provider_identity) =
serve_task.await.expect("serve task should not panic");
match (response, serve_result) {
(Ok(macula_rust::frame::CallResponse::Result { payload, .. }), Ok(())) => {
let matches = payload == Value::text("answer from phone1 (frankfurt)");
println!(
"OBSERVED: cross-station CALL/RESULT succeeded -- Milan's CALL reached \
Frankfurt's provider and the RESULT came back, content matches = {matches}"
);
}
(Ok(other), serve_result) => {
println!(
"OBSERVED: cross-station CALL got a response but not the expected RESULT: \
{other:?} (serve_result={serve_result:?})"
);
}
(Err(e), serve_result) => {
println!(
"OBSERVED: cross-station CALL failed -- {e} (serve_result={serve_result:?}). \
If this fails the same way the streaming test's DATA phase did, the CALL \
path shares the same cross-station gap; if it succeeds, signaling built on \
CALL rather than STREAM_OPEN+DATA is on solid ground."
);
}
}
provider_session
.close(
"normal",
Some("cross-station signaling test done"),
&provider_identity,
)
.await;
caller_session
.close(
"normal",
Some("cross-station signaling test done"),
&caller_identity,
)
.await;
}
#[tokio::test]
#[ignore = "requires network access to a live macula-station"]
async fn direct_dial_advertise_resolve_and_call_round_trip_against_the_real_fleet() {
let provider_identity = KeyPair::generate_with_default_puzzle();
let caller_identity = KeyPair::generate_with_default_puzzle();
let mut provider_session =
connection::connect(MILAN_HOST, MILAN_PORT, Trust::WebPki, &provider_identity)
.await
.expect("provider handshake should succeed");
let mut resolve_session =
connection::connect(STATION_HOST, STATION_PORT, Trust::WebPki, &caller_identity)
.await
.expect("resolve-side handshake should succeed");
let realm: [u8; 32] = rand::random();
let procedure = format!(
"macula_rust.test_direct_dial.{}",
hex::encode(rand::random::<[u8; 8]>())
);
macula_rust::direct_dial::advertise_direct(
&mut provider_session,
&provider_identity,
realm,
&procedure,
std::time::Duration::from_secs(120),
)
.await
.expect("advertise_direct should publish the DHT record");
let target_procedure = procedure.clone();
let lookup = move |_realm: &[u8; 32], proc: &str| -> Option<connection::CallHandler> {
if proc != target_procedure {
return None;
}
let handler: connection::CallHandler = std::sync::Arc::new(|payload: Value| {
Box::pin(async move {
let n = match payload.get("n") {
Some(Value::Int(n)) => *n,
_ => return Err("missing or non-integer field \"n\"".to_string()),
};
Ok(Value::Int(n * 2))
}) as connection::BoxFuture<'static, Result<Value, String>>
});
Some(handler)
};
let serve_task = tokio::spawn(async move {
let result = provider_session
.serve_one_call(
lookup,
&provider_identity,
std::time::Duration::from_secs(20),
)
.await;
(result, provider_session, provider_identity)
});
let payload = Value::Map(vec![(Value::text("n"), Value::Int(21))]);
let response = macula_rust::direct_dial::call(
&mut resolve_session,
&caller_identity,
realm,
&procedure,
payload,
std::time::Duration::from_secs(15),
)
.await
.expect("direct-dial call should resolve, dial, and complete");
let (serve_result, provider_session, provider_identity) =
serve_task.await.expect("serve task should not panic");
serve_result.expect("provider should serve the direct-dialed inbound CALL");
match response {
macula_rust::frame::CallResponse::Result { payload, .. } => {
assert_eq!(
payload,
Value::Int(42),
"21 * 2 should reply with RESULT 42"
);
}
other => panic!("expected a RESULT, got {other:?}"),
}
println!(
"OBSERVED: direct-dial resolved+dialed a station DIFFERENT from the resolve session's own, and got a real RESULT for procedure={procedure}"
);
provider_session
.close(
"normal",
Some("direct-dial provider test done"),
&provider_identity,
)
.await;
resolve_session
.close(
"normal",
Some("direct-dial resolve-side test done"),
&caller_identity,
)
.await;
}
#[tokio::test]
#[ignore = "requires network access to a live macula-station"]
async fn keep_advertised_direct_republishes_against_the_real_fleet() {
let publisher_identity = KeyPair::generate_with_default_puzzle();
let reader_identity = KeyPair::generate_with_default_puzzle();
let mut reader_session =
connection::connect(STATION_HOST, STATION_PORT, Trust::WebPki, &reader_identity)
.await
.expect("reader handshake should succeed");
let mut loop_session = connection::connect(
STATION_HOST,
STATION_PORT,
Trust::WebPki,
&publisher_identity,
)
.await
.expect("loop session handshake should succeed");
let realm: [u8; 32] = rand::random();
let procedure = format!(
"macula_rust.test_keep_advertised.{}",
hex::encode(rand::random::<[u8; 8]>())
);
let uri = macula_rust::dht::discovery_uri(realm, &procedure);
let key = macula_rust::dht::procedure_key(&uri);
let (stop_tx, stop_rx) = tokio::sync::oneshot::channel::<()>();
let loop_procedure = procedure.clone();
let loop_task = tokio::spawn(async move {
macula_rust::direct_dial::keep_advertised_direct(
&mut loop_session,
&publisher_identity,
realm,
&loop_procedure,
std::time::Duration::from_secs(120),
std::time::Duration::from_millis(500),
async move {
let _ = stop_rx.await;
},
|e| eprintln!("keep_advertised_direct tick failed (non-fatal): {e}"),
)
.await;
(loop_session, publisher_identity)
});
tokio::time::sleep(std::time::Duration::from_millis(300)).await;
let first = macula_rust::dht::find_record(&mut reader_session, &reader_identity, key)
.await
.expect("first tick should already be visible");
tokio::time::sleep(std::time::Duration::from_millis(700)).await;
let second = macula_rust::dht::find_record(&mut reader_session, &reader_identity, key)
.await
.expect("second tick should be visible");
assert!(
second.created_at > first.created_at,
"expected the second tick's created_at ({}) to be strictly after the first's ({})",
second.created_at,
first.created_at
);
println!(
"OBSERVED: created_at advanced from {} to {} across two KeepAdvertisedDirect ticks",
first.created_at, second.created_at
);
stop_tx
.send(())
.expect("loop task should still be listening for stop");
let (loop_session, publisher_identity) =
tokio::time::timeout(std::time::Duration::from_secs(5), loop_task)
.await
.expect("keep_advertised_direct should return promptly after stop")
.expect("loop task should not panic");
println!("OBSERVED: keep_advertised_direct returned promptly after stop");
reader_session
.close(
"normal",
Some("keep_advertised_direct reader test done"),
&reader_identity,
)
.await;
loop_session
.close(
"normal",
Some("keep_advertised_direct loop session done"),
&publisher_identity,
)
.await;
}