use std::time::Duration;
use rsiprtp::ice::{Candidate, IceRole};
use rsiprtp::sdp::builder::{MediaBuilder, SdpBuilder};
use rsiprtp::sdp::ice_attrs;
use rsiprtp::sdp::parser::SessionDescription;
use rsiprtp::session::{
CallEndReason, CallEvent, CallManager, CallState, Dialog, IceAnswerInputs, IceRemoteParams,
IceSession, ManagerConfig, ManagerEvent,
};
const SHORT: Duration = Duration::from_millis(500);
const RUN_CHECKS_BUDGET: Duration = Duration::from_secs(2);
const PROBE_BUDGET: Duration = Duration::from_secs(2);
fn extract_ice_params(sdp: &SessionDescription) -> IceRemoteParams {
let audio = sdp.audio_media().expect("audio media");
let (ufrag, pwd) =
ice_attrs::read_ice_credentials(audio).expect("offer/answer carries ICE credentials");
let candidates = ice_attrs::read_candidates(audio);
assert!(
!candidates.is_empty(),
"peer SDP must carry at least one a=candidate line"
);
IceRemoteParams {
ufrag,
pwd,
candidates,
}
}
fn build_offer(
default: &Candidate,
local: &rsiprtp::session::IceLocalParams,
) -> SessionDescription {
let local_ip = default.address.ip();
let mut sdp = SdpBuilder::new(local_ip)
.session_name("rsiprtp ICE test")
.add_media(MediaBuilder::audio(default.address.port()).pcmu())
.build();
ice_attrs::apply_default_candidate(&mut sdp, 0, default);
let audio = sdp
.media
.get_mut(0)
.expect("audio media just inserted by builder");
ice_attrs::write_ice_credentials(audio, &local.ufrag, &local.pwd);
ice_attrs::write_candidates(audio, &local.candidates);
ice_attrs::write_rtcp_mux(audio);
sdp
}
fn dialog_pair() -> (Dialog, Dialog) {
let call_id = "ice-basic-call".to_string();
let from_tag = "alice-tag".to_string();
let to_tag = "bob-tag".to_string();
let alice_uri = "sip:alice@127.0.0.1".to_string();
let bob_uri = "sip:bob@127.0.0.1".to_string();
let uac = Dialog::new_uac(
call_id.clone(),
from_tag.clone(),
to_tag.clone(),
alice_uri.clone(),
bob_uri.clone(),
1,
);
let uas = Dialog::new_uas(call_id, from_tag, to_tag, bob_uri, alice_uri, 1);
(uac, uas)
}
#[tokio::test(flavor = "multi_thread")]
async fn ice_basic_two_managers_loopback_handshake() {
let mut a_manager = CallManager::new(ManagerConfig::default());
let mut a_ice = IceSession::gather(IceRole::Controlling, vec![], SHORT, SHORT)
.await
.expect("A: gather host candidates");
let a_call_id = a_manager.create_call("sip:bob@127.0.0.1".to_string());
let a_default = a_ice
.default_candidate()
.expect("A: at least one host candidate")
.clone();
let offer = build_offer(&a_default, a_ice.local());
let offer_wire = offer.to_string();
let offer_parsed = SessionDescription::parse(&offer_wire).expect("offer SDP parses");
let mut b_manager = CallManager::new(ManagerConfig::default());
let (uac_dialog, uas_dialog) = dialog_pair();
let b_call_id = b_manager
.accept_inbound_invite(uas_dialog, &offer_parsed)
.expect("B: accept_inbound_invite");
let b_events = b_manager.drain_events();
assert!(
b_events
.iter()
.any(|e| matches!(e, ManagerEvent::IncomingCall(id) if id == &b_call_id)),
"B: accept_inbound_invite must emit IncomingCall, got {:?}",
b_events
);
let mut b_ice = IceSession::gather(IceRole::Controlled, vec![], SHORT, SHORT)
.await
.expect("B: gather host candidates");
let b_default = b_ice
.default_candidate()
.expect("B: at least one host candidate")
.clone();
let answer = {
let inputs = IceAnswerInputs::new(&b_default, b_ice.local());
b_manager
.build_answer_for(&b_call_id, &inputs)
.expect("B: build_answer_for")
};
let answer_wire = answer.to_string();
let answer_parsed = SessionDescription::parse(&answer_wire).expect("answer SDP parses");
let answer_audio = answer_parsed
.audio_media()
.expect("answer carries audio media");
assert!(
ice_attrs::read_rtcp_mux(answer_audio),
"answer must propagate a=rtcp-mux from the offer"
);
let a_answer_dialog = uac_dialog.clone();
assert!(
a_manager.handle_invite_success(
&a_call_id,
a_answer_dialog,
&answer_parsed,
None,
std::time::Instant::now(),
),
"A: handle_invite_success accepts the answer"
);
let a_events = a_manager.drain_events();
assert!(
a_events.iter().any(|e| matches!(
e,
ManagerEvent::CallStateChanged(id, CallState::Established) if id == &a_call_id
)),
"A: handle_invite_success must emit CallStateChanged(_, Established), got {:?}",
a_events
);
let a_remote = extract_ice_params(&answer_parsed);
let b_remote = extract_ice_params(&offer_parsed);
let join = tokio::time::timeout(RUN_CHECKS_BUDGET, async move {
let (a_res, b_res) = tokio::join!(a_ice.run_checks(a_remote), b_ice.run_checks(b_remote));
(a_ice, b_ice, a_res, b_res)
})
.await
.expect("ICE checks complete within budget");
let (a_ice, b_ice, a_res, b_res) = join;
a_res.expect("A: run_checks ok");
b_res.expect("B: run_checks ok");
assert!(
a_manager
.get_call(&a_call_id)
.expect("A: call exists")
.media()
.is_some(),
"A: call has a media session after handle_invite_success"
);
assert!(b_manager.answer_call(&b_call_id), "B: answer_call");
assert!(
b_manager
.get_call(&b_call_id)
.expect("B: call exists")
.media()
.is_some(),
"B: call has a media session after build_answer_for + answer_call"
);
let socket_a = a_ice.rtp_socket().expect("A: rtp socket");
let socket_b = b_ice.rtp_socket().expect("B: rtp socket");
let peer_a = a_ice.peer_addr().expect("A: peer addr");
let peer_b = b_ice.peer_addr().expect("B: peer addr");
assert_eq!(
peer_a,
socket_b.local_addr().expect("B socket local addr"),
"A's peer must be B's bound socket"
);
assert_eq!(
peer_b,
socket_a.local_addr().expect("A socket local addr"),
"B's peer must be A's bound socket"
);
let probe_a_to_b = b"hello-from-a";
socket_a
.send_to(probe_a_to_b, peer_a)
.await
.expect("A: send to B");
let mut buf = [0u8; 64];
let (n, _from) = tokio::time::timeout(PROBE_BUDGET, socket_b.recv_from(&mut buf))
.await
.expect("B: probe arrives within budget")
.expect("B: recv_from ok");
assert_eq!(&buf[..n], probe_a_to_b, "A->B probe bytes match");
let probe_b_to_a = b"hello-from-b";
socket_b
.send_to(probe_b_to_a, peer_b)
.await
.expect("B: send to A");
let (n, _from) = tokio::time::timeout(PROBE_BUDGET, socket_a.recv_from(&mut buf))
.await
.expect("A: probe arrives within budget")
.expect("A: recv_from ok");
assert_eq!(&buf[..n], probe_b_to_a, "B->A probe bytes match");
}
#[tokio::test(flavor = "multi_thread")]
async fn ice_run_checks_rejects_bad_message_integrity_and_call_cleans_up() {
let mut a_ice = IceSession::gather(IceRole::Controlling, vec![], SHORT, SHORT)
.await
.expect("A: gather");
let a_default = a_ice
.default_candidate()
.expect("A: at least one host candidate")
.clone();
let mut b_manager = CallManager::new(ManagerConfig::default());
let mut b_ice = IceSession::gather(IceRole::Controlled, vec![], SHORT, SHORT)
.await
.expect("B: gather");
let offer = build_offer(&a_default, a_ice.local());
let offer_wire = offer.to_string();
let offer_parsed = SessionDescription::parse(&offer_wire).expect("offer parses");
let (_uac_dialog, uas_dialog) = dialog_pair();
let b_call_id = b_manager
.accept_inbound_invite(uas_dialog, &offer_parsed)
.expect("B: accept_inbound_invite");
let _ = b_manager.drain_events();
let a_remote = IceRemoteParams {
ufrag: b_ice.local().ufrag.clone(),
pwd: b_ice.local().pwd.clone(),
candidates: b_ice.local().candidates.clone(),
};
let b_remote = IceRemoteParams {
ufrag: a_ice.local().ufrag.clone(),
pwd: "this-is-not-the-real-pwd".to_string(),
candidates: a_ice.local().candidates.clone(),
};
assert!(!a_remote.candidates.is_empty());
assert!(!b_remote.candidates.is_empty());
let join = tokio::time::timeout(RUN_CHECKS_BUDGET, async move {
let (a_res, b_res) = tokio::join!(a_ice.run_checks(a_remote), b_ice.run_checks(b_remote));
(a_res, b_res)
})
.await
.expect("checks complete within budget");
let (_a_res, b_res) = join;
let b_err = b_res.expect_err("B: run_checks must fail when its remote pwd is wrong");
let b_msg = format!("{}", b_err);
assert!(
b_msg.contains("validated") || b_msg.contains("checks failed"),
"B: error must mention validation failure, got {:?}",
b_msg
);
let dialog_id = b_manager
.reject_inbound_invite(&b_call_id)
.expect("B: reject_inbound_invite returns dialog id");
let call = b_manager
.get_call(&b_call_id)
.expect("B: call still present");
assert_eq!(call.state(), CallState::Terminated, "B: call terminated");
assert_eq!(call.dialog_id(), Some(&dialog_id));
let events = b_manager.drain_events();
assert!(
events.iter().any(|e| matches!(
e,
ManagerEvent::CallEvent(id, CallEvent::Ended(CallEndReason::Error)) if id == &b_call_id
)),
"B: rejection must emit Ended(Error), got {:?}",
events
);
assert!(b_manager.reject_inbound_invite(&b_call_id).is_none());
}