use crate::sip::headers::*;
use crate::sip::{Method, StatusCode, Version};
use crate::transaction::key::{TransactionKey, TransactionRole};
use crate::transaction::transaction::Transaction;
use crate::transaction::TransactionType;
use crate::transport::transport_layer::DomainResolver;
use crate::transport::udp::UdpConnection;
use crate::transport::{SipConnection, TransportLayer};
use crate::{EndpointBuilder, Result};
use std::time::Duration;
use tokio::time::{sleep, timeout};
use tokio_util::sync::CancellationToken;
struct NoopResolver;
#[async_trait::async_trait]
impl DomainResolver for NoopResolver {
async fn resolve(&self, _target: &crate::transport::SipAddr) -> Result<crate::transport::SipAddr> {
Err(crate::Error::DnsResolutionError("noop resolver".into()))
}
}
fn make_invite(branch: &str, to_tag: &str) -> crate::sip::Request {
crate::sip::Request {
method: Method::Invite,
uri: crate::sip::Uri::try_from("sip:bob@example.com").unwrap(),
headers: vec![
Via::new(&format!("SIP/2.0/UDP 127.0.0.1:5060;branch={}", branch)).into(),
CSeq::new("1 INVITE").into(),
From::new("Alice <sip:alice@example.com>;tag=aliceTag123").into(),
To::new(&format!("Bob <sip:bob@example.com>;tag={}", to_tag)).into(),
CallId::new("drop-test-call-id@example.com").into(),
MaxForwards::new("70").into(),
]
.into(),
version: Version::V2,
body: Default::default(),
}
}
fn make_ack_for_invite(invite: &crate::sip::Request) -> crate::sip::Request {
let mut headers = invite.headers.clone();
for header in headers.iter_mut() {
if let crate::sip::Header::CSeq(cseq) = header {
*cseq = crate::sip::headers::CSeq::new("1 ACK").into();
}
}
crate::sip::Request {
method: Method::Ack,
uri: invite.uri.clone(),
headers,
version: invite.version.clone(),
body: Default::default(),
}
}
#[tokio::test]
async fn test_cleanup_server_invite_completed_registers_finished() -> crate::Result<()> {
let endpoint = super::create_test_endpoint(Some("127.0.0.1:0")).await?;
let invite = make_invite("z9hG4bK-drop-unit-1", "serverTagAAA");
let key = TransactionKey::from_request(&invite, TransactionRole::Server)?;
let tx = Transaction::new_server(key.clone(), invite.clone(), endpoint.inner.clone(), None);
assert_eq!(tx.transaction_type, TransactionType::ServerInvite);
drop(tx);
assert!(
!endpoint.inner.finished_transactions.contains_key(&key),
"finished_transactions should be empty when dropped in Trying without a response"
);
Ok(())
}
#[tokio::test]
async fn test_cleanup_server_invite_completed_keeps_waiting_ack() -> crate::Result<()> {
let endpoint = super::create_test_endpoint(Some("127.0.0.1:0")).await?;
let invite = make_invite("z9hG4bK-drop-unit-2", "serverTagBBB");
let key = TransactionKey::from_request(&invite, TransactionRole::Server)?;
let mut tx = Transaction::new_server(key.clone(), invite.clone(), endpoint.inner.clone(), None);
let resp = crate::sip::Response {
status_code: StatusCode::ServiceUnavailable,
version: Version::V2,
headers: invite.headers.clone(),
body: Default::default(),
};
tx.last_response = Some(resp.clone());
tx.state = crate::transaction::TransactionState::Completed;
let dialog_id = crate::dialog::DialogId::try_from((&resp, TransactionRole::Server))?;
endpoint
.inner
.waiting_ack
.insert(dialog_id.clone(), key.clone());
drop(tx);
sleep(Duration::from_millis(50)).await;
assert!(
endpoint.inner.waiting_ack.contains_key(&dialog_id),
"waiting_ack must NOT be removed when ServerInvite is dropped in Completed state"
);
let finished = endpoint.inner.finished_transactions.get(&key);
assert!(
finished.is_some(),
"finished_transactions must contain the response for dropped ServerInvite in Completed"
);
if let Some(Some(crate::sip::SipMessage::Response(r))) = finished.map(|v| v.value().clone()) {
assert_eq!(r.status_code, StatusCode::ServiceUnavailable);
} else {
panic!("finished_transactions entry should be a Response with 503");
}
Ok(())
}
#[tokio::test]
async fn test_cleanup_server_invite_terminated_removes_waiting_ack() -> crate::Result<()> {
let endpoint = super::create_test_endpoint(Some("127.0.0.1:0")).await?;
let invite = make_invite("z9hG4bK-drop-unit-3", "serverTagCCC");
let key = TransactionKey::from_request(&invite, TransactionRole::Server)?;
let mut tx = Transaction::new_server(key.clone(), invite.clone(), endpoint.inner.clone(), None);
let resp = crate::sip::Response {
status_code: StatusCode::BusyHere,
version: Version::V2,
headers: invite.headers.clone(),
body: Default::default(),
};
tx.last_response = Some(resp.clone());
tx.state = crate::transaction::TransactionState::Terminated;
let dialog_id = crate::dialog::DialogId::try_from((&resp, TransactionRole::Server))?;
endpoint
.inner
.waiting_ack
.insert(dialog_id.clone(), key.clone());
drop(tx);
sleep(Duration::from_millis(50)).await;
assert!(
!endpoint.inner.waiting_ack.contains_key(&dialog_id),
"waiting_ack must be removed when ServerInvite terminates normally"
);
assert!(
endpoint.inner.finished_transactions.contains_key(&key),
"finished_transactions should contain the response for Timer J absorption"
);
Ok(())
}
#[tokio::test]
async fn test_server_invite_drop_retransmission_and_ack() {
let token = CancellationToken::new();
let server_conn = UdpConnection::create_connection("127.0.0.1:0".parse().unwrap(), None, None)
.await
.expect("create server connection");
let server_conn_sip: SipConnection = server_conn.clone().into();
let server_addr = server_conn_sip.get_addr().clone();
let tl = TransportLayer::new(token.child_token());
tl.add_transport(server_conn_sip.clone());
let endpoint = EndpointBuilder::new()
.with_user_agent("rsipstack-test")
.with_transport_layer(tl)
.build();
let client_conn = UdpConnection::create_connection("127.0.0.1:0".parse().unwrap(), None, None)
.await
.expect("create client connection");
let client_conn_sip: SipConnection = client_conn.clone().into();
let endpoint_inner = endpoint.inner.clone();
let serve_handle = tokio::spawn(async move {
let _ = endpoint_inner.serve().await;
});
let branch = "z9hG4bK-drop-int-branch1";
let invite = make_invite(branch, "serverTagInt1");
client_conn_sip
.send(invite.clone().into(), Some(&server_addr))
.await
.expect("send invite");
let mut incoming = endpoint.incoming_transactions().expect("incoming");
let mut tx = timeout(Duration::from_secs(2), incoming.recv())
.await
.expect("timeout waiting for incoming transaction")
.expect("no incoming transaction");
assert_eq!(tx.original.method, Method::Invite);
let tx_key = tx.key.clone();
tx.reply(StatusCode::ServiceUnavailable)
.await
.expect("reply 503");
drop(tx);
sleep(Duration::from_millis(100)).await;
assert!(
endpoint.inner.finished_transactions.contains_key(&tx_key),
"finished_transactions should contain the 503 response after drop"
);
assert_eq!(
endpoint.inner.waiting_ack.len(),
1,
"waiting_ack should still have 1 entry for the dropped ServerInvite"
);
let mut buf = vec![0u8; 4096];
let (len, _) = timeout(Duration::from_secs(2), client_conn.recv_raw(&mut buf))
.await
.expect("timeout waiting for 503 response")
.expect("recv failed");
let resp_str = String::from_utf8_lossy(&buf[..len]);
assert!(
resp_str.contains("503"),
"first response should be 503, got: {}",
&resp_str[..resp_str.len().min(200)]
);
client_conn_sip
.send(invite.clone().into(), Some(&server_addr))
.await
.expect("retransmit invite");
let (len, _) = timeout(Duration::from_secs(2), client_conn.recv_raw(&mut buf))
.await
.expect("timeout waiting for retransmitted 503 response")
.expect("recv failed on retransmission");
let resp_str = String::from_utf8_lossy(&buf[..len]);
assert!(
resp_str.contains("503"),
"retransmitted INVITE should get 503 auto-reply, got: {}",
&resp_str[..resp_str.len().min(200)]
);
let ack = make_ack_for_invite(&invite);
client_conn_sip
.send(ack.clone().into(), Some(&server_addr))
.await
.expect("send ack");
sleep(Duration::from_millis(200)).await;
assert_eq!(
endpoint.inner.waiting_ack.len(),
0,
"waiting_ack should be empty after ACK is absorbed"
);
match timeout(Duration::from_millis(300), client_conn.recv_raw(&mut buf)).await {
Err(_) => { }
Ok(Ok((len, _))) => {
let resp_str = String::from_utf8_lossy(&buf[..len]);
panic!(
"should not receive anything after ACK, got: {}",
&resp_str[..resp_str.len().min(200)]
);
}
Ok(Err(e)) => {
panic!("unexpected error waiting after ACK: {:?}", e);
}
}
serve_handle.abort();
}
#[tokio::test]
async fn test_cancel_after_final_response_does_not_reach_tu() {
let token = CancellationToken::new();
let server_conn = UdpConnection::create_connection("127.0.0.1:0".parse().unwrap(), None, None)
.await
.expect("create server connection");
let server_conn_sip: SipConnection = server_conn.clone().into();
let server_addr = server_conn_sip.get_addr().clone();
let tl = TransportLayer::new(token.child_token());
tl.add_transport(server_conn_sip);
let endpoint = EndpointBuilder::new()
.with_user_agent("rsipstack-test")
.with_transport_layer(tl)
.build();
let client_conn = UdpConnection::create_connection("127.0.0.1:0".parse().unwrap(), None, None)
.await
.expect("create client connection");
let client_conn_sip: SipConnection = client_conn.clone().into();
let endpoint_inner = endpoint.inner.clone();
let serve_handle = tokio::spawn(async move {
let _ = endpoint_inner.serve().await;
});
let invite = make_invite("z9hG4bK-late-cancel", "serverTagLateCancel");
client_conn_sip
.send(invite.clone().into(), Some(&server_addr))
.await
.expect("send invite");
let mut incoming = endpoint.incoming_transactions().expect("incoming");
let mut tx = timeout(Duration::from_secs(2), incoming.recv())
.await
.expect("timeout waiting for incoming transaction")
.expect("no incoming transaction");
tx.reply(StatusCode::TemporarilyUnavailable)
.await
.expect("reply 480");
let mut buf = vec![0u8; 4096];
let (len, _) = timeout(Duration::from_secs(2), client_conn.recv_raw(&mut buf))
.await
.expect("timeout waiting for 480 response")
.expect("recv failed");
let response = String::from_utf8_lossy(&buf[..len]);
assert!(response.contains("480"), "expected 480, got: {response}");
let mut cancel_headers = invite.headers.clone();
for header in cancel_headers.iter_mut() {
if let crate::sip::Header::CSeq(cseq) = header {
*cseq = CSeq::new("1 CANCEL").into();
}
}
let cancel = crate::sip::Request {
method: Method::Cancel,
uri: invite.uri.clone(),
headers: cancel_headers,
version: invite.version.clone(),
body: Default::default(),
};
client_conn_sip
.send(cancel.into(), Some(&server_addr))
.await
.expect("send cancel");
let delivered = timeout(Duration::from_millis(100), tx.receive()).await;
assert!(
delivered.is_err(),
"CANCEL after a final response must not reach the TU"
);
let (len, _) = timeout(Duration::from_secs(2), client_conn.recv_raw(&mut buf))
.await
.expect("timeout waiting for CANCEL 200 response")
.expect("recv failed");
let response = String::from_utf8_lossy(&buf[..len]);
assert!(
response.contains("200"),
"CANCEL should receive 200 OK, got: {response}"
);
assert_eq!(
tx.last_response
.as_ref()
.map(|response| &response.status_code),
Some(&StatusCode::TemporarilyUnavailable),
"late CANCEL must not replace the INVITE final response"
);
serve_handle.abort();
}
#[tokio::test]
async fn test_server_invite_normal_termination_registers_finished() {
let token = CancellationToken::new();
let server_conn = UdpConnection::create_connection("127.0.0.1:0".parse().unwrap(), None, None)
.await
.expect("create server connection");
let server_conn_sip: SipConnection = server_conn.clone().into();
let server_addr = server_conn_sip.get_addr().clone();
let tl = TransportLayer::new(token.child_token());
tl.add_transport(server_conn_sip.clone());
let endpoint = EndpointBuilder::new()
.with_user_agent("rsipstack-test")
.with_transport_layer(tl)
.build();
let client_conn = UdpConnection::create_connection("127.0.0.1:0".parse().unwrap(), None, None)
.await
.expect("create client connection");
let client_conn_sip: SipConnection = client_conn.clone().into();
let endpoint_inner = endpoint.inner.clone();
let serve_handle = tokio::spawn(async move {
let _ = endpoint_inner.serve().await;
});
let branch = "z9hG4bK-normal-term-branch";
let invite = make_invite(branch, "serverTagNormal");
client_conn_sip
.send(invite.clone().into(), Some(&server_addr))
.await
.expect("send invite");
let mut incoming = endpoint.incoming_transactions().expect("incoming");
let server_loop = async {
let mut tx = timeout(Duration::from_secs(2), incoming.recv())
.await
.expect("timeout waiting for incoming")
.expect("no incoming");
assert_eq!(tx.original.method, Method::Invite);
let tx_key = tx.key.clone();
tx.reply(StatusCode::ServiceUnavailable)
.await
.expect("reply 503");
while let Some(msg) = tx.receive().await {
if let crate::sip::SipMessage::Request(req) = msg {
if req.method == Method::Ack {
break;
}
}
}
drop(tx);
sleep(Duration::from_millis(100)).await;
assert!(
endpoint.inner.finished_transactions.contains_key(&tx_key),
"finished_transactions should contain the response after normal termination"
);
};
let client_loop = async {
let mut buf = vec![0u8; 4096];
let (len, _) = timeout(Duration::from_secs(2), client_conn.recv_raw(&mut buf))
.await
.expect("timeout waiting for 503")
.expect("recv failed");
assert!(String::from_utf8_lossy(&buf[..len]).contains("503"));
sleep(Duration::from_millis(50)).await;
let ack = make_ack_for_invite(&invite);
client_conn_sip
.send(ack.into(), Some(&server_addr))
.await
.expect("send ack");
};
tokio::join!(server_loop, client_loop);
serve_handle.abort();
}
#[tokio::test]
async fn test_finished_transactions_replies_correct_status() {
let token = CancellationToken::new();
let server_conn = UdpConnection::create_connection("127.0.0.1:0".parse().unwrap(), None, None)
.await
.expect("create server connection");
let server_conn_sip: SipConnection = server_conn.clone().into();
let server_addr = server_conn_sip.get_addr().clone();
let tl = TransportLayer::new(token.child_token());
tl.add_transport(server_conn_sip.clone());
let endpoint = EndpointBuilder::new()
.with_user_agent("rsipstack-test")
.with_transport_layer(tl)
.build();
let client_conn = UdpConnection::create_connection("127.0.0.1:0".parse().unwrap(), None, None)
.await
.expect("create client connection");
let client_conn_sip: SipConnection = client_conn.clone().into();
let endpoint_inner = endpoint.inner.clone();
let serve_handle = tokio::spawn(async move {
let _ = endpoint_inner.serve().await;
});
let branch = "z9hG4bK-status-check-branch";
let invite = make_invite(branch, "serverTagStatus");
client_conn_sip
.send(invite.clone().into(), Some(&server_addr))
.await
.expect("send invite");
let mut incoming = endpoint.incoming_transactions().expect("incoming");
let mut tx = timeout(Duration::from_secs(2), incoming.recv())
.await
.expect("timeout waiting for incoming")
.expect("no incoming");
tx.reply(StatusCode::BusyHere).await.expect("reply 486");
drop(tx);
sleep(Duration::from_millis(100)).await;
let mut buf = vec![0u8; 4096];
let (len, _) = timeout(Duration::from_secs(2), client_conn.recv_raw(&mut buf))
.await
.expect("timeout waiting for 486")
.expect("recv failed");
let resp = String::from_utf8_lossy(&buf[..len]);
assert!(resp.contains("486"), "should receive 486");
sleep(Duration::from_millis(50)).await;
client_conn_sip
.send(invite.clone().into(), Some(&server_addr))
.await
.expect("retransmit invite");
let (len, _) = timeout(Duration::from_secs(2), client_conn.recv_raw(&mut buf))
.await
.expect("timeout waiting for retransmitted 486")
.expect("recv failed");
let resp = String::from_utf8_lossy(&buf[..len]);
assert!(
resp.contains("486"),
"retransmitted INVITE should get 486 auto-reply"
);
serve_handle.abort();
}
#[tokio::test]
async fn test_cleanup_server_invite_confirmed_drop_removes_waiting_ack() -> crate::Result<()> {
let cancel_token = CancellationToken::new();
let tl = TransportLayer::new_with_domain_resolver(cancel_token, Box::new(NoopResolver));
let endpoint = EndpointBuilder::new()
.with_user_agent("rsipstack-test")
.with_transport_layer(tl)
.build();
let invite = make_invite("z9hG4bK-confirmed-drop", "serverTagConfirmedDrop");
let key = TransactionKey::from_request(&invite, TransactionRole::Server)?;
let mut tx = Transaction::new_server(key.clone(), invite.clone(), endpoint.inner.clone(), None);
let resp = crate::sip::Response {
status_code: StatusCode::ServiceUnavailable,
version: Version::V2,
headers: invite.headers.clone(),
body: Default::default(),
};
tx.last_response = Some(resp.clone());
tx.state = crate::transaction::TransactionState::Confirmed;
let dialog_id = crate::dialog::DialogId::try_from((&resp, TransactionRole::Server))?;
endpoint
.inner
.waiting_ack
.insert(dialog_id.clone(), key.clone());
drop(tx);
sleep(Duration::from_millis(50)).await;
assert!(
!endpoint.inner.waiting_ack.contains_key(&dialog_id),
"waiting_ack must be removed when ServerInvite is dropped in Confirmed state (ACK already received)"
);
Ok(())
}
#[tokio::test]
async fn test_timer_cleanup_removes_orphaned_waiting_ack() -> crate::Result<()> {
use crate::transaction::endpoint::EndpointOption;
let cancel_token = CancellationToken::new();
let tl = TransportLayer::new_with_domain_resolver(cancel_token.clone(), Box::new(NoopResolver));
let endpoint = EndpointBuilder::new()
.with_user_agent("rsipstack-test")
.with_transport_layer(tl)
.with_option(EndpointOption {
t1x64: Duration::from_millis(100),
..Default::default()
})
.build();
let inner = endpoint.inner.clone();
let serve_handle = tokio::spawn(async move {
let _ = inner.serve().await;
});
let invite = make_invite("z9hG4bK-timer-cleanup", "serverTagTimerCleanup");
let key = TransactionKey::from_request(&invite, TransactionRole::Server)?;
let mut tx = Transaction::new_server(key.clone(), invite.clone(), endpoint.inner.clone(), None);
let resp = crate::sip::Response {
status_code: StatusCode::ServiceUnavailable,
version: Version::V2,
headers: invite.headers.clone(),
body: Default::default(),
};
tx.last_response = Some(resp.clone());
tx.state = crate::transaction::TransactionState::Completed;
let dialog_id = crate::dialog::DialogId::try_from((&resp, TransactionRole::Server))?;
endpoint
.inner
.waiting_ack
.insert(dialog_id.clone(), key.clone());
drop(tx);
sleep(Duration::from_millis(50)).await;
assert!(
endpoint.inner.waiting_ack.contains_key(&dialog_id),
"waiting_ack should still exist after drop from Completed (waiting for ACK)"
);
assert!(
endpoint.inner.finished_transactions.contains_key(&key),
"finished_transactions should exist after drop"
);
sleep(Duration::from_millis(500)).await;
assert!(
!endpoint.inner.finished_transactions.contains_key(&key),
"finished_transactions should be cleaned up by TimerCleanup"
);
assert!(
!endpoint.inner.waiting_ack.contains_key(&dialog_id),
"waiting_ack must be cleaned up by TimerCleanup (orphaned entry safety net)"
);
serve_handle.abort();
cancel_token.cancel();
Ok(())
}