#![cfg(feature = "iicp-tcp")]
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use std::time::Duration;
use base64::{engine::general_purpose::URL_SAFE_NO_PAD, Engine as _};
use ed25519_dalek::{Signer, SigningKey};
use iicp_client::confidentiality::{
decrypt_payload_with_context, decrypt_response, encrypt_payload_with_context, encrypt_response,
};
use iicp_client::node::{IicpNode, NodeConfig};
use iicp_client::relay_session::{HttpPollWorkerSession, RelaySession, RelaySessionRegistry};
use iicp_client::CxPublicKey;
use serde_json::{json, Value};
use std::collections::HashMap;
use x25519_dalek::{PublicKey as X25519PublicKey, StaticSecret};
const INTENT: &str = "urn:iicp:intent:llm:chat:v1";
fn free_port() -> u16 {
std::net::TcpListener::bind("127.0.0.1:0")
.unwrap()
.local_addr()
.unwrap()
.port()
}
async fn spawn_relay() -> u16 {
let mut cfg = NodeConfig::new("relay-node", "http://relay.local", INTENT);
cfg.relay_capable = true;
cfg.relay_accept_port = free_port();
spawn_relay_with_config(cfg).await
}
async fn spawn_relay_with_config(cfg: NodeConfig) -> u16 {
let port = free_port();
let node = Arc::new(IicpNode::new(cfg));
tokio::spawn(async move {
let handler =
move |_req: iicp_client::node::TaskRequest| async move { Ok(json!({"echo": true})) };
let _ = node
.serve(handler, &format!("127.0.0.1:{port}"), None)
.await;
});
let client = reqwest::Client::new();
for _ in 0..50 {
if client
.get(format!("http://127.0.0.1:{port}/iicp/health"))
.timeout(Duration::from_millis(300))
.send()
.await
.is_ok()
{
break;
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
port
}
fn signed_ticket(worker_id: &str, relay_id: &str) -> (String, String) {
static NEXT_JTI: AtomicU64 = AtomicU64::new(1);
let sk = SigningKey::from_bytes(&[9u8; 32]);
let payload = serde_json::json!({
"v": 1, "typ": "relay-bind-ticket", "jti": format!("{:032x}", NEXT_JTI.fetch_add(1, Ordering::Relaxed)),
"iss": "test", "sub": worker_id, "aud": relay_id,
"iat": 1, "exp": 9_999_999_999_i64,
})
.to_string();
let b64 = URL_SAFE_NO_PAD.encode(payload.as_bytes());
let mut msg = b"iicp:relay-bind-ticket:v1\n".to_vec();
msg.extend_from_slice(b64.as_bytes());
let sig = sk.sign(&msg);
(
format!("{}.{}", b64, hex::encode(sig.to_bytes())),
hex::encode(sk.verifying_key().to_bytes()),
)
}
async fn bind(
client: &reqwest::Client,
port: u16,
worker_id: &str,
models: &[&str],
) -> reqwest::Response {
client
.post(format!("http://127.0.0.1:{port}/v1/relay/bind"))
.json(&json!({"worker_id": worker_id, "intent": INTENT, "models": models}))
.send()
.await
.unwrap()
}
#[tokio::test]
async fn session_forward_pull_result_roundtrip() {
let sess =
HttpPollWorkerSession::new("w-browser".into(), INTENT.into(), vec!["tinyllama".into()]);
let worker = {
let sess = sess.clone();
tokio::spawn(async move {
let call = sess.next_call(Duration::from_secs(5)).await.unwrap();
assert_eq!(call["task"]["payload"]["q"], 1);
sess.on_response(
call["call_id"].as_str().unwrap(),
json!({"result": {"a": 2}}),
);
})
};
let result = sess
.forward_task(&json!({"payload": {"q": 1}}), 5)
.await
.unwrap();
worker.await.unwrap();
assert_eq!(result, json!({"result": {"a": 2}}));
}
#[tokio::test]
async fn session_next_call_times_out_to_none() {
let sess = HttpPollWorkerSession::new("w-idle".into(), String::new(), vec![]);
assert!(sess.next_call(Duration::from_millis(50)).await.is_none());
}
#[tokio::test]
async fn session_liveness_window_and_close() {
let sess = HttpPollWorkerSession::with_liveness_window(
"w-live".into(),
String::new(),
vec![],
Duration::from_millis(50),
);
assert!(sess.is_alive());
tokio::time::sleep(Duration::from_millis(80)).await;
assert!(!sess.is_alive()); let fresh = HttpPollWorkerSession::new("w-fresh".into(), String::new(), vec![]);
fresh.close();
assert!(!fresh.is_alive());
}
#[test]
fn registry_get_by_token() {
let reg = RelaySessionRegistry::new();
let sess = HttpPollWorkerSession::new("w-tok".into(), String::new(), vec![]);
let token = sess.session_token.clone();
reg.bind("w-tok".into(), RelaySession::HttpPoll(sess));
assert!(reg.get_by_token(&token).is_some());
assert!(reg.get_by_token("wrong").is_none());
assert!(reg.get_by_token("").is_none());
}
#[tokio::test]
async fn bind_returns_token_and_cors() {
let port = spawn_relay().await;
let client = reqwest::Client::new();
let resp = bind(&client, port, "w-bind-1", &["tinyllama-1.1b"]).await;
assert_eq!(resp.status(), 200);
assert_eq!(
resp.headers().get("access-control-allow-origin").unwrap(),
"*"
);
let body: Value = resp.json().await.unwrap();
assert!(body["session_token"].as_str().unwrap().len() >= 32);
assert_eq!(body["worker_endpoint_path"], "/v1/relay-for/w-bind-1");
}
#[tokio::test]
async fn alive_rebind_rejected_409() {
let port = spawn_relay().await;
let client = reqwest::Client::new();
assert_eq!(bind(&client, port, "w-rebind", &[]).await.status(), 200);
let resp2 = bind(&client, port, "w-rebind", &[]).await;
assert_eq!(resp2.status(), 409);
let body: Value = resp2.json().await.unwrap();
assert_eq!(body["error"]["code"], "IICP-E038");
}
#[tokio::test]
async fn strict_bind_ticket_mode_accepts_valid_http_poll_ticket_and_rejects_wrong_worker() {
let (good_ticket, pub_hex) = signed_ticket("w-http-ticket", "relay-node");
let (bad_ticket, _) = signed_ticket("attacker", "relay-node");
let mut cfg = NodeConfig::new("relay-node", "http://relay.local", INTENT);
cfg.relay_capable = true;
cfg.relay_accept_port = free_port();
cfg.relay_bind_ticket_public_key_hex = Some(pub_hex);
cfg.relay_require_bind_ticket = true;
let port = spawn_relay_with_config(cfg).await;
let client = reqwest::Client::new();
let ok = client
.post(format!("http://127.0.0.1:{port}/v1/relay/bind"))
.json(&json!({"worker_id": "w-http-ticket", "intent": INTENT, "models": [], "bind_ticket": good_ticket}))
.send()
.await
.unwrap();
assert_eq!(ok.status(), 200);
let accepted: Value = ok.json().await.unwrap();
let session_token = accepted["session_token"].as_str().unwrap();
let unbind = client
.post(format!("http://127.0.0.1:{port}/v1/relay/unbind"))
.bearer_auth(session_token)
.json(&json!({}))
.send()
.await
.unwrap();
assert_eq!(unbind.status(), 204);
let replay = client
.post(format!("http://127.0.0.1:{port}/v1/relay/bind"))
.json(&json!({"worker_id": "w-http-ticket", "intent": INTENT, "models": [], "bind_ticket": good_ticket}))
.send()
.await
.unwrap();
assert_eq!(replay.status(), 409);
let replay_body: Value = replay.json().await.unwrap();
assert!(replay_body["error"]["message"]
.as_str()
.unwrap()
.contains("replayed"));
let wrong = client
.post(format!("http://127.0.0.1:{port}/v1/relay/bind"))
.json(&json!({"worker_id": "w-http-ticket-2", "intent": INTENT, "models": [], "bind_ticket": bad_ticket}))
.send()
.await
.unwrap();
assert_eq!(wrong.status(), 401);
let body: Value = wrong.json().await.unwrap();
assert_eq!(body["error"]["code"], "IICP-E040");
let missing = bind(&client, port, "w-http-ticket-missing", &[]).await;
assert_eq!(missing.status(), 401);
}
#[tokio::test]
async fn pull_and_result_require_bearer() {
let port = spawn_relay().await;
let client = reqwest::Client::new();
let r = client
.get(format!("http://127.0.0.1:{port}/v1/relay/pull"))
.send()
.await
.unwrap();
assert_eq!(r.status(), 401);
let r = client
.get(format!("http://127.0.0.1:{port}/v1/relay/pull"))
.bearer_auth("nope")
.send()
.await
.unwrap();
assert_eq!(r.status(), 401);
let r = client
.post(format!("http://127.0.0.1:{port}/v1/relay/result"))
.bearer_auth("nope")
.json(&json!({"call_id": "x", "result": {}}))
.send()
.await
.unwrap();
assert_eq!(r.status(), 401);
}
#[tokio::test]
async fn full_dispatch_roundtrip_via_relay_for() {
let port = spawn_relay().await;
let client = reqwest::Client::new();
let resp = bind(&client, port, "w-roundtrip", &["tinyllama-1.1b"]).await;
assert_eq!(resp.status(), 200);
let body: Value = resp.json().await.unwrap();
let token = body["session_token"].as_str().unwrap().to_string();
let worker = {
let client = client.clone();
tokio::spawn(async move {
let pull = client
.get(format!("http://127.0.0.1:{port}/v1/relay/pull"))
.bearer_auth(&token)
.timeout(Duration::from_secs(35))
.send()
.await
.unwrap();
assert_eq!(pull.status(), 200);
let call: Value = pull.json().await.unwrap();
assert_eq!(call["task"]["task_id"], "t-1");
client
.post(format!("http://127.0.0.1:{port}/v1/relay/result"))
.bearer_auth(&token)
.json(&json!({
"call_id": call["call_id"],
"result": {"result": {"text": "MESH OK from browser"}}
}))
.send()
.await
.unwrap();
})
};
let dispatch = client
.post(format!(
"http://127.0.0.1:{port}/v1/relay-for/w-roundtrip/v1/task"
))
.timeout(Duration::from_secs(40))
.json(&json!({
"task_id": "t-1",
"intent": INTENT,
"payload": {"messages": [{"role": "user", "content": "hi"}]},
"constraints": {"timeout_ms": 120000, "qos": "best_effort"}
}))
.send()
.await
.unwrap();
worker.await.unwrap();
assert_eq!(dispatch.status(), 200);
assert_eq!(
dispatch
.headers()
.get("access-control-allow-origin")
.unwrap(),
"*"
);
let resp: Value = dispatch.json().await.unwrap();
assert_eq!(resp["status"], "completed");
assert_eq!(resp["result"]["text"], "MESH OK from browser");
}
#[tokio::test]
async fn strict_bind_keeps_encrypted_request_and_response_opaque_through_relay() {
let worker_id = "w-cx-roundtrip";
let task_id = "t-cx-relay-1";
let secret_text = "relay must never observe this plaintext";
let (bind_ticket, bind_public_key) = signed_ticket(worker_id, "relay-node");
let mut cfg = NodeConfig::new("relay-node", "http://relay.local", INTENT);
cfg.relay_capable = true;
cfg.relay_accept_port = free_port();
cfg.relay_bind_ticket_public_key_hex = Some(bind_public_key);
cfg.relay_require_bind_ticket = true;
let port = spawn_relay_with_config(cfg).await;
let client = reqwest::Client::new();
let private_key = StaticSecret::from([42u8; 32]);
let public_key = X25519PublicKey::from(&private_key);
let cx_public_key = CxPublicKey {
algorithm: "X25519".to_string(),
encoding: Some("base64url".to_string()),
key: URL_SAFE_NO_PAD.encode(public_key.as_bytes()),
key_id: "cx-relay-roundtrip".to_string(),
features: vec!["response_encryption_v1".to_string()],
};
let (request_envelope, consumer_secret) = encrypt_payload_with_context(
&json!({"secret": secret_text}),
&cx_public_key,
task_id,
INTENT,
)
.unwrap();
let bind = client
.post(format!("http://127.0.0.1:{port}/v1/relay/bind"))
.json(&json!({
"worker_id": worker_id,
"intent": INTENT,
"models": [],
"bind_ticket": bind_ticket,
}))
.send()
.await
.unwrap();
assert_eq!(bind.status(), 200);
let bind_body: Value = bind.json().await.unwrap();
let session_token = bind_body["session_token"].as_str().unwrap().to_string();
let worker_client = client.clone();
let private_bytes = private_key.to_bytes();
let worker = tokio::spawn(async move {
let pull = worker_client
.get(format!("http://127.0.0.1:{port}/v1/relay/pull"))
.bearer_auth(&session_token)
.timeout(Duration::from_secs(35))
.send()
.await
.unwrap();
assert_eq!(pull.status(), 200);
let call: Value = pull.json().await.unwrap();
let task = call["task"].clone();
assert!(task.get("payload").is_none());
let envelope: HashMap<String, Value> =
serde_json::from_value(task["iicp_conf"].clone()).unwrap();
let (opened, worker_secret) =
decrypt_payload_with_context(&envelope, &private_bytes).unwrap();
assert_eq!(opened, json!({"secret": secret_text}));
let plain_response = json!({
"task_id": task_id,
"status": "success",
"result": {"text": "encrypted relay response"},
});
let response_envelope = encrypt_response(&plain_response, &worker_secret, task_id).unwrap();
worker_client
.post(format!("http://127.0.0.1:{port}/v1/relay/result"))
.bearer_auth(&session_token)
.json(&json!({
"call_id": call["call_id"],
"result": {"iicp_conf_resp": response_envelope},
}))
.send()
.await
.unwrap();
task.to_string()
});
let dispatch = client
.post(format!(
"http://127.0.0.1:{port}/v1/relay-for/{worker_id}/v1/task"
))
.timeout(Duration::from_secs(40))
.json(&json!({
"task_id": task_id,
"intent": INTENT,
"iicp_conf": request_envelope,
"cx_response_encryption": "required",
}))
.send()
.await
.unwrap();
let relay_visible = worker.await.unwrap();
assert_eq!(dispatch.status(), 200);
let response: Value = dispatch.json().await.unwrap();
assert!(response.get("result").is_none());
let response_envelope: HashMap<String, Value> =
serde_json::from_value(response["iicp_conf_resp"].clone()).unwrap();
let opened_response = decrypt_response(&response_envelope, &consumer_secret, task_id).unwrap();
assert_eq!(
opened_response["result"]["text"],
"encrypted relay response"
);
assert!(!relay_visible.contains(secret_text));
assert!(!response.to_string().contains("encrypted relay response"));
}
#[tokio::test]
async fn relay_for_unknown_worker_404() {
let port = spawn_relay().await;
let client = reqwest::Client::new();
let resp = client
.post(format!(
"http://127.0.0.1:{port}/v1/relay-for/w-ghost/v1/task"
))
.json(&json!({"task_id": "t-x"}))
.send()
.await
.unwrap();
assert_eq!(resp.status(), 404);
let body: Value = resp.json().await.unwrap();
assert_eq!(body["error"]["code"], "IICP-E030");
}
#[tokio::test]
async fn relay_for_health_reflects_session() {
let port = spawn_relay().await;
let client = reqwest::Client::new();
assert_eq!(
bind(&client, port, "w-health", &["m1", "m2"])
.await
.status(),
200
);
let resp = client
.get(format!(
"http://127.0.0.1:{port}/v1/relay-for/w-health/iicp/health"
))
.send()
.await
.unwrap();
assert_eq!(resp.status(), 200);
let health: Value = resp.json().await.unwrap();
assert_eq!(health["status"], "ok");
assert_eq!(health["via_relay"], true);
assert_eq!(health["models"], json!(["m1", "m2"]));
}
#[tokio::test]
async fn unbind_releases_worker_id() {
let port = spawn_relay().await;
let client = reqwest::Client::new();
let resp = bind(&client, port, "w-unbind", &[]).await;
let body: Value = resp.json().await.unwrap();
let token = body["session_token"].as_str().unwrap().to_string();
let u = client
.post(format!("http://127.0.0.1:{port}/v1/relay/unbind"))
.bearer_auth(&token)
.json(&json!({}))
.send()
.await
.unwrap();
assert_eq!(u.status(), 204);
assert_eq!(bind(&client, port, "w-unbind", &[]).await.status(), 200);
}
#[tokio::test]
async fn options_preflight() {
let port = spawn_relay().await;
let client = reqwest::Client::new();
let resp = client
.request(
reqwest::Method::OPTIONS,
format!("http://127.0.0.1:{port}/v1/relay/bind"),
)
.send()
.await
.unwrap();
assert_eq!(resp.status(), 204);
assert_eq!(
resp.headers().get("access-control-allow-origin").unwrap(),
"*"
);
assert!(resp
.headers()
.get("access-control-allow-headers")
.unwrap()
.to_str()
.unwrap()
.contains("Authorization"));
}
#[tokio::test]
async fn options_preflight_on_task() {
let port = spawn_relay().await;
let client = reqwest::Client::new();
let r = client
.request(
reqwest::Method::OPTIONS,
format!("http://127.0.0.1:{port}/v1/task"),
)
.send()
.await
.unwrap();
assert_eq!(r.status(), 204);
assert_eq!(r.headers().get("access-control-allow-origin").unwrap(), "*");
}
#[tokio::test]
async fn health_and_task_carry_cors() {
let port = spawn_relay().await;
let client = reqwest::Client::new();
let h = client
.get(format!("http://127.0.0.1:{port}/iicp/health"))
.send()
.await
.unwrap();
assert_eq!(h.headers().get("access-control-allow-origin").unwrap(), "*");
let t = client
.post(format!("http://127.0.0.1:{port}/v1/task"))
.json(&serde_json::json!({"task_id":"t-cors","intent":"urn:iicp:intent:llm:chat:v1","payload":{"messages":[]}}))
.send()
.await
.unwrap();
assert_eq!(t.headers().get("access-control-allow-origin").unwrap(), "*");
}
#[test]
fn session_cap_excludes_rebind() {
use iicp_client::relay_session::{HttpPollWorkerSession, RelaySession, RelaySessionRegistry};
let reg = RelaySessionRegistry::new();
for i in 0..256 {
reg.bind(
format!("w-{i}"),
RelaySession::HttpPoll(HttpPollWorkerSession::new(
format!("w-{i}"),
String::new(),
vec![],
)),
);
}
assert_eq!(reg.count(), 256);
assert!(reg.at_capacity("w-new")); assert!(!reg.at_capacity("w-0")); }