#![cfg(all(
not(target_arch = "wasm32"),
feature = "testing-endpoints",
feature = "test-relay-client"
))]
use anyhow::{bail, ensure, Context, Result};
use base64::{engine::general_purpose::URL_SAFE_NO_PAD, Engine as _};
use ed25519_dalek::{Signer as _, SigningKey};
use openrtc::client::TransportConfig;
use openrtc::native::{ControlPlane, DeviceSigner, Features, RoomArchitectureMode, RoomOptions};
use openrtc::Client;
use serde_json::{json, Value};
use std::{env, sync::Arc, time::Duration};
struct TestSigner(SigningKey);
impl DeviceSigner for TestSigner {
fn public_jwk(&self, _app_tag: &str) -> Result<Value> {
Ok(json!({ "kty": "OKP", "crv": "Ed25519",
"x": URL_SAFE_NO_PAD.encode(self.0.verifying_key().to_bytes()) }))
}
fn sign(&self, _app_tag: &str, challenge: &[u8]) -> Result<Vec<u8>> {
Ok(self.0.sign(challenge).to_bytes().to_vec())
}
}
fn local_origin(name: &str, scheme: &str) -> Result<String> {
let value = env::var(name).with_context(|| format!("missing {name}"))?;
let url = reqwest::Url::parse(&value)?;
ensure!(
url.scheme() == scheme
&& matches!(url.host_str(), Some("127.0.0.1" | "localhost" | "[::1]"))
&& url.port().is_some_and(|port| port > 0)
&& url.username().is_empty()
&& url.password().is_none()
&& url.query().is_none()
&& url.fragment().is_none()
&& url.path() == "/",
"Rust consumer requires explicit loopback origins"
);
Ok(value)
}
#[derive(Clone, Copy, PartialEq, Eq)]
enum Scenario {
Smoke,
Stability,
Renewal,
Rejoin,
GatewayOutage,
AdministrativeRevocation,
}
async fn run_consumer(scenario: Scenario) -> Result<()> {
let stability = scenario == Scenario::Stability;
let outage = scenario == Scenario::GatewayOutage;
let renewal = scenario == Scenario::Renewal || outage;
let source_digest = env::var("OPENRTC_RUST_CONSUMER_SOURCE_DIGEST")?;
ensure!(
source_digest.len() == 64
&& source_digest.bytes().all(|byte| byte.is_ascii_hexdigit())
&& Some(source_digest.as_str()) == option_env!("OPENRTC_RUST_CONSUMER_SOURCE_DIGEST"),
"Rust consumer source is stale or unbound"
);
let run_id = env::var("OPENRTC_RUST_CONSUMER_RUN_ID")?;
ensure!(
!run_id.is_empty()
&& run_id.len() <= 80
&& run_id
.bytes()
.all(|byte| byte.is_ascii_alphanumeric() || byte == b'-'),
"invalid run id"
);
let room_id = format!("rust-service-{run_id}");
let api_key = "pk_test_0000000000000000000000000000000000000000";
let control_url = local_origin("OPENRTC_CONTROL_PLANE_URL", "http")?;
let gateway_url = local_origin("OPENRTC_COORDINATION_GATEWAY_URL", "http")?;
let relay_url = local_origin("OPENRTC_TEST_IROH_RELAY_URL", "https")?;
let mut seed = [0u8; 32];
getrandom::getrandom(&mut seed)
.map_err(|error| anyhow::anyhow!("generate test signer: {error}"))?;
let control =
ControlPlane::anonymous(api_key, Arc::new(TestSigner(SigningKey::from_bytes(&seed))))?
.with_testing_endpoints(&control_url, &gateway_url)?;
for phase in 0..if scenario == Scenario::Rejoin { 2 } else { 1 } {
let room = control
.join_room(
&room_id,
"rust-consumer",
RoomOptions {
max_peers: Some(8),
architecture: RoomArchitectureMode::Authority,
features: Features {
iroh_relay: true,
..Features::default()
},
..RoomOptions::default()
},
)
.await?;
let client = room
.compose_client(
Client::builder(api_key.to_string(), Box::new(|| None))?.transport_config(
TransportConfig {
relay: true,
webrtc: None,
..TransportConfig::default()
},
),
)
.await?;
let mut states = client.connection_state_updates();
let mut messages = client.subscribe_native_peer_data();
let endpoint = iroh::Endpoint::builder(iroh::endpoint::presets::Minimal)
.relay_mode(iroh::RelayMode::custom([
relay_url.parse::<iroh::RelayUrl>()?
]))
.ca_tls_config(iroh::tls::CaTlsConfig::insecure_skip_verify())
.alpns(vec![openrtc::native_node::PlutoniumProtocol::ALPN.to_vec()])
.bind_addr("127.0.0.1:0".parse::<std::net::SocketAddr>()?)?
.bind()
.await?;
let result: Result<()> = async {
tokio::time::timeout(Duration::from_secs(10), endpoint.online()).await?;
client
.adopt_endpoint_with_router_mode(endpoint.clone(), true)
.await?;
let mut ticket = client
.endpoint_ticket_with_token(&format!("v2:room:{room_id}"), 8)
.await?;
let old_token = if renewal {
use openrtc::session_token::{decode_token_payload, expiring_ticket, split_ticket};
let (endpoint_ticket, suffix) = split_ticket(&ticket);
let payload = decode_token_payload(suffix.context("missing local ticket payload")?)
.context("invalid local ticket payload")?;
let expires = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)?.as_millis() as u64 + 90_000;
client.register_token_until(payload.token.clone(), payload.scope.to_string(), 8, expires);
ticket = expiring_ticket(endpoint_ticket, &payload.token, payload.scope, 8, Some(expires));
Some((payload.token, expires))
} else {
None
};
client
.update_presence(&room_id, "Independent Rust consumer", &ticket, None)
.await?;
let connection = loop {
let state = states.recv().await?;
if state.routable {
break state;
}
};
ensure!(
connection.active_transport_stable_id.is_some(),
"missing live transport proof"
);
let server_ticket = || async {
let server = room.signaling().list_devices(&room_id, None).await?
.into_iter().find(|device| device.device_id == "server")
.context("assigned service is missing")?;
let ticket = server.ticket.context("assigned service ticket is missing")?;
let (_, suffix) = openrtc::session_token::split_ticket(&ticket);
openrtc::session_token::decode_token_payload(suffix.context("missing service ticket suffix")?)
.context("invalid assigned service ticket")
};
let old_server_ticket = if renewal { Some(server_ticket().await?) } else { None };
let assert_stable = |state: &openrtc::client::StateSnapshot| -> Result<()> {
ensure!(state.connection_id == connection.connection_id, "unexpected service edge");
ensure!(
state.routable && !state.replacement_in_progress
&& state.active_transport_stable_id == connection.active_transport_stable_id
&& state.transport_generation == connection.transport_generation
&& state.route_generation == connection.route_generation
&& state.remote_node_id == connection.remote_node_id,
"service generation lost stability: state={} reason={}",
state.state, state.readiness_reason
);
Ok(())
};
let exchanges: u64 = if stability { 50 } else if outage { 2 } else if renewal { 22 } else { 1 };
let mut elapsed = {
let traffic = async {
let mut started = tokio::time::Instant::now();
for sequence in 0..exchanges {
let offset = if sequence == 49 { 300 } else if outage { sequence * 5 } else { sequence.saturating_sub(1) * 5 };
tokio::select! {
message = messages.recv() => bail!("unsolicited/duplicate service payload while idle: {}", message.is_ok()),
_ = tokio::time::sleep_until(started + Duration::from_secs(offset)) => {}
}
let request = json!({ "type": "rust-service-request", "runId": run_id,
"value": "rust-native-to-node", "sequence": sequence });
let deadline = if sequence == 0 { 15 } else { 2 };
tokio::time::timeout(Duration::from_secs(deadline), async {
client.send_peer(&connection.connection_id, &serde_json::to_vec(&request)?).await?;
let response = messages.recv().await?;
ensure!(response.connection_id == connection.connection_id, "reply from another connection");
ensure!(serde_json::from_slice::<Value>(&response.payload)? == json!({
"type": "rust-service-response", "runId": run_id,
"value": "node-to-rust-native", "sequence": sequence
}), "protected Rust reply mismatch at sequence {sequence}");
Ok::<(), anyhow::Error>(())
}).await.with_context(|| format!("local service round trip {sequence} exceeded {deadline} seconds"))??;
if sequence == 0 {
eprintln!("OPENRTC_RUST_SERVICE_CONVERGED initialRoundTripMs={}", started.elapsed().as_millis());
started = tokio::time::Instant::now();
}
}
Ok::<_, anyhow::Error>(started.elapsed())
};
tokio::pin!(traffic);
loop {
tokio::select! {
biased;
state = states.recv() => assert_stable(&state?)?,
result = &mut traffic => break result?,
}
}
};
let current = client.connection_states().await;
ensure!(
current.len() == 1,
"consumer must have exactly one assigned service route"
);
assert_stable(¤t[0])?;
if outage {
let outage_started = tokio::time::Instant::now();
let expiry = old_token.as_ref().context("missing expiring outage ticket")?.1;
let now = std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH)?.as_millis() as u64;
tokio::time::sleep(Duration::from_millis(expiry.saturating_sub(now) + 100)).await;
let denied = client.send_peer(&connection.connection_id, &serde_json::to_vec(&json!({
"type": "rust-service-request", "runId": run_id,
"value": "expired-must-not-deliver", "sequence": 99
}))?).await;
eprintln!("OPENRTC_RUST_SERVICE_EXPIRED_SEND localWriteAccepted={}", denied.is_ok());
ensure!(tokio::time::timeout(Duration::from_secs(2), messages.recv()).await.is_err(),
"expired application request produced a response");
while states.try_recv().is_ok() {}
println!("OPENRTC_RUST_SERVICE_OUTAGE_EXPIRED_DENIED");
let recovered = tokio::time::timeout(Duration::from_secs(45), async {
loop {
let state = states.recv().await?;
if state.routable && client.connection_states().await.iter().any(|current|
current.connection_id == state.connection_id && current.routable
&& current.active_transport_stable_id == state.active_transport_stable_id) {
break Ok::<_, anyhow::Error>(state);
}
}
}).await.context("service did not recover after gateway restoration")??;
ensure!(recovered.remote_node_id == connection.remote_node_id, "outage changed service identity");
tokio::time::timeout(Duration::from_secs(15), async {
client.send_peer(&recovered.connection_id, &serde_json::to_vec(&json!({
"type": "rust-service-request", "runId": run_id,
"value": "rust-native-to-node", "sequence": 2
}))?).await?;
let response = messages.recv().await?;
ensure!(response.connection_id == recovered.connection_id, "recovery reply from wrong route");
ensure!(serde_json::from_slice::<Value>(&response.payload)? == json!({
"type": "rust-service-response", "runId": run_id,
"value": "node-to-rust-native", "sequence": 2
}), "recovered protected reply mismatch");
Ok::<(), anyhow::Error>(())
}).await.context("recovered service payload timed out")??;
elapsed += outage_started.elapsed();
}
if let Some((old_token, expires)) = old_token {
ensure!(std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH)?.as_millis()
> u128::from(expires), "renewal proof did not cross the original ticket expiry");
ensure!(client.validate_session_token(&old_token).is_err(),
"retired service ticket remains accepted");
}
if let Some(old) = old_server_ticket {
let replacement = server_ticket().await?;
let old_expiry = old.expires_at_ms.context("old service ticket must expire")?;
ensure!(std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH)?.as_millis()
> u128::from(old_expiry), "renewal did not cross original server ticket expiry");
ensure!(replacement.token != old.token && replacement.expires_at_ms > old.expires_at_ms,
"gateway did not project the renewed server ticket");
}
if scenario == Scenario::Rejoin && phase == 0 {
let (_, _, mut held_send, _) = client.open_peer_bi(&connection.connection_id, None).await?;
let revoked = client.revoke_tokens_by_scope(&format!("v2:room:{room_id}")).await;
ensure!(revoked.contains(&connection.connection_id), "scope revocation missed the service route");
ensure!(held_send.write_all(b"revoked-stream-must-not-deliver").await.is_err(),
"revoked stream remained writable");
ensure!(client.send_peer(&connection.connection_id, b"revoked-message-must-not-deliver").await.is_err(),
"revoked route remained writable");
ensure!(!client.connection_states().await.iter().any(|state| state.routable),
"revoked consumer retained a routable connection");
}
if scenario == Scenario::AdministrativeRevocation {
use std::io::{BufRead, Write};
let started = tokio::time::Instant::now();
let (_, _, mut held_send, _) = client.open_peer_bi(&connection.connection_id, None).await?;
println!("OPENRTC_ADMIN_REVOCATION_READY {}", room.principal_id);
std::io::stdout().flush()?;
let acknowledgement = tokio::task::spawn_blocking(|| -> std::io::Result<String> {
let mut line = String::new();
std::io::stdin().lock().read_line(&mut line)?;
Ok(line)
}).await??;
ensure!(acknowledgement.trim() == "OPENRTC_ADMIN_REVOKED", "administrative API failed");
tokio::time::timeout(Duration::from_secs(5), async {
while client.connection_states().await.iter().any(|state| state.routable) {
tokio::time::sleep(Duration::from_millis(20)).await;
}
}).await.context("administrative revocation retained a routable service edge")?;
ensure!(held_send.write_all(b"administratively-revoked-stream").await.is_err(),
"administratively revoked held stream remained writable");
ensure!(client.send_peer(&connection.connection_id, b"administratively-revoked-message").await.is_err(),
"administratively revoked service remained writable");
ensure!(tokio::time::timeout(Duration::from_secs(1), messages.recv()).await.is_err(),
"unexpected payload after administrative revocation");
ensure!(!client.connection_states().await.iter().any(|state| state.routable),
"revoked service route reappeared without authorization");
println!("OPENRTC_ADMIN_REVOCATION_DENIED");
elapsed += started.elapsed();
}
println!(
"OPENRTC_RUST_CONSUMER_OK {}",
json!({ "runId": run_id,
"sourceDigest": source_digest, "nodeId": endpoint.id().to_string(),
"exchanges": exchanges + u64::from(outage), "elapsedMs": elapsed.as_millis(),
"quietSeconds": if stability { 60 } else { 0 } })
);
Ok(())
}
.await;
room.close().await;
endpoint.close().await;
result?;
}
Ok(())
}
#[tokio::test]
#[ignore = "run only through test:rust-service:emulator"]
async fn public_rust_consumer() -> Result<()> {
match tokio::time::timeout(Duration::from_secs(60), run_consumer(Scenario::Smoke)).await {
Ok(result) => result,
Err(_) => bail!("independent Rust service consumer timed out"),
}
}
#[tokio::test]
#[ignore = "run only through test:rust-service:stability"]
async fn public_rust_consumer_five_minute() -> Result<()> {
match tokio::time::timeout(Duration::from_secs(360), run_consumer(Scenario::Stability)).await {
Ok(result) => result,
Err(_) => bail!("five-minute Rust service stability timed out"),
}
}
#[tokio::test]
#[ignore = "run only through test:rust-service:renewal"]
async fn public_rust_consumer_renews_ticket() -> Result<()> {
match tokio::time::timeout(Duration::from_secs(150), run_consumer(Scenario::Renewal)).await {
Ok(result) => result,
Err(_) => bail!("Rust service ticket renewal timed out"),
}
}
#[tokio::test]
#[ignore = "run only through test:rust-service:rejoin"]
async fn public_rust_consumer_revokes_and_rejoins() -> Result<()> {
match tokio::time::timeout(Duration::from_secs(100), run_consumer(Scenario::Rejoin)).await {
Ok(result) => result,
Err(_) => bail!("Rust service revocation/rejoin timed out"),
}
}
#[tokio::test]
#[ignore = "run only through test:rust-service:outage"]
async fn public_rust_consumer_gateway_outage() -> Result<()> {
match tokio::time::timeout(
Duration::from_secs(180),
run_consumer(Scenario::GatewayOutage),
)
.await
{
Ok(result) => result,
Err(_) => bail!("Rust service gateway outage timed out"),
}
}
#[tokio::test]
#[ignore = "run only through test:rust-service:admin-revocation"]
async fn public_rust_consumer_administrative_revocation() -> Result<()> {
tokio::time::timeout(
Duration::from_secs(60),
run_consumer(Scenario::AdministrativeRevocation),
)
.await
.context("Rust service administrative revocation timed out")?
}